From 78e62af9d81fdcb608253eeba6513aac768f81c6 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Tue, 16 Dec 2025 06:03:55 +0000 Subject: [PATCH] nix,spindle: sync workflow files on `sh.tangled.git.refUpdate` Spindle will sync git repo when new repo is registered Spindle will listen to `sh.tangled.git.refUpdate` event from knot stream and sync its local git repo instead. Spindle's git repo will sparse-checkout only `/.tangled/workflows` directory. Spindle now requires git version >=2.49 for `--revision` flag in `git clone` command. References: - - Signed-off-by: Seongmin Lee --- go.mod | 1 + go.sum | 2 ++ nix/gomod2nix.toml | 3 +++ spindle/server.go | 71 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----- spindle/tap.go | 14 ++++++++++++-- nix/modules/spindle.nix | 4 ++++ spindle/config/config.go | 4 ++++ spindle/git/git.go | 73 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 8 file(s) changed, 165 insertion(s)(+), 7 deletion(s)(-) diff --git a/go.mod b/go.mod --- a/go.mod +++ b/go.mod @@ -33,6 +33,7 @@ github.com/gorilla/feeds v1.2.0 github.com/gorilla/sessions v1.4.0 github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674 + github.com/hashicorp/go-version v1.8.0 github.com/hiddeco/sshsig v0.2.0 github.com/hpcloud/tail v1.0.0 github.com/ipfs/go-cid v0.5.0 diff --git a/go.sum b/go.sum --- a/go.sum +++ b/go.sum @@ -337,6 +337,8 @@ github.com/hashicorp/go-secure-stdlib/strutil v0.1.2/go.mod h1:Gou2R9+il93BqX25LAKCLuM+y9U2T4hlwvT1yprcna4= github.com/hashicorp/go-sockaddr v1.0.7 h1:G+pTkSO01HpR5qCxg7lxfsFEZaG+C0VssTy/9dbT+Fw= github.com/hashicorp/go-sockaddr v1.0.7/go.mod h1:FZQbEYa1pxkQ7WLpyXJ6cbjpT8q0YgQaK/JakXqGyWw= +github.com/hashicorp/go-version v1.8.0 h1:KAkNb1HAiZd1ukkxDFGmokVZe1Xy9HG6NUp+bPle2i4= +github.com/hashicorp/go-version v1.8.0/go.mod h1:fltr4n8CU8Ke44wwGCBoEymUuxUHl09ZGVZPK5anwXA= github.com/hashicorp/golang-lru v1.0.2 h1:dV3g9Z/unq5DpblPpw+Oqcv4dU/1omnb4Ok8iPY6p1c= github.com/hashicorp/golang-lru v1.0.2/go.mod h1:iADmTwqILo4mZ8BN3D2Q6+9jd8WM5uGBxy+E8yxSoD4= github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= diff --git a/nix/gomod2nix.toml b/nix/gomod2nix.toml --- a/nix/gomod2nix.toml +++ b/nix/gomod2nix.toml @@ -356,6 +356,9 @@ [mod."github.com/hashicorp/go-sockaddr"] version = "v1.0.7" hash = "sha256-p6eDOrGzN1jMmT/F/f/VJMq0cKNFhUcEuVVwTE6vSrs=" + [mod."github.com/hashicorp/go-version"] + version = "v1.8.0" + hash = "sha256-KXtqERmYrWdpqPCViWcHbe6jnuH7k16bvBIcuJuevj8=" [mod."github.com/hashicorp/golang-lru"] version = "v1.0.2" hash = "sha256-yy+5botc6T5wXgOe2mfNXJP3wr+MkVlUZ2JBkmmrA48=" diff --git a/spindle/server.go b/spindle/server.go --- a/spindle/server.go +++ b/spindle/server.go @@ -8,12 +8,14 @@ "log/slog" "maps" "net/http" + "path/filepath" "sync" "time" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/service/tap" "github.com/go-chi/chi/v5" + "github.com/hashicorp/go-version" "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" "tangled.org/core/eventconsumer/cursor" @@ -25,6 +27,7 @@ "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/engines/nixery" + "tangled.org/core/spindle/git" "tangled.org/core/spindle/models" "tangled.org/core/spindle/queue" "tangled.org/core/spindle/secrets" @@ -56,12 +59,16 @@ func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engine) (*Spindle, error) { logger := log.FromContext(ctx) - d, err := db.Make(ctx, cfg.Server.DBPath) + if err := ensureGitVersion(); err != nil { + return nil, fmt.Errorf("ensuring git version: %w", err) + } + + d, err := db.Make(ctx, cfg.Server.DBPath()) if err != nil { return nil, fmt.Errorf("failed to setup db: %w", err) } - e, err := rbac2.NewEnforcer(cfg.Server.DBPath) + e, err := rbac2.NewEnforcer(cfg.Server.DBPath()) if err != nil { return nil, fmt.Errorf("failed to setup rbac enforcer: %w", err) } @@ -84,11 +91,11 @@ } logger.Info("using openbao secrets provider", "proxy_address", cfg.Server.Secrets.OpenBao.ProxyAddr, "mount", cfg.Server.Secrets.OpenBao.Mount) case "sqlite", "": - vault, err = secrets.NewSQLiteManager(cfg.Server.DBPath, secrets.WithTableName("secrets")) + vault, err = secrets.NewSQLiteManager(cfg.Server.DBPath(), secrets.WithTableName("secrets")) if err != nil { return nil, fmt.Errorf("failed to setup sqlite secrets provider: %w", err) } - logger.Info("using sqlite secrets provider", "path", cfg.Server.DBPath) + logger.Info("using sqlite secrets provider", "path", cfg.Server.DBPath()) default: return nil, fmt.Errorf("unknown secrets provider: %s", cfg.Server.Secrets.Provider) } @@ -120,7 +127,7 @@ } logger.Info("owner set", "did", cfg.Server.Owner) - cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath) + cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath()) if err != nil { return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err) } @@ -312,7 +319,10 @@ } func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, msg eventconsumer.Message) error { + l := log.FromContext(ctx).With("handler", "processKnotStream") + l = l.With("src", src.Key(), "msg.Nsid", msg.Nsid, "msg.Rkey", msg.Rkey) if msg.Nsid == tangled.PipelineNSID { + return nil tpl := tangled.Pipeline{} err := json.Unmarshal(msg.EventJson, &tpl) if err != nil { @@ -413,7 +423,58 @@ } else { s.l.Error("failed to enqueue pipeline: queue is full") } + } else if msg.Nsid == tangled.GitRefUpdateNSID { + event := tangled.GitRefUpdate{} + if err := json.Unmarshal(msg.EventJson, &event); err != nil { + l.Error("error unmarshalling", "err", err) + return err + } + l = l.With("repoDid", event.RepoDid, "repoName", event.RepoName) + + // resolve repo name to rkey + // TODO: git.refUpdate should respond with rkey instead of repo name + repo, err := s.db.GetRepoWithName(syntax.DID(event.RepoDid), event.RepoName) + if err != nil { + return fmt.Errorf("get repo with did and name (%s/%s): %w", event.RepoDid, event.RepoName, err) + } + + // NOTE: we are blindly trusting the knot that it will return only repos it own + repoCloneUri := s.newRepoCloneUrl(src.Key(), event.RepoDid, event.RepoName) + repoPath := s.newRepoPath(repo.Did, repo.Rkey) + if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, event.NewSha); err != nil { + return fmt.Errorf("sync git repo: %w", err) + } + l.Info("synced git repo") + + // TODO: plan the pipeline } + return nil +} + +// newRepoPath creates a path to store repository by its did and rkey. +// The path format would be: `/data/repos/did:plc:foo/sh.tangled.repo/repo-rkey +func (s *Spindle) newRepoPath(did syntax.DID, rkey syntax.RecordKey) string { + return filepath.Join(s.cfg.Server.RepoDir(), did.String(), tangled.RepoNSID, rkey.String()) +} + +func (s *Spindle) newRepoCloneUrl(knot, did, name string) string { + scheme := "https://" + if s.cfg.Server.Dev { + scheme = "http://" + } + return fmt.Sprintf("%s%s/%s/%s", scheme, knot, did, name) +} + +const RequiredVersion = "2.49.0" + +func ensureGitVersion() error { + v, err := git.Version() + if err != nil { + return fmt.Errorf("fetching git version: %w", err) + } + if v.LessThan(version.Must(version.NewVersion(RequiredVersion))) { + return fmt.Errorf("installed git version %q is not supported, Spindle requires git version >= %q", v, RequiredVersion) + } return nil } diff --git a/spindle/tap.go b/spindle/tap.go --- a/spindle/tap.go +++ b/spindle/tap.go @@ -10,6 +10,7 @@ "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" "tangled.org/core/spindle/db" + "tangled.org/core/spindle/git" "tangled.org/core/tapc" ) @@ -225,12 +226,14 @@ return nil } - if err := s.db.PutRepo(&db.Repo{ + repo := &db.Repo{ Did: evt.Record.Did, Rkey: evt.Record.Rkey, Name: record.Name, Knot: record.Knot, - }); err != nil { + } + + if err := s.db.PutRepo(repo); err != nil { return fmt.Errorf("adding repo to db: %w", err) } @@ -241,6 +244,13 @@ // add this knot to the event consumer src := eventconsumer.NewKnotSource(record.Knot) s.ks.AddSource(context.Background(), src) + + // setup sparse sync + repoCloneUri := s.newRepoCloneUrl(repo.Knot, repo.Did.String(), repo.Name) + repoPath := s.newRepoPath(repo.Did, repo.Rkey) + if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, ""); err != nil { + return fmt.Errorf("setting up sparse-clone git repo: %w", err) + } l.Info("added repo", "repo", evt.Record.AtUri()) return nil diff --git a/nix/modules/spindle.nix b/nix/modules/spindle.nix --- a/nix/modules/spindle.nix +++ b/nix/modules/spindle.nix @@ -1,5 +1,6 @@ { config, + pkgs, lib, ... }: let @@ -132,6 +133,9 @@ description = "spindle service"; after = ["network.target" "docker.service"]; wantedBy = ["multi-user.target"]; + path = [ + pkgs.git + ]; serviceConfig = { LogsDirectory = "spindle"; StateDirectory = "spindle"; diff --git a/spindle/config/config.go b/spindle/config/config.go --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -29,6 +29,10 @@ return syntax.DID(fmt.Sprintf("did:web:%s", s.Hostname)) } +func (s Server) RepoDir() string { + return filepath.Join(s.DataDir, "repos") +} + func (s Server) DBPath() string { return filepath.Join(s.DataDir, "spindle.db") } diff --git a/spindle/git/git.go b/spindle/git/git.go new file mode 100644 --- /dev/null +++ b/spindle/git/git.go @@ -0,0 +1,73 @@ +package git + +import ( + "bytes" + "context" + "fmt" + "os" + "os/exec" + "strings" + + "github.com/hashicorp/go-version" +) + +func Version() (*version.Version, error) { + var buf bytes.Buffer + cmd := exec.Command("git", "version") + cmd.Stdout = &buf + cmd.Stderr = os.Stderr + err := cmd.Run() + if err != nil { + return nil, err + } + fields := strings.Fields(buf.String()) + if len(fields) < 3 { + return nil, fmt.Errorf("invalid git version: %s", buf.String()) + } + + // version string is like: "git version 2.29.3" or "git version 2.29.3.windows.1" + versionString := fields[2] + if pos := strings.Index(versionString, "windows"); pos >= 1 { + versionString = versionString[:pos-1] + } + return version.NewVersion(versionString) +} + +const WorkflowDir = `/.tangled/workflows` + +func SparseSyncGitRepo(ctx context.Context, cloneUri, path, rev string) error { + exist, err := isDir(path) + if err != nil { + return err + } + if rev == "" { + rev = "HEAD" + } + if !exist { + if err := exec.Command("git", "clone", "--no-checkout", "--depth=1", "--filter=tree:0", "--revision="+rev, cloneUri, path).Run(); err != nil { + return fmt.Errorf("git clone: %w", err) + } + if err := exec.Command("git", "-C", path, "sparse-checkout", "set", "--no-cone", WorkflowDir).Run(); err != nil { + return fmt.Errorf("git sparse-checkout set: %w", err) + } + } else { + if err := exec.Command("git", "-C", path, "fetch", "--depth=1", "--filter=tree:0", "origin", rev).Run(); err != nil { + return fmt.Errorf("git pull: %w", err) + } + } + if err := exec.Command("git", "-C", path, "checkout", rev).Run(); err != nil { + return fmt.Errorf("git checkout: %w", err) + } + return nil +} + +func isDir(path string) (bool, error) { + info, err := os.Stat(path) + if err == nil && info.IsDir() { + return true, nil + } + if os.IsNotExist(err) { + return false, nil + } + return false, err +} -- tangled.sh