diff --git a/spindle/db/collaborators.go b/spindle/db/collaborators.go index 2a5344c08..24848f888 100644 --- a/spindle/db/collaborators.go +++ b/spindle/db/collaborators.go @@ -113,3 +113,28 @@ func (d *DB) ListCollaboratorsByRepoDid(repoDid syntax.DID) ([]RepoCollaborator, } return out, nil } + +func (d *DB) ListKnotCollaboratorsByRepoDid(repoDid syntax.DID) ([]RepoCollaborator, error) { + rows, err := d.Query( + `select owner_did, rkey, subject, repo_did from repo_collaborators + where repo_did = ? and owner_did = ?`, + repoDid.String(), repoDid.String(), + ) + if err != nil { + return nil, fmt.Errorf("list knot collaborators for %s: %w", repoDid, err) + } + defer rows.Close() + + var out []RepoCollaborator + for rows.Next() { + c, err := scanCollab(rows) + if err != nil { + return nil, err + } + out = append(out, *c) + } + if err := rows.Err(); err != nil { + return nil, err + } + return out, nil +} diff --git a/spindle/db/db.go b/spindle/db/db.go index 72ee4a966..c93b8fcb4 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -57,6 +57,20 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { unique(owner, rkey) ); + + create table if not exists feed_cursors ( + knot text primary key, + seq integer not null + ); + + create table if not exists feed_refs ( + repo_did text not null, + rkey text not null, + sha text not null, + + primary key (repo_did, rkey) + ); + create table if not exists repo_collaborators ( id integer primary key autoincrement, owner_did text not null, diff --git a/spindle/db/feed.go b/spindle/db/feed.go new file mode 100644 index 000000000..9300d9774 --- /dev/null +++ b/spindle/db/feed.go @@ -0,0 +1,72 @@ +package db + +import ( + "database/sql" + "errors" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +func (d *DB) LoadFeedCursor(knot string) (int64, error) { + var seq int64 + err := d.QueryRow(`select seq from feed_cursors where knot = ?`, knot).Scan(&seq) + if errors.Is(err, sql.ErrNoRows) { + return 0, nil + } + return seq, err +} + +func (d *DB) StoreFeedCursor(knot string, seq int64) error { + _, err := d.Exec( + `insert into feed_cursors (knot, seq) values (?, ?) + on conflict(knot) do update set seq = excluded.seq`, + knot, seq, + ) + return err +} + +func (d *DB) FeedRefSha(repoDid syntax.DID, rkey string) (string, bool, error) { + var sha string + err := d.QueryRow( + `select sha from feed_refs where repo_did = ? and rkey = ?`, + repoDid.String(), rkey, + ).Scan(&sha) + if errors.Is(err, sql.ErrNoRows) { + return "", false, nil + } + if err != nil { + return "", false, err + } + return sha, true, nil +} + +func (d *DB) PutFeedRef(repoDid syntax.DID, rkey, sha string) error { + _, err := d.Exec( + `insert into feed_refs (repo_did, rkey, sha) values (?, ?, ?) + on conflict(repo_did, rkey) do update set sha = excluded.sha`, + repoDid.String(), rkey, sha, + ) + return err +} + +func (d *DB) SeedFeedRef(repoDid syntax.DID, rkey, sha string) error { + _, err := d.Exec( + `insert into feed_refs (repo_did, rkey, sha) values (?, ?, ?) + on conflict(repo_did, rkey) do nothing`, + repoDid.String(), rkey, sha, + ) + return err +} + +func (d *DB) DeleteFeedRef(repoDid syntax.DID, rkey string) error { + _, err := d.Exec( + `delete from feed_refs where repo_did = ? and rkey = ?`, + repoDid.String(), rkey, + ) + return err +} + +func (d *DB) DeleteFeedRefsByRepoDid(repoDid syntax.DID) error { + _, err := d.Exec(`delete from feed_refs where repo_did = ?`, repoDid.String()) + return err +} diff --git a/spindle/db/repos.go b/spindle/db/repos.go index 79179a8f3..54aaaf67c 100644 --- a/spindle/db/repos.go +++ b/spindle/db/repos.go @@ -205,3 +205,9 @@ func (d *DB) CountReposByOwner(owner syntax.DID) (int, error) { err := d.QueryRow(`select count(*) from repos where owner = ?`, owner.String()).Scan(&n) return n, err } + +func (d *DB) CountReposByKnot(knot string) (int64, error) { + var n int64 + err := d.QueryRow(`select count(*) from repos where knot = ?`, knot).Scan(&n) + return n, err +} diff --git a/spindle/feed/feed.go b/spindle/feed/feed.go new file mode 100644 index 000000000..dc2f167a2 --- /dev/null +++ b/spindle/feed/feed.go @@ -0,0 +1,118 @@ +package feed + +import ( + "context" + "log/slog" + "sync" + + "tangled.org/core/hostutil" + knotfeed "tangled.org/core/knotfeed" +) + +type Hooks struct { + LoadCursor func(ctx context.Context, knot string) (int64, error) + StoreCursor func(ctx context.Context, knot string, seq int64) error + Handle func(ctx context.Context, knot string, msg knotfeed.Message) error + OutdatedReplay func(ctx context.Context, knot string) int64 + OnConnectError func(knot string, err error) +} + +type Feed struct { + logger *slog.Logger + hooks Hooks + + mu sync.Mutex + subs map[string]*subscriber +} + +type subscriber struct { + cancel context.CancelFunc +} + +func New(logger *slog.Logger, hooks Hooks) *Feed { + return &Feed{ + logger: logger, + hooks: hooks, + subs: make(map[string]*subscriber), + } +} + +func (f *Feed) Start(ctx context.Context, knots []string) { + for _, knot := range knots { + f.Subscribe(ctx, knot) + } +} + +func (f *Feed) Subscribe(ctx context.Context, knot string) { + host, noTLS, err := hostutil.ParseHostname(knot) + if err != nil { + f.logger.Error("unsubscribable knot host", "knot", knot, "err", err) + return + } + + f.mu.Lock() + if _, ok := f.subs[knot]; ok { + f.mu.Unlock() + return + } + runCtx, cancel := context.WithCancel(ctx) + sub := &subscriber{cancel: cancel} + f.subs[knot] = sub + f.mu.Unlock() + + consumer := &knotfeed.Consumer{ + Host: host, + NoTLS: noTLS, + Logger: f.logger, + LoadCursor: func(ctx context.Context) (int64, error) { + return f.hooks.LoadCursor(ctx, knot) + }, + StoreCursor: func(ctx context.Context, seq int64) error { + return f.hooks.StoreCursor(ctx, knot, seq) + }, + Handle: func(ctx context.Context, msg knotfeed.Message) error { + return f.hooks.Handle(ctx, knot, msg) + }, + OutdatedReplay: func(ctx context.Context) int64 { + if f.hooks.OutdatedReplay == nil { + return 0 + } + return f.hooks.OutdatedReplay(ctx, knot) + }, + } + if f.hooks.OnConnectError != nil { + consumer.OnConnectError = func(err error) { + f.hooks.OnConnectError(knot, err) + } + } + + go f.run(runCtx, knot, sub, consumer) +} + +func (f *Feed) Unsubscribe(knot string) { + f.mu.Lock() + sub, ok := f.subs[knot] + if !ok { + f.mu.Unlock() + return + } + delete(f.subs, knot) + f.mu.Unlock() + + sub.cancel() +} + +func (f *Feed) run(ctx context.Context, knot string, sub *subscriber, consumer *knotfeed.Consumer) { + defer f.drop(knot, sub) + if err := consumer.Run(ctx); err != nil && ctx.Err() == nil { + f.logger.Error("knot feed exited", "knot", knot, "err", err) + } +} + +func (f *Feed) drop(knot string, sub *subscriber) { + f.mu.Lock() + defer f.mu.Unlock() + if f.subs[knot] == sub { + delete(f.subs, knot) + } +} diff --git a/spindle/feed/feed_test.go b/spindle/feed/feed_test.go new file mode 100644 index 000000000..d3ddb90c5 --- /dev/null +++ b/spindle/feed/feed_test.go @@ -0,0 +1,77 @@ +package feed + +import ( + "context" + "errors" + "io" + "log/slog" + "testing" + "time" +) + +func TestFeedRefcountsAndReleasesSubscriptions(t *testing.T) { + entered := make(chan string, 8) + held := make(chan struct{}) + f := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Hooks{ + LoadCursor: func(ctx context.Context, knot string) (int64, error) { + select { + case entered <- knot: + case <-ctx.Done(): + return 0, ctx.Err() + } + select { + case <-held: + case <-ctx.Done(): + } + return 0, errors.New("no network in tests") + }, + }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + const knot = "knot.nel.pet" + f.Subscribe(ctx, knot) + f.Subscribe(ctx, knot) + select { + case got := <-entered: + if got != knot { + t.Fatalf("first session loaded the cursor for %q, want %q", got, knot) + } + case <-time.After(2 * time.Second): + t.Fatal("no session started for a subscribed knot") + } + + f.Unsubscribe(knot) + f.Unsubscribe(knot) + close(held) + + f.Subscribe(ctx, knot) + select { + case got := <-entered: + if got != knot { + t.Fatalf("second session loaded the cursor for %q, want %q", got, knot) + } + case <-time.After(2 * time.Second): + t.Fatalf("a fresh session didn't start after the knot was released and resubscribed") + } +} + +func TestFeedRejectsUnsubscribableHost(t *testing.T) { + entered := make(chan string, 1) + f := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Hooks{ + LoadCursor: func(ctx context.Context, knot string) (int64, error) { + entered <- knot + return 0, errors.New("no network in tests") + }, + }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + f.Subscribe(ctx, "192.0.2.1") + f.Unsubscribe("192.0.2.1") + select { + case knot := <-entered: + t.Fatalf("a session started for an unparsable host %q", knot) + case <-time.After(100 * time.Millisecond): + } +} diff --git a/spindle/knotfeed.go b/spindle/knotfeed.go new file mode 100644 index 000000000..46662e8bc --- /dev/null +++ b/spindle/knotfeed.go @@ -0,0 +1,489 @@ +package spindle + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "io" + "log/slog" + "math/rand/v2" + "net/http" + "net/url" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/gitutil" + "tangled.org/core/hostutil" + knotfeed "tangled.org/core/knotfeed" + kgit "tangled.org/core/knotserver/git" + "tangled.org/core/log" + "tangled.org/core/rbac" + "tangled.org/core/spindle/db" + "tangled.org/core/workflow" +) + +const collaboratorInviteNSID = "sh.tangled.repo.collaboratorInvite" + +const ( + reconcileInterval = 10 * time.Minute + reconcileJitter = time.Minute + maxListPages = 256 + + feedWipeAttempts = 3 + feedWipeRetryWait = 100 * time.Millisecond +) + +const ( + changedFileOverheadBytes = 48 + changedFilesMaxBytes = 524288 + changedFilesMaxCount = 8192 +) + +var knotHTTPClient = &http.Client{Timeout: 30 * time.Second} + +type refRecord struct { + Rkey string + Sha string +} + +func (s *Spindle) handleKnotFeed(ctx context.Context, knot string, msg knotfeed.Message) error { + if msg.Type != knotfeed.TypeCommit || msg.Commit == nil { + return nil + } + var errs []error + var invites []syntax.DID + seen := make(map[syntax.DID]struct{}) + for _, op := range msg.Commit.Records { + repoDid := syntax.DID(msg.Commit.Repo) + switch op.Collection { + case knotfeed.GitRefCollection: + if err := s.handleRefOp(ctx, knot, repoDid, op); err != nil { + errs = append(errs, err) + } + case collaboratorInviteNSID: + if _, ok := seen[repoDid]; !ok { + seen[repoDid] = struct{}{} + invites = append(invites, repoDid) + } + } + } + for _, repoDid := range invites { + if err := s.reconcileCollaborators(ctx, knot, repoDid); err != nil { + errs = append(errs, err) + } + } + return errors.Join(errs...) +} + +func (s *Spindle) owningRepo(l *slog.Logger, knot string, repoDid syntax.DID, event string) (*db.Repo, bool, error) { + repo, err := s.db.GetRepoByDid(repoDid) + switch { + case errors.Is(err, sql.ErrNoRows): + l.Info(fmt.Sprintf("skipping %s for unknown repo", event)) + return nil, false, nil + case err != nil: + return nil, false, fmt.Errorf("lookup repo %s: %w", repoDid, err) + case repo.Knot != knot: + l.Info(fmt.Sprintf("dropping %s from non-owning knot", event), "repoKnot", repo.Knot) + return nil, false, nil + } + return repo, true, nil +} + +func (s *Spindle) handleRefOp(ctx context.Context, knot string, repoDid syntax.DID, op knotfeed.RecordOp) error { + l := log.FromContext(ctx).With("repo", repoDid, "rkey", op.Rkey) + + repo, proceed, err := s.owningRepo(l, knot, repoDid, "ref update") + if !proceed || err != nil { + return err + } + + refname, ok := knotfeed.UnescapeRkey(op.Rkey) + if !ok { + l.Info("skipping ref update with undecodable rkey") + return nil + } + if !isMaterializedRef(refname) { + l.Info("skipping ref update outside heads and tags", "ref", refname) + return nil + } + + if err := s.ensureRefState(ctx, knot, repoDid); err != nil { + return fmt.Errorf("bootstrapping ref state: %w", err) + } + + if op.Deleted() { + if err := s.db.DeleteFeedRef(repoDid, op.Rkey); err != nil { + return fmt.Errorf("forgetting ref state: %w", err) + } + l.Info("ref deleted, nothing to trigger", "ref", refname) + return nil + } + + record, err := knotfeed.DecodeRefRecord(op.Bytes) + if err != nil { + return fmt.Errorf("decoding ref record: %w", err) + } + + oldSha, _, err := s.db.FeedRefSha(repoDid, op.Rkey) + if err != nil { + return fmt.Errorf("reading ref state: %w", err) + } + // deliver push webhooks independently of CI; a skip-ci push option + // should not suppress webhook notifications + s.wh.FirePush(ctx, repo, string(record.Editor), refname, oldSha.String(), record.Sha.String()) + + if kgit.HasSkipCIPushOption(record.PushOptions) { + if err := s.db.PutFeedRef(repoDid, op.Rkey, record.Sha); err != nil { + return fmt.Errorf("recording ref state: %w", err) + } + l.Info("push requested ci skip, skipping the event", "ref", refname) + return nil + } + + repoCloneUri := s.newRepoCloneUrl(knot, repoDid) + repoPath := s.newRepoPath(repoDid) + if err := gitutil.SparseSync(ctx, repoCloneUri, repoPath, record.Sha, sparseWorkflowDir); err != nil { + return fmt.Errorf("sync git repo: %w", err) + } + l.Info("synced git repo") + + changedFiles := changedFilesUnderBudget(l, repoPath, oldSha, record.Sha) + + if err := s.db.PutFeedRef(repoDid, op.Rkey, record.Sha); err != nil { + return fmt.Errorf("recording ref state: %w", err) + } + + triggerRepo, err := s.buildTriggerRepo(ctx, repo) + if err != nil { + return fmt.Errorf("building trigger repo: %w", err) + } + + trigger := tangled.Pipeline_TriggerMetadata{ + Kind: string(workflow.TriggerKindPush), + Push: &tangled.Pipeline_PushTriggerData{ + Ref: refname, + OldSha: oldSha, + NewSha: record.Sha, + }, + Repo: triggerRepo, + } + + pipelineId, err := s.runPipeline(ctx, repoDid, trigger, changedFiles, repoCloneUri, repoPath, record.Sha, nil, triggerRepo) + if err != nil { + return err + } + if pipelineId.Rkey == "" { + l.Info("no workflow matched 'push' trigger, skipping the event") + return nil + } + l.Info("pipeline triggered", "pipeline", pipelineId.AtUri()) + return nil +} + +func isMaterializedRef(refname string) bool { + return strings.HasPrefix(refname, "refs/heads/") || strings.HasPrefix(refname, "refs/tags/") +} + +func (s *Spindle) ensureRefState(ctx context.Context, knot string, repoDid syntax.DID) error { + if _, seeded := s.refStateSeeded.LoadOrStore(repoDid.String(), struct{}{}); seeded { + return nil + } + refs, err := s.refRecords(ctx, knot, repoDid) + if err != nil { + s.refStateSeeded.Delete(repoDid.String()) + return err + } + for _, ref := range refs { + if err := s.db.SeedFeedRef(repoDid, ref.Rkey, ref.Sha); err != nil { + s.refStateSeeded.Delete(repoDid.String()) + return fmt.Errorf("seeding ref state: %w", err) + } + } + return nil +} + +func changedFilesUnderBudget(l *slog.Logger, repoPath, oldSha, newSha string) []string { + if oldSha == "" { + return nil + } + gr, err := kgit.Open(repoPath, newSha) + if err != nil { + l.Warn("cannot open synced repo for changed files", "err", err) + return nil + } + paths, err := gr.ChangedFilesBetween(oldSha, newSha) + if err != nil { + l.Warn("changed files unavailable between revisions", "oldSha", oldSha, "newSha", newSha, "err", err) + return nil + } + return admitChangedFiles(paths) +} + +func admitChangedFiles(paths []string) []string { + spent := 0 + admitted := make([]string, 0, min(len(paths), changedFilesMaxCount)) + for _, path := range paths { + spent += changedFileOverheadBytes + len(path) + if spent > changedFilesMaxBytes || len(admitted) == changedFilesMaxCount { + return nil + } + admitted = append(admitted, path) + } + return admitted +} + +func (s *Spindle) reconcileCollaborators(ctx context.Context, knot string, repoDid syntax.DID) error { + l := log.FromContext(ctx).With("repo", repoDid) + + _, proceed, err := s.owningRepo(l, knot, repoDid, "collaborator reconcile") + if !proceed || err != nil { + return err + } + + desired, err := s.collaboratorList(ctx, knot, repoDid) + if err != nil { + return fmt.Errorf("listing collaborators: %w", err) + } + tracked, err := s.db.ListKnotCollaboratorsByRepoDid(repoDid) + if err != nil { + return fmt.Errorf("listing tracked collaborators: %w", err) + } + + desiredSet := make(map[syntax.DID]struct{}, len(desired)) + for _, subject := range desired { + desiredSet[subject] = struct{}{} + } + trackedSet := make(map[syntax.DID]struct{}, len(tracked)) + for _, c := range tracked { + trackedSet[c.Subject] = struct{}{} + } + + for subject := range desiredSet { + if _, ok := trackedSet[subject]; ok { + continue + } + if err := s.e.AddCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { + l.Error("failed to add collaborator policy", "subject", subject, "err", err) + continue + } + if err := s.db.AddKnotCollaborator(repoDid, subject); err != nil { + l.Error("failed to track collaborator", "subject", subject, "err", err) + continue + } + l.Info("added knot-managed collaborator", "subject", subject) + } + for subject := range trackedSet { + if _, ok := desiredSet[subject]; ok { + continue + } + if err := s.e.RemoveCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { + l.Error("failed to remove collaborator policy", "subject", subject, "err", err) + continue + } + if err := s.db.DeleteRepoCollaboratorBySubjectRepo(subject, repoDid); err != nil { + l.Error("failed to delete collaborator row", "subject", subject, "err", err) + continue + } + l.Info("removed knot-managed collaborator", "subject", subject) + } + return nil +} + +func (s *Spindle) reconcileCollaboratorsLoop(ctx context.Context) { + s.reconcileAllCollaborators(ctx) + for { + delay := reconcileInterval + time.Duration(rand.Int64N(int64(reconcileJitter))) + select { + case <-ctx.Done(): + return + case <-time.After(delay): + } + s.reconcileAllCollaborators(ctx) + } +} + +func (s *Spindle) reconcileAllCollaborators(ctx context.Context) { + repos, err := s.db.AllRepos() + if err != nil { + s.l.Warn("failed to load repos for collaborator reconcile", "err", err) + return + } + for _, repo := range repos { + if repo.RepoDid == "" { + continue + } + if err := s.reconcileCollaborators(ctx, repo.Knot, repo.RepoDid); err != nil { + s.l.Warn("collaborator reconcile failed", "repo", repo.RepoDid, "err", err) + } + } +} + +func (s *Spindle) feedOutdatedReplay(ctx context.Context, knot string) int64 { + repos, err := s.db.AllRepos() + if err != nil { + s.l.Warn("failed to load repos after outdated cursor", "knot", knot, "err", err) + return 0 + } + reset := 0 + for _, repo := range repos { + if repo.Knot != knot || repo.RepoDid == "" { + continue + } + if err := s.wipeFeedRefs(ctx, repo.RepoDid); err != nil { + s.l.Error("failed to reset ref state, old shas may be stale until the next push", "repo", repo.RepoDid, "err", err) + continue + } + s.refStateSeeded.Delete(repo.RepoDid.String()) + reset++ + } + if reset > 0 { + s.l.Warn("knot cannot replay from our cursor, ref state will be re-seeded live", "knot", knot, "repos", reset) + } + return 0 +} + +func (s *Spindle) wipeFeedRefs(ctx context.Context, repoDid syntax.DID) error { + var err error + for attempt := range feedWipeAttempts { + if attempt > 0 { + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(feedWipeRetryWait << attempt): + } + } + err = s.db.DeleteFeedRefsByRepoDid(repoDid) + if err == nil { + return nil + } + } + return err +} + +func (s *Spindle) collaboratorList(ctx context.Context, knot string, repoDid syntax.DID) ([]syntax.DID, error) { + if s.listCollaborators != nil { + return s.listCollaborators(ctx, knot, repoDid) + } + + base, err := knotEndpoint(knot, "/xrpc/sh.tangled.repo.listCollaborators") + if err != nil { + return nil, err + } + var subjects []syntax.DID + cursor := "" + for page := 0; ; page++ { + if page >= maxListPages { + return nil, fmt.Errorf("collaborators for %s exceed %d pages", repoDid, maxListPages) + } + q := url.Values{} + q.Set("subject", repoDid.String()) + q.Set("limit", "1000") + if cursor != "" { + q.Set("cursor", cursor) + } + var out struct { + Items []struct { + Subject string `json:"subject"` + } `json:"items"` + Cursor string `json:"cursor"` + } + if err := getJSON(ctx, base+"?"+q.Encode(), &out); err != nil { + return nil, err + } + for _, item := range out.Items { + did, err := syntax.ParseDID(item.Subject) + if err != nil { + return nil, fmt.Errorf("parsing collaborator subject %s: %w", item.Subject, err) + } + subjects = append(subjects, did) + } + if out.Cursor == "" { + return subjects, nil + } + cursor = out.Cursor + } +} + +func (s *Spindle) refRecords(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) { + if s.listRefRecords != nil { + return s.listRefRecords(ctx, knot, repoDid) + } + + base, err := knotEndpoint(knot, "/xrpc/com.atproto.repo.listRecords") + if err != nil { + return nil, err + } + var refs []refRecord + cursor := "" + pages := 0 + for { + pages++ + if pages > maxListPages { + return nil, fmt.Errorf("ref records for %s exceed %d pages", repoDid, maxListPages) + } + q := url.Values{} + q.Set("repo", repoDid.String()) + q.Set("collection", knotfeed.GitRefCollection) + q.Set("limit", "100") + if cursor != "" { + q.Set("cursor", cursor) + } + var out struct { + Records []struct { + Uri string `json:"uri"` + Value struct { + Sha string `json:"sha"` + } `json:"value"` + } `json:"records"` + Cursor string `json:"cursor"` + } + if err := getJSON(ctx, base+"?"+q.Encode(), &out); err != nil { + return nil, err + } + for _, rec := range out.Records { + uri, err := syntax.ParseATURI(rec.Uri) + if err != nil { + return nil, fmt.Errorf("parsing ref record uri %s: %w", rec.Uri, err) + } + refs = append(refs, refRecord{Rkey: uri.RecordKey().String(), Sha: rec.Value.Sha}) + } + if out.Cursor == "" { + return refs, nil + } + cursor = out.Cursor + } +} + +func knotEndpoint(knot, path string) (string, error) { + host, noTLS, err := hostutil.ParseHostname(knot) + if err != nil { + return "", fmt.Errorf("parsing knot host %s: %w", knot, err) + } + scheme := "https" + if noTLS { + scheme = "http" + } + return fmt.Sprintf("%s://%s%s", scheme, host, path), nil +} + +func getJSON(ctx context.Context, u string, out any) error { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) + if err != nil { + return err + } + resp, err := knotHTTPClient.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return fmt.Errorf("%s %s answered %s: %s", http.MethodGet, u, resp.Status, body) + } + return json.NewDecoder(resp.Body).Decode(out) +} diff --git a/spindle/knotfeed_test.go b/spindle/knotfeed_test.go new file mode 100644 index 000000000..fc92e9add --- /dev/null +++ b/spindle/knotfeed_test.go @@ -0,0 +1,393 @@ +package spindle + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "os" + "os/exec" + "path/filepath" + "slices" + "strings" + "testing" + + "github.com/bluesky-social/indigo/atproto/syntax" + cbg "github.com/whyrusleeping/cbor-gen" + "tangled.org/core/knotfeed" + "tangled.org/core/rbac" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/webhook" +) + +const ( + testKnot = "knot.nel.pet" + testForeign = "barnacle.nel.pet" + testRepoDid = syntax.DID("did:plc:limpet") + testSubject = syntax.DID("did:plc:boltless") + testRkeyMain = "refs~2fheads~2fmain" + testShaOld = "1111111111111111111111111111111111111111" + testShaNew = "2222222222222222222222222222222222222222" +) + +func newTestFeedSpindle(t *testing.T) *Spindle { + t.Helper() + d, e := newTestSpindleDB(t) + quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) + s := &Spindle{db: d, e: e, l: quiet, wh: webhook.New(d, false)} + if err := d.AddRepo(db.Repo{Knot: testKnot, Owner: "did:plc:akshay", Rkey: "3kqrstuvwxyz", RepoDid: testRepoDid}); err != nil { + t.Fatalf("AddRepo: %v", err) + } + return s +} + +func feedCommit(repo string, ops ...knotfeed.RecordOp) knotfeed.Message { + return knotfeed.Message{ + Type: knotfeed.TypeCommit, + Commit: &knotfeed.Commit{ + Repo: repo, + Seq: 42, + Records: ops, + }, + } +} + +func permissionsOf(t *testing.T, s *Spindle, repoDid syntax.DID) []string { + t.Helper() + perms := slices.Sorted(slices.Values(s.e.GetPermissionsInRepo(testSubject.String(), rbac.ThisServer, repoDid.String()))) + return perms +} + +func TestKnotFeedReconcileCollaborators(t *testing.T) { + cases := []struct { + name string + desired []syntax.DID + seed bool + granted bool + }{ + {"a listed collaborator is granted", []syntax.DID{testSubject}, false, true}, + {"a tracked collaborator absent from the knot is revoked", nil, true, false}, + {"a collaborator the knot still lists keeps the grant", []syntax.DID{testSubject}, true, true}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + s := newTestFeedSpindle(t) + if tc.seed { + if err := s.e.AddCollaborator(testSubject.String(), rbac.ThisServer, testRepoDid.String()); err != nil { + t.Fatalf("seed policy: %v", err) + } + if err := s.db.AddKnotCollaborator(testRepoDid, testSubject); err != nil { + t.Fatalf("seed row: %v", err) + } + } + s.listCollaborators = func(ctx context.Context, knot string, repoDid syntax.DID) ([]syntax.DID, error) { + return tc.desired, nil + } + + if err := s.reconcileCollaborators(context.Background(), testKnot, testRepoDid); err != nil { + t.Fatalf("reconcileCollaborators: %v", err) + } + + want := []string(nil) + if tc.granted { + want = []string{"repo:collaborator", "repo:push", "repo:settings"} + } + if perms := permissionsOf(t, s, testRepoDid); !slices.Equal(perms, want) { + t.Errorf("permissions = %v, want %v", perms, want) + } + held := slices.ContainsFunc(mustCollaborators(t, s), func(c db.RepoCollaborator) bool { return c.Subject == testSubject }) + if held != tc.granted { + t.Errorf("repo_collaborators row present = %v, want %v", held, tc.granted) + } + }) + } +} + +func mustCollaborators(t *testing.T, s *Spindle) []db.RepoCollaborator { + t.Helper() + collabs, err := s.db.ListCollaboratorsByRepoDid(testRepoDid) + if err != nil { + t.Fatalf("ListCollaboratorsByRepoDid: %v", err) + } + return collabs +} + +func TestKnotFeedReconcileIgnoresForeignKnotAndUnknownRepo(t *testing.T) { + s := newTestFeedSpindle(t) + var asked int + s.listCollaborators = func(ctx context.Context, knot string, repoDid syntax.DID) ([]syntax.DID, error) { + asked++ + return []syntax.DID{testSubject}, nil + } + + if err := s.reconcileCollaborators(context.Background(), testForeign, testRepoDid); err != nil { + t.Fatalf("foreign knot reconcile: %v", err) + } + if err := s.reconcileCollaborators(context.Background(), testKnot, syntax.DID("did:plc:unknown")); err != nil { + t.Fatalf("unknown repo reconcile: %v", err) + } + if asked != 0 { + t.Errorf("listCollaborators asked %d times, want 0", asked) + } + if perms := permissionsOf(t, s, testRepoDid); len(perms) != 0 { + t.Errorf("permissions = %v, want none", perms) + } +} + +func TestKnotFeedRefOps(t *testing.T) { + t.Run("a delete op drops the ref and triggers nothing", func(t *testing.T) { + s := newTestFeedSpindle(t) + s.listRefRecords = func(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) { + return []refRecord{{Rkey: testRkeyMain, Sha: testShaOld}}, nil + } + if err := s.db.PutFeedRef(testRepoDid, testRkeyMain, testShaOld); err != nil { + t.Fatalf("seed ref state: %v", err) + } + + if err := s.handleKnotFeed(context.Background(), testKnot, feedCommit(testRepoDid.String(), + knotfeed.RecordOp{Action: "delete", Collection: knotfeed.GitRefCollection, Rkey: testRkeyMain}, + )); err != nil { + t.Fatalf("handleKnotFeed: %v", err) + } + if _, ok, err := s.db.FeedRefSha(testRepoDid, testRkeyMain); err != nil || ok { + t.Errorf("ref state row present = %v (err %v), want gone", ok, err) + } + }) + + t.Run("a skip-ci push updates ref state and triggers nothing", func(t *testing.T) { + s := newTestFeedSpindle(t) + s.listRefRecords = func(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) { + return nil, nil + } + if err := s.db.PutFeedRef(testRepoDid, testRkeyMain, testShaOld); err != nil { + t.Fatalf("seed ref state: %v", err) + } + + if err := s.handleKnotFeed(context.Background(), testKnot, feedCommit(testRepoDid.String(), + knotfeed.RecordOp{Action: "create", Collection: knotfeed.GitRefCollection, Rkey: testRkeyMain, Bytes: encodeRefRecord(testShaNew, "skip-ci")}, + )); err != nil { + t.Fatalf("handleKnotFeed: %v", err) + } + if sha, _, err := s.db.FeedRefSha(testRepoDid, testRkeyMain); err != nil || sha != testShaNew { + t.Errorf("ref state sha = %q (err %v), want %q", sha, err, testShaNew) + } + }) + + t.Run("a create op that never materializes skips bootstrap", func(t *testing.T) { + for _, tc := range []struct { + name string + repo string + rkey string + }{ + {"unknown repo", "did:plc:unknown", testRkeyMain}, + {"ref outside heads and tags", testRepoDid.String(), "refs~2fnotes~2fwip"}, + } { + t.Run(tc.name, func(t *testing.T) { + s := newTestFeedSpindle(t) + s.listRefRecords = func(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) { + t.Errorf("bootstrap ran for %s", tc.name) + return nil, nil + } + + if err := s.handleKnotFeed(context.Background(), testKnot, feedCommit(tc.repo, + knotfeed.RecordOp{Action: "create", Collection: knotfeed.GitRefCollection, Rkey: tc.rkey, Bytes: encodeRefRecord(testShaNew)}, + )); err != nil { + t.Fatalf("handleKnotFeed: %v", err) + } + }) + } + }) + + t.Run("an undecodable record is an error rather than a silent drop", func(t *testing.T) { + s := newTestFeedSpindle(t) + s.listRefRecords = func(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) { + return nil, nil + } + + if err := s.handleKnotFeed(context.Background(), testKnot, feedCommit(testRepoDid.String(), + knotfeed.RecordOp{Action: "create", Collection: knotfeed.GitRefCollection, Rkey: testRkeyMain, Bytes: []byte{0xff}}, + )); err == nil { + t.Fatal("expected an error for an undecodable ref record") + } + }) + + t.Run("a collaboratorInvite op reconciles the repo", func(t *testing.T) { + s := newTestFeedSpindle(t) + var reconciled bool + s.listCollaborators = func(ctx context.Context, knot string, repoDid syntax.DID) ([]syntax.DID, error) { + reconciled = true + return nil, nil + } + + if err := s.handleKnotFeed(context.Background(), testKnot, feedCommit(testRepoDid.String(), + knotfeed.RecordOp{Action: "create", Collection: collaboratorInviteNSID, Rkey: testSubject.String()}, + )); err != nil { + t.Fatalf("handleKnotFeed: %v", err) + } + if !reconciled { + t.Error("collaboratorInvite didn't run a reconcile") + } + }) +} + +func TestKnotFeedSeedsRefStateOnce(t *testing.T) { + s := newTestFeedSpindle(t) + rkeyTag := "refs~2ftags~2fv1" + fetched := 0 + s.listRefRecords = func(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) { + if fetched++; fetched == 1 { + return nil, errors.New("knot unreachable") + } + return []refRecord{{Rkey: testRkeyMain, Sha: testShaOld}, {Rkey: rkeyTag, Sha: testShaNew}}, nil + } + + if err := s.ensureRefState(context.Background(), testKnot, testRepoDid); err == nil { + t.Fatal("expected the first bootstrap attempt to surface the fetch error") + } + if err := s.ensureRefState(context.Background(), testKnot, testRepoDid); err != nil { + t.Fatalf("ensureRefState: %v", err) + } + for _, want := range []struct{ rkey, sha string }{ + {testRkeyMain, testShaOld}, + {rkeyTag, testShaNew}, + } { + if sha, ok, err := s.db.FeedRefSha(testRepoDid, want.rkey); err != nil || !ok || sha != want.sha { + t.Errorf("ref state (%s, %s) = (%q, %v, err %v)", want.rkey, want.sha, sha, ok, err) + } + } + if err := s.ensureRefState(context.Background(), testKnot, testRepoDid); err != nil || fetched != 2 { + t.Fatalf("re-bootstrap fetched %d times (err %v), want 2 fetches and no error", fetched, err) + } +} + +func TestAdmitChangedFilesUnderBudget(t *testing.T) { + shaSized := func(n int) []string { + paths := make([]string, n) + for i := range n { + paths[i] = "src/f.go" + } + return paths + } + + if got := admitChangedFiles([]string{"a.go", "b/c.go"}); !slices.Equal(got, []string{"a.go", "b/c.go"}) { + t.Errorf("small listing = %v, want both paths", got) + } + if got := admitChangedFiles(shaSized(11000)); got != nil { + t.Errorf("listing over the entry cap = %v, want nil", got) + } + wide := shaSized(8192) + if got := admitChangedFiles(wide); !slices.Equal(got, wide) { + t.Errorf("listing at the entry cap = %d paths, want all %d", len(got), len(wide)) + } + big := make([]string, 11000) + for i := range big { + big[i] = string(make([]byte, 64)) + } + if got := admitChangedFiles(big); got != nil { + t.Errorf("listing over the byte budget = %d paths, want nil", len(got)) + } +} + +func TestChangedFilesUnderBudgetGuardsFirstSight(t *testing.T) { + quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) + if got := changedFilesUnderBudget(quiet, "/nonexistent/repo", "", "2222222222222222222222222222222222222222"); got != nil { + t.Errorf("first sight of a ref = %v, want nil without touching the repo", got) + } +} + +func encodeRefRecord(sha string, options ...string) []byte { + var out bytes.Buffer + writeText := func(s string) { + cbg.CborWriteHeader(&out, cbg.MajTextString, uint64(len(s))) + out.WriteString(s) + } + entries := uint64(1) + if len(options) > 0 { + entries = 2 + } + cbg.CborWriteHeader(&out, cbg.MajMap, entries) + writeText("sha") + writeText(sha) + if len(options) > 0 { + writeText("x-tngl-push-options") + cbg.CborWriteHeader(&out, cbg.MajArray, uint64(len(options))) + for _, option := range options { + writeText(option) + } + } + return out.Bytes() +} + +func TestChangedFilesUnderBudgetWithRealRepo(t *testing.T) { + if _, err := exec.LookPath("git"); err != nil { + t.Skip("git not available") + } + repoPath := filepath.Join(t.TempDir(), "repo") + run := func(args ...string) string { + out, err := exec.Command("git", args...).CombinedOutput() + if err != nil { + t.Fatalf("git %v: %v\n%s", args, err, out) + } + return strings.TrimSpace(string(out)) + } + run("init", "-q", "-b", "main", repoPath) + run("-C", repoPath, "config", "user.email", "ci@nel.pet") + run("-C", repoPath, "config", "user.name", "ci") + write := func(name, content string) { + path := filepath.Join(repoPath, name) + os.MkdirAll(filepath.Dir(path), 0o755) + if err := os.WriteFile(path, []byte(content), 0o644); err != nil { + t.Fatalf("write %s: %v", name, err) + } + } + commit := func() string { + run("-C", repoPath, "add", "-A") + run("-C", repoPath, "commit", "-q", "-m", "wip") + return run("-C", repoPath, "rev-parse", "HEAD") + } + + write("README.md", "one\n") + shaFirst := commit() + write(".tangled/workflows/ci.yaml", "when:\n event: push\n") + write("src/app.go", "package main\n") + shaSecond := commit() + run("-C", repoPath, "mv", "src/app.go", "src/main.go") + shaThird := commit() + + quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) + for _, tc := range []struct { + name string + oldSha string + newSha string + want []string + }{ + {"a push commit lists its added and modified files", shaFirst, shaSecond, []string{".tangled/workflows/ci.yaml", "src/app.go"}}, + {"a rename lists both the old and new path", shaSecond, shaThird, []string{"src/app.go", "src/main.go"}}, + {"a no-op revision range lists nothing", shaThird, shaThird, nil}, + } { + t.Run(tc.name, func(t *testing.T) { + got := changedFilesUnderBudget(quiet, repoPath, tc.oldSha, tc.newSha) + slices.Sort(got) + if !slices.Equal(got, tc.want) { + t.Errorf("changed files = %v, want %v", got, tc.want) + } + }) + } +} + +func TestCollaboratorListBailsPastMaxListPages(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + json.NewEncoder(w).Encode(map[string]any{"items": []any{}, "cursor": "more"}) + })) + defer srv.Close() + + s := &Spindle{} + host := "localhost:" + strings.TrimPrefix(srv.Listener.Addr().String(), "127.0.0.1:") + if _, err := s.collaboratorList(context.Background(), host, testRepoDid); err == nil { + t.Fatal("collaboratorList answered no error past the page cap") + } +} diff --git a/spindle/knotstream_test.go b/spindle/knotstream_test.go deleted file mode 100644 index 9beab00f6..000000000 --- a/spindle/knotstream_test.go +++ /dev/null @@ -1,70 +0,0 @@ -package spindle - -import ( - "context" - "encoding/json" - "errors" - "io" - "log/slog" - "slices" - "testing" - - "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/samber/lo" - "tangled.org/core/eventconsumer" - "tangled.org/core/eventstream" - knotdb "tangled.org/core/knotserver/db" - "tangled.org/core/log" - "tangled.org/core/rbac" - "tangled.org/core/spindle/db" -) - -func TestKnotStreamTracksCollaboratorsOnlyFromTheKnotHostingTheRepo(t *testing.T) { - const knot, foreignKnot = "knot.nel.pet", "barnacle.nel.pet" - const repoDid, subject = syntax.DID("did:plc:limpet"), syntax.DID("did:plc:boltless") - type step struct { - host string - op knotdb.AclOp - member bool - } - event := func(st step) eventstream.Event { - record := lo.Ternary[any](st.member, knotdb.KnotMemberUpdate{Op: st.op, Subject: subject.String()}, knotdb.RepoCollaboratorUpdate{Op: st.op, Subject: subject.String(), Repo: repoDid.String()}) - return eventstream.Event{Rkey: "evt", Nsid: lo.Ternary(st.member, knotdb.KnotMemberUpdateNSID, knotdb.RepoCollaboratorUpdateNSID), EventJson: lo.Must(json.Marshal(record))} - } - - for _, tc := range []struct { - name string - steps []step - wantErr, granted bool - }{ - {"an add grants the policy and writes the row", []step{{knot, knotdb.AclOpAdd, false}}, false, true}, - {"a remove revokes the policy and drops the row", []step{{knot, knotdb.AclOpAdd, false}, {knot, knotdb.AclOpRemove, false}}, false, false}, - {"an add from a knot that doesn't host the repo is dropped quietly", []step{{foreignKnot, knotdb.AclOpAdd, false}}, false, false}, - {"a remove from a knot that doesn't host the repo can't revoke the grant", []step{{knot, knotdb.AclOpAdd, false}, {foreignKnot, knotdb.AclOpRemove, false}}, false, true}, - {"an unrecognized op is an error rather than a silent drop", []step{{knot, knotdb.AclOp("bogus"), false}}, true, false}, - {"a knot memberUpdate is ignored because knot membership is not repo access", []step{{knot, knotdb.AclOpAdd, true}}, false, false}, - } { - t.Run(tc.name, func(t *testing.T) { - d, e := newTestSpindleDB(t) - quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) - ctx, s := log.IntoContext(context.Background(), quiet), &Spindle{db: d, e: e, l: quiet} - lo.Must0(d.AddRepo(db.Repo{Knot: knot, Owner: "did:plc:akshay", Rkey: "3kqrstuvwxyz", RepoDid: repoDid})) - - err := lo.Reduce(tc.steps, func(acc error, st step, _ int) error { - return errors.Join(acc, s.processKnotStream(ctx, eventconsumer.Source{Kind: eventconsumer.KindKnot, Host: st.host}, event(st))) - }, nil) - if (err != nil) != tc.wantErr { - t.Fatalf("processKnotStream err = %v, want an error: %v", err, tc.wantErr) - } - - wantPerms := lo.Ternary(tc.granted, []string{"repo:collaborator", "repo:push", "repo:settings"}, nil) - perms := slices.Sorted(slices.Values(s.e.GetPermissionsInRepo(subject.String(), rbac.ThisServer, repoDid.String()))) - if !slices.Equal(perms, wantPerms) { - t.Errorf("permissions = %v, want %v", perms, wantPerms) - } - if held := slices.ContainsFunc(lo.Must(s.db.ListCollaboratorsByRepoDid(repoDid)), func(c db.RepoCollaborator) bool { return c.Subject == subject }); held != tc.granted { - t.Errorf("repo_collaborators row present = %v, want %v", held, tc.granted) - } - }) - } -} diff --git a/spindle/server.go b/spindle/server.go index e42275ead..fc2408193 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -3,10 +3,8 @@ package spindle import ( "context" "crypto/sha256" - "database/sql" _ "embed" "encoding/binary" - "encoding/json" "errors" "fmt" "io" @@ -27,13 +25,9 @@ import ( "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "tangled.org/core/api/tangled" - "tangled.org/core/eventconsumer" - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/eventstream" "tangled.org/core/gitutil" "tangled.org/core/idresolver" "tangled.org/core/jetstream" - knotdb "tangled.org/core/knotserver/db" kgit "tangled.org/core/knotserver/git" "tangled.org/core/log" "tangled.org/core/notifier" @@ -46,6 +40,7 @@ import ( "tangled.org/core/spindle/engine" "tangled.org/core/spindle/engines/dummy" "tangled.org/core/spindle/engines/nixery" + "tangled.org/core/spindle/feed" "tangled.org/core/spindle/mill" "tangled.org/core/spindle/mill/executor" "tangled.org/core/spindle/models" @@ -80,35 +75,38 @@ type executorClient interface { } type Spindle struct { - jc *jetstream.JetstreamClient - tap *Tap - embedTap *embeddedTap - db *db.DB - e *rbac.Enforcer - l *slog.Logger - n *notifier.Notifier - wh *webhook.Service - metrics *observability.Metrics - engs map[string]models.Engine - jobWake chan struct{} - jobWorkers sync.WaitGroup - cfg *config.Config - ks *eventconsumer.Consumer - res *idresolver.Resolver - verify repoverify.Verifier - vault secrets.Manager - cache storage.Storage - motd []byte - motdMu sync.RWMutex - rootCtx context.Context - rootCancel context.CancelFunc - store artifactstore.Store - stores *artifactstore.Stores - reader artifactstore.Reader - qm *quota.Manager + jc *jetstream.JetstreamClient + tap *Tap + embedTap *embeddedTap + db *db.DB + e *rbac.Enforcer + l *slog.Logger + n *notifier.Notifier + wh *webhook.Service + metrics *observability.Metrics + engs map[string]models.Engine + jobWake chan struct{} + jobWorkers sync.WaitGroup + cfg *config.Config + feed *feed.Feed + refStateSeeded sync.Map + res *idresolver.Resolver + verify repoverify.Verifier + vault secrets.Manager + cache storage.Storage + motd []byte + motdMu sync.RWMutex + rootCtx context.Context + rootCancel context.CancelFunc + store artifactstore.Store + stores *artifactstore.Stores + reader artifactstore.Reader + qm *quota.Manager // set only when this spindle hosts the mill or joins one as an executor - mill *mill.Mill - exec executorClient + mill *mill.Mill + exec executorClient + listCollaborators func(ctx context.Context, knot string, repoDid syntax.DID) ([]syntax.DID, error) + listRefRecords func(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) } func newCacheStore(ctx context.Context, cfg *config.Config) (storage.Storage, error) { @@ -299,41 +297,24 @@ func New(ctx context.Context, cfg *config.Config, d *db.DB, engines map[string]m } logger.Info("owner set", "did", cfg.Server.Owner) - cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath) - if err != nil { - return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err) - } - err = jc.StartJetstream(lifecycleCtx, spindle.ingest()) if err != nil { return nil, fmt.Errorf("failed to start jetstream consumer: %w", err) } + spindle.feed = feed.New(log.SubLogger(logger, "knotfeed"), feed.Hooks{ + LoadCursor: func(ctx context.Context, knot string) (int64, error) { + return spindle.db.LoadFeedCursor(knot) + }, + StoreCursor: func(ctx context.Context, knot string, seq int64) error { + return spindle.db.StoreFeedCursor(knot, seq) + }, + Handle: spindle.handleKnotFeed, + OutdatedReplay: spindle.feedOutdatedReplay, + OnConnectError: func(knot string, err error) { + logger.Warn("cannot reach knot firehose", "knot", knot, "err", err) + }, + }) - // spindle listen to knot stream for sh.tangled.git.refUpdate - // which will sync the local workflow files in spindle and enqueues the - // pipeline job for on-push workflows - ccfg := eventconsumer.NewConsumerConfig() - ccfg.Logger = log.SubLogger(logger, "eventconsumer") - ccfg.ProcessFunc = spindle.processKnotStream - ccfg.CursorStore = cursorStore - if cfg.Server.Dev { - ccfg.RetryInterval = 5 * time.Second - ccfg.MaxRetryInterval = 10 * time.Second - } else { - ccfg.RetryInterval = 1 * time.Minute - ccfg.MaxRetryInterval = 10 * time.Minute - } - knownKnots, err := d.Knots() - if err != nil { - return nil, err - } - for _, knot := range knownKnots { - logger.Info("adding source start", "knot", knot) - src := eventconsumer.NewKnotSource(knot) - eventconsumer.MigrateLegacyCursor(cursorStore, src) - ccfg.Sources[src] = struct{}{} - } - spindle.ks = eventconsumer.NewConsumer(*ccfg) if cfg.Server.Tap.Embed { pw, err := randomAdminPassword() if err != nil { @@ -449,9 +430,15 @@ func (s *Spindle) Start(ctx context.Context) error { s.resumeWipes(ctx) go func() { - s.l.Info("starting knot event consumer") - s.ks.Start(runCtx) + s.l.Info("starting knot firehose feed") + knots, err := s.db.Knots() + if err != nil { + s.l.Error("listing known knots for the feed", "err", err) + return + } + s.feed.Start(runCtx, knots) }() + go s.reconcileCollaboratorsLoop(ctx) s.l.Info("starting tap client", "url", s.cfg.Server.Tap.Url) s.tap.Start(tapCtx) @@ -821,154 +808,6 @@ func (s *Spindle) XrpcRouter() http.Handler { return x.Router() } -func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) error { - ctx, span := observability.Tracer().Start(ctx, "knot.ingest") - defer span.End() - - metrics := s.metrics - err := s.processKnotStreamInner(ctx, src, msg) - if err != nil { - metrics.RecordEventIngestion("knot", "error") - span.SetStatus(codes.Error, "failed to process knot stream event") - } else { - metrics.RecordEventIngestion("knot", "success") - span.SetStatus(codes.Ok, "success") - } - return err -} - -func (s *Spindle) processKnotStreamInner(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) 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 == knotdb.RepoCollaboratorUpdateNSID { - return s.ingestKnotCollaborator(ctx, l, src, msg) - } - 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("repo", event.Repo, "ref", event.Ref, "newSha", event.NewSha) - l.Debug("debug") - - repoDid := syntax.DID(event.Repo) - repo, err := s.db.GetRepoByDid(repoDid) - if err != nil { - return fmt.Errorf("unknown repoDid %s: %w", repoDid, err) - } - - if ban, err := s.db.IsBanned(repoDid, repo.Owner); err != nil { - return fmt.Errorf("checking bans: %w", err) - } else if ban != nil { - l.Warn("dropping push for banned repo") - return nil - } - - if src.Host != repo.Knot { - return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, repo.Knot) - } - - // deliver push webhooks independently of CI; a skip-ci push option - // should not suppress webhook notifications - s.wh.FirePush(ctx, repo, event.CommitterDid, event.Ref, event.OldSha, event.NewSha) - - if kgit.HasSkipCIPushOption(event.PushOptions) { - l.Info("push event requested ci skip, skipping the event") - return nil - } - - // NOTE: we are blindly trusting the knot that it will return only repos it own - repoCloneUri := s.newRepoCloneUrl(src.Host, repoDid) - repoPath := s.newRepoPath(repoDid) - if err := gitutil.SparseSync(ctx, repoCloneUri, repoPath, event.NewSha, sparseWorkflowDir); err != nil { - return fmt.Errorf("sync git repo: %w", err) - } - l.Info("synced git repo") - - triggerRepo, err := s.buildTriggerRepo(ctx, repo) - if err != nil { - return fmt.Errorf("building trigger repo: %w", err) - } - - trigger := tangled.Pipeline_TriggerMetadata{ - Kind: string(workflow.TriggerKindPush), - Push: &tangled.Pipeline_PushTriggerData{ - Ref: event.Ref, - OldSha: event.OldSha, - NewSha: event.NewSha, - }, - Repo: triggerRepo, - } - - pipelineId, err := s.runPipeline(ctx, repoDid, trigger, event.ChangedFiles, repoCloneUri, repoPath, event.NewSha, nil, triggerRepo) - if err != nil { - return err - } - if pipelineId.Rkey == "" { - l.Info("no workflow matched 'push' trigger, skipping the event") - return nil - } - l.Info("pipeline triggered", "pipeline", pipelineId.AtUri()) - } - - return nil -} - -func (s *Spindle) ingestKnotCollaborator(ctx context.Context, l *slog.Logger, src eventconsumer.Source, msg eventstream.Event) error { - var rec knotdb.RepoCollaboratorUpdate - if err := json.Unmarshal(msg.EventJson, &rec); err != nil { - l.Error("error unmarshalling collaboratorUpdate", "err", err) - return err - } - - subject, err := syntax.ParseDID(rec.Subject) - if err != nil { - l.Info("skipping collaboratorUpdate with malformed subject", "subject", rec.Subject, "err", err) - return nil - } - repoDid, err := syntax.ParseDID(rec.Repo) - if err != nil { - l.Info("skipping collaboratorUpdate with malformed repo", "repo", rec.Repo, "err", err) - return nil - } - - repo, err := s.db.GetRepoByDid(repoDid) - if errors.Is(err, sql.ErrNoRows) { - l.Info("skipping collaboratorUpdate for unknown repo", "repo", repoDid) - return nil - } - if err != nil { - return fmt.Errorf("lookup repo %s: %w", repoDid, err) - } - if src.Host != repo.Knot { - l.Warn("dropping collaboratorUpdate from non-owning knot", "src", src.Host, "repoKnot", repo.Knot) - return nil - } - - switch rec.Op { - case knotdb.AclOpAdd: - if err := s.e.AddCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { - return fmt.Errorf("add collaborator policy: %w", err) - } - if err := s.db.AddKnotCollaborator(repoDid, subject); err != nil { - return fmt.Errorf("track collaborator: %w", err) - } - l.Info("added knot-managed collaborator", "subject", subject, "repo", repoDid) - case knotdb.AclOpRemove: - if err := s.e.RemoveCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { - return fmt.Errorf("remove collaborator policy: %w", err) - } - if err := s.db.DeleteRepoCollaboratorBySubjectRepo(subject, repoDid); err != nil { - return fmt.Errorf("delete collaborator row: %w", err) - } - l.Info("removed knot-managed collaborator", "subject", subject, "repo", repoDid) - default: - return fmt.Errorf("collaboratorUpdate unknown op %q", rec.Op) - } - return nil -} - // buildTriggerRepo gathers trigger metadata, resolving default branch from the knot func (s *Spindle) buildTriggerRepo(ctx context.Context, repo *db.Repo) (*tangled.Pipeline_TriggerRepo, error) { rkey := string(repo.Rkey) diff --git a/spindle/server.go.master b/spindle/server.go.master deleted file mode 100644 index 7de463428..000000000 --- a/spindle/server.go.master +++ /dev/null @@ -1,930 +0,0 @@ -package spindle - -import ( - "context" - "database/sql" - _ "embed" - "encoding/json" - "errors" - "fmt" - "log/slog" - "maps" - "net/http" - "path/filepath" - "sync" - "time" - - "github.com/bluesky-social/indigo/atproto/syntax" - indigoxrpc "github.com/bluesky-social/indigo/xrpc" - "github.com/go-chi/chi/v5" - "github.com/go-git/go-git/v5/plumbing/object" - "github.com/hashicorp/go-version" - "tangled.org/core/api/tangled" - "tangled.org/core/eventconsumer" - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/eventstream" - "tangled.org/core/idresolver" - "tangled.org/core/jetstream" - knotdb "tangled.org/core/knotserver/db" - kgit "tangled.org/core/knotserver/git" - "tangled.org/core/log" - "tangled.org/core/notifier" - "tangled.org/core/rbac" - "tangled.org/core/repoident" - "tangled.org/core/repoverify" - "tangled.org/core/spindle/config" - "tangled.org/core/spindle/db" - "tangled.org/core/spindle/engine" - "tangled.org/core/spindle/engines/dummy" - "tangled.org/core/spindle/engines/nixery" - "tangled.org/core/spindle/git" - "tangled.org/core/spindle/models" - "tangled.org/core/spindle/secrets" - "tangled.org/core/spindle/xrpc" - "tangled.org/core/tid" - "tangled.org/core/workflow" - "tangled.org/core/xrpc/serviceauth" -) - -//go:embed motd -var defaultMotd []byte - -const ( - rbacDomain = "thisserver" -) - -type Spindle struct { - jc *jetstream.JetstreamClient - tap *Tap - embedTap *embeddedTap - db *db.DB - e *rbac.Enforcer - l *slog.Logger - n *notifier.Notifier - engs map[string]models.Engine - cfg *config.Config - ks *eventconsumer.Consumer - res *idresolver.Resolver - verify repoverify.Verifier - vault secrets.Manager - motd []byte - motdMu sync.RWMutex - rootCtx context.Context - jobWake chan struct{} -} - -// New creates a new Spindle server with the provided configuration and engines. -func New(ctx context.Context, cfg *config.Config, d *db.DB, engines map[string]models.Engine) (*Spindle, error) { - logger := log.FromContext(ctx) - - e, err := rbac.NewEnforcer(cfg.Server.DBPath) - if err != nil { - return nil, fmt.Errorf("failed to setup rbac enforcer: %w", err) - } - e.E.EnableAutoSave(true) - - n := notifier.New() - - var vault secrets.Manager - switch cfg.Server.Secrets.Provider { - case "openbao": - if cfg.Server.Secrets.OpenBao.ProxyAddr == "" { - return nil, fmt.Errorf("openbao proxy address is required when using openbao secrets provider") - } - vault, err = secrets.NewOpenBaoManager( - cfg.Server.Secrets.OpenBao.ProxyAddr, - logger, - secrets.WithMountPath(cfg.Server.Secrets.OpenBao.Mount), - ) - if err != nil { - return nil, fmt.Errorf("failed to setup openbao secrets provider: %w", err) - } - 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")) - 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) - default: - return nil, fmt.Errorf("unknown secrets provider: %s", cfg.Server.Secrets.Provider) - } - - if err := runStartupMigrations(ctx, d, cfg.Server.Tap.Embed, cfg.Server.Tap.DBPath, logger); err != nil { - return nil, fmt.Errorf("failed to run startup migrations: %w", err) - } - - collections := []string{ - tangled.SpindleMemberNSID, - tangled.RepoNSID, - tangled.RepoCollaboratorNSID, - tangled.RepoPullNSID, - } - jc, err := jetstream.NewJetstreamClient(cfg.Server.JetstreamEndpoint, "spindle", collections, nil, log.SubLogger(logger, "jetstream"), d, true, true) - if err != nil { - return nil, fmt.Errorf("failed to setup jetstream client: %w", err) - } - jc.AddDid(cfg.Server.Owner) - // pull records are created by arbitrary users too, same hack as in tap - jc.ExemptCollection(tangled.RepoPullNSID) - - // Check if the spindle knows about any Dids; - dids, err := d.GetAllDids() - if err != nil { - return nil, fmt.Errorf("failed to get all dids: %w", err) - } - for _, d := range dids { - jc.AddDid(d) - } - - knownRepos, err := d.AllRepos() - if err != nil { - return nil, fmt.Errorf("failed to get known repos: %w", err) - } - for _, r := range knownRepos { - if r.Owner != "" { - jc.AddDid(r.Owner.String()) - } - } - - resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl) - - spindle := &Spindle{ - jc: jc, - e: e, - db: d, - l: logger, - n: &n, - engs: engines, - cfg: cfg, - res: resolver, - verify: repoverify.New(resolver, cfg.Server.Dev), - vault: vault, - motd: defaultMotd, - rootCtx: ctx, - jobWake: make(chan struct{}, 1), - } - - err = e.AddSpindle(rbacDomain) - if err != nil { - return nil, fmt.Errorf("failed to set rbac domain: %w", err) - } - err = spindle.configureOwner() - if err != nil { - return nil, err - } - logger.Info("owner set", "did", cfg.Server.Owner) - - cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath) - if err != nil { - return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err) - } - - err = jc.StartJetstream(ctx, spindle.ingest()) - if err != nil { - return nil, fmt.Errorf("failed to start jetstream consumer: %w", err) - } - - // spindle listen to knot stream for sh.tangled.git.refUpdate - // which will sync the local workflow files in spindle and enqueues the - // pipeline job for on-push workflows - ccfg := eventconsumer.NewConsumerConfig() - ccfg.Logger = log.SubLogger(logger, "eventconsumer") - ccfg.ProcessFunc = spindle.processKnotStream - ccfg.CursorStore = cursorStore - ccfg.WorkerCount = 16 - ccfg.QueueSize = 200 - if cfg.Server.Dev { - ccfg.RetryInterval = 5 * time.Second - ccfg.MaxRetryInterval = 10 * time.Second - } else { - ccfg.RetryInterval = 1 * time.Minute - ccfg.MaxRetryInterval = 10 * time.Minute - } - knownKnots, err := d.Knots() - if err != nil { - return nil, err - } - for _, knot := range knownKnots { - logger.Info("adding source start", "knot", knot) - src := eventconsumer.NewKnotSource(knot) - eventconsumer.MigrateLegacyCursor(cursorStore, src) - ccfg.Sources[src] = struct{}{} - } - spindle.ks = eventconsumer.NewConsumer(*ccfg) - - if cfg.Server.Tap.Embed { - pw, err := randomAdminPassword() - if err != nil { - return nil, err - } - cfg.Server.Tap.AdminPassword = pw - logger.Info("embedded tap: using random admin password") - } - spindle.tap = NewTapClient(spindle) - - return spindle, nil -} - -// DB returns the database instance. -func (s *Spindle) DB() *db.DB { - return s.db -} - -// Engines returns the map of available engines. -func (s *Spindle) Engines() map[string]models.Engine { - return s.engs -} - -// Vault returns the secrets manager instance. -func (s *Spindle) Vault() secrets.Manager { - return s.vault -} - -// Notifier returns the notifier instance. -func (s *Spindle) Notifier() *notifier.Notifier { - return s.n -} - -// Enforcer returns the RBAC enforcer instance. -func (s *Spindle) Enforcer() *rbac.Enforcer { - return s.e -} - -// SetMotdContent sets custom MOTD content, replacing the embedded default. -func (s *Spindle) SetMotdContent(content []byte) { - s.motdMu.Lock() - defer s.motdMu.Unlock() - s.motd = content -} - -// GetMotdContent returns the current MOTD content. -func (s *Spindle) GetMotdContent() []byte { - s.motdMu.RLock() - defer s.motdMu.RUnlock() - return s.motd -} - -// Start starts the Spindle server (blocking). -func (s *Spindle) Start(ctx context.Context) error { - // starts a job queue runner in the background - s.StartJobWorkers(ctx) - - // Stop vault token renewal if it implements Stopper - if stopper, ok := s.vault.(secrets.Stopper); ok { - defer stopper.Stop() - } - - tapCtx, tapCancel := context.WithCancel(ctx) - - if s.cfg.Server.Tap.Embed { - emb, err := startEmbeddedTap(tapCtx, s.cfg, log.SubLogger(s.l, "embedtap")) - if err != nil { - tapCancel() - return fmt.Errorf("starting embedded tap: %w", err) - } - s.embedTap = emb - defer func() { - tapCancel() - s.embedTap.Shutdown() - }() - - go s.watchTapDrain(tapCtx, tapCancel) - } else { - defer tapCancel() - } - - go func() { - s.l.Info("starting knot event consumer") - s.ks.Start(ctx) - }() - - s.l.Info("starting tap client", "url", s.cfg.Server.Tap.Url) - s.tap.Start(tapCtx) - - s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr) - return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router()) -} - -func (s *Spindle) declareTapInterest(ctx context.Context) { - repos, err := s.db.AllRepos() - if err != nil { - s.l.Warn("tap declare: failed to load known repos", "err", err) - return - } - seen := make(map[syntax.DID]struct{}, len(repos)) - dids := make([]syntax.DID, 0, len(repos)) - for _, r := range repos { - if r.Owner == "" { - continue - } - if _, ok := seen[r.Owner]; ok { - continue - } - seen[r.Owner] = struct{}{} - dids = append(dids, r.Owner) - } - if err := s.tap.AddOwnerDIDs(ctx, dids); err != nil { - s.l.Warn("tap declare: AddRepos rejected", "count", len(dids), "err", err) - return - } - s.l.Info("tap declare: known owner DIDs registered", "count", len(dids)) -} - -func Run(ctx context.Context) error { - cfg, err := config.Load(ctx) - if err != nil { - return fmt.Errorf("failed to load config: %w", err) - } - - if err := ensureGitVersion(); err != nil { - return fmt.Errorf("ensuring git version: %w", err) - } - - d, err := db.Make(ctx, cfg.Server.DBPath) - if err != nil { - return fmt.Errorf("failed to setup db: %w", err) - } - - nixeryEng, err := nixery.New(ctx, cfg) - if err != nil { - return err - } - - microvmEng, err := newMicrovmEngine(ctx, cfg, d) - if err != nil { - return err - } - - s, err := New(ctx, cfg, d, map[string]models.Engine{ - "nixery": nixeryEng, - "microvm": microvmEng, - "dummy": dummy.New(log.FromContext(ctx)), - }) - if err != nil { - return err - } - - return s.Start(ctx) -} - -func (s *Spindle) Router() http.Handler { - mux := chi.NewRouter() - - mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { - w.Write(s.GetMotdContent()) - }) - mux.HandleFunc("/events", s.Events) - mux.HandleFunc("/logs/{knot}/{rkey}/{name}", s.Logs) - - mux.Mount("/xrpc", s.XrpcRouter()) - return mux -} - -func (s *Spindle) XrpcRouter() http.Handler { - serviceAuth := serviceauth.NewServiceAuth(s.l, s.res.Directory(), s.cfg.Server.Did().String()) - - l := log.SubLogger(s.l, "xrpc") - - x := xrpc.Xrpc{ - Logger: l, - Db: s.db, - Enforcer: s.e, - Engines: s.engs, - Config: s.cfg, - Resolver: s.res, - Vault: s.vault, - Notifier: s.Notifier(), - ServiceAuth: serviceAuth, - Trigger: s, - } - - return x.Router() -} - -func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) 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 == knotdb.RepoCollaboratorUpdateNSID { - return s.ingestKnotCollaborator(ctx, l, src, msg) - } - 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("repo", event.Repo, "ref", event.Ref, "newSha", event.NewSha) - l.Debug("debug") - - repoDid := syntax.DID(event.Repo) - repo, err := s.db.GetRepoByDid(repoDid) - if err != nil { - return fmt.Errorf("unknown repoDid %s: %w", repoDid, err) - } - - if src.Host != repo.Knot { - return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, repo.Knot) - } - - if kgit.HasSkipCIPushOption(event.PushOptions) { - l.Info("push event requested ci skip, skipping the event") - return nil - } - - // NOTE: we are blindly trusting the knot that it will return only repos it own - repoCloneUri := s.newRepoCloneUrl(src.Host, repoDid) - repoPath := s.newRepoPath(repoDid) - - triggerRepo, err := s.buildTriggerRepo(ctx, repo) - if err != nil { - return fmt.Errorf("building trigger repo: %w", err) - } - - trigger := tangled.Pipeline_TriggerMetadata{ - Kind: string(workflow.TriggerKindPush), - Push: &tangled.Pipeline_PushTriggerData{ - Ref: event.Ref, - OldSha: event.OldSha, - NewSha: event.NewSha, - }, - Repo: triggerRepo, - } - - pipelineId, err := s.runPipeline(ctx, repoDid, trigger, event.ChangedFiles, repoCloneUri, repoPath, event.NewSha, nil, triggerRepo) - if err != nil { - return err - } - if pipelineId.Rkey == "" { - l.Info("no workflow matched 'push' trigger, skipping the event") - return nil - } - l.Info("pipeline triggered", "pipeline", pipelineId.AtUri()) - } - - return nil -} - -func (s *Spindle) ingestKnotCollaborator(ctx context.Context, l *slog.Logger, src eventconsumer.Source, msg eventstream.Event) error { - var rec knotdb.RepoCollaboratorUpdate - if err := json.Unmarshal(msg.EventJson, &rec); err != nil { - l.Error("error unmarshalling collaboratorUpdate", "err", err) - return err - } - - subject, err := syntax.ParseDID(rec.Subject) - if err != nil { - l.Info("skipping collaboratorUpdate with malformed subject", "subject", rec.Subject, "err", err) - return nil - } - repoDid, err := syntax.ParseDID(rec.Repo) - if err != nil { - l.Info("skipping collaboratorUpdate with malformed repo", "repo", rec.Repo, "err", err) - return nil - } - - repo, err := s.db.GetRepoByDid(repoDid) - if errors.Is(err, sql.ErrNoRows) { - l.Info("skipping collaboratorUpdate for unknown repo", "repo", repoDid) - return nil - } - if err != nil { - return fmt.Errorf("lookup repo %s: %w", repoDid, err) - } - if src.Host != repo.Knot { - l.Warn("dropping collaboratorUpdate from non-owning knot", "src", src.Host, "repoKnot", repo.Knot) - return nil - } - - switch rec.Op { - case knotdb.AclOpAdd: - if err := s.e.AddCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { - return fmt.Errorf("add collaborator policy: %w", err) - } - if err := s.db.AddKnotCollaborator(repoDid, subject); err != nil { - return fmt.Errorf("track collaborator: %w", err) - } - l.Info("added knot-managed collaborator", "subject", subject, "repo", repoDid) - case knotdb.AclOpRemove: - if err := s.e.RemoveCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { - return fmt.Errorf("remove collaborator policy: %w", err) - } - if err := s.db.DeleteRepoCollaboratorBySubjectRepo(subject, repoDid); err != nil { - return fmt.Errorf("delete collaborator row: %w", err) - } - l.Info("removed knot-managed collaborator", "subject", subject, "repo", repoDid) - default: - return fmt.Errorf("collaboratorUpdate unknown op %q", rec.Op) - } - return nil -} - -// buildTriggerRepo gathers trigger metadata, resolving default branch from the knot -func (s *Spindle) buildTriggerRepo(ctx context.Context, repo *db.Repo) (*tangled.Pipeline_TriggerRepo, error) { - rkey := string(repo.Rkey) - repoDid := repo.RepoDid.String() - return s.buildTriggerRepoFrom(ctx, repo.Knot, repo.Owner.String(), rkey, repoDid), nil -} - -func (s *Spindle) buildTriggerRepoFrom(ctx context.Context, knot, did, rkey, repoDid string) *tangled.Pipeline_TriggerRepo { - scheme := "https" - if s.cfg.Server.Dev { - scheme = "http" - } - client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, knot)} - - // this should maybe (?) be in the refUpdate event itself to save a roundtrip - defaultBranch := "" - if out, err := tangled.RepoGetDefaultBranch(ctx, client, repoDid); err == nil { - defaultBranch = out.Name - } - - var rkeyPtr *string - if rkey != "" { - rkeyPtr = &rkey - } - return &tangled.Pipeline_TriggerRepo{ - Did: did, - Knot: knot, - Repo: rkeyPtr, - RepoDid: &repoDid, - DefaultBranch: defaultBranch, - } -} - -func (s *Spindle) resolvePipelineSourceRepo(ctx context.Context, trigger *tangled.Pipeline_TriggerMetadata) (*tangled.Pipeline_TriggerRepo, error) { - if trigger == nil { - return nil, nil - } - if trigger.SourceRepo == nil || *trigger.SourceRepo == "" { - return trigger.Repo, nil - } - repoDid, err := syntax.ParseDID(*trigger.SourceRepo) - if err != nil { - return nil, fmt.Errorf("parse sourceRepo %s: %w", *trigger.SourceRepo, err) - } - return s.resolveSourceRepoInfo(ctx, repoDid) -} - -// resolveSourceRepoInfo resolves trigger-repo metadata for a source repo DID. -func (s *Spindle) resolveSourceRepoInfo(ctx context.Context, repoDid syntax.DID) (*tangled.Pipeline_TriggerRepo, error) { - repo, err := s.db.GetRepoByDid(repoDid) - if err == nil { - return s.buildTriggerRepo(ctx, repo) - } - - // verify repo, we don't want git sync to point to arbitrary endpoints - res, err := s.verify(ctx, repoident.RepoDid(repoDid)) - if err != nil { - return nil, fmt.Errorf("verify sourceRepo %s: %w", repoDid, err) - } - return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host, res.OwnerDid.String(), res.Rkey, repoDid.String()), nil -} - -// runPipeline compiles and enqueues the pipeline for the given revision. -// sourceRepo is the resolved repo the code was checked out from, forwarded to -// processPipeline for env vars. -func (s *Spindle) runPipeline(ctx context.Context, repoDid syntax.DID, trigger tangled.Pipeline_TriggerMetadata, changedFiles []string, repoCloneUri, repoPath, rev string, only []string, sourceRepo *tangled.Pipeline_TriggerRepo) (models.PipelineId, error) { - l := log.FromContext(ctx) - - compiler := workflow.Compiler{ - ChangedFiles: changedFiles, - Trigger: trigger, - } - - rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev) - if err != nil { - return models.PipelineId{}, fmt.Errorf("loading pipeline: %w", err) - } - if len(rawPipeline) == 0 { - return models.PipelineId{}, nil - } - - tpl := compiler.Compile(compiler.Parse(rawPipeline)) - // todo(dawn): pass compile error to workflow log - for _, w := range compiler.Diagnostics.Errors { - l.Error(w.String()) - } - for _, w := range compiler.Diagnostics.Warnings { - l.Warn(w.String()) - } - - if len(only) > 0 { - tpl.Workflows = filterWorkflows(tpl.Workflows, only) - } - if len(tpl.Workflows) == 0 { - return models.PipelineId{}, nil - } - - pipelineId := models.PipelineId{ - Knot: trigger.Repo.Knot, - Rkey: tid.TID(), - } - if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { - return models.PipelineId{}, fmt.Errorf("creating pipeline event: %w", err) - } - err = s.processPipeline(repoDid, tpl, pipelineId, sourceRepo) - return pipelineId, err -} - -// filterWorkflows filters workflows to the requested names -func filterWorkflows(workflows []*tangled.Pipeline_Workflow, only []string) []*tangled.Pipeline_Workflow { - allowed := make(map[string]struct{}, len(only)) - for _, n := range only { - allowed[n] = struct{}{} - } - var filtered []*tangled.Pipeline_Workflow - for _, w := range workflows { - if w == nil { - continue - } - if _, ok := allowed[w.Name]; ok { - filtered = append(filtered, w) - } - } - return filtered -} - -// TriggerManual dispatches a pipeline at sha, authorized against and recorded -// under repoDid. sourceRepo, pull, and inputs are optional trigger payload. -func (s *Spindle) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull xrpc.PullContext, inputs []*tangled.Pipeline_Pair) (syntax.ATURI, error) { - repo, err := s.db.GetRepoByDid(repoDid) - if err != nil { - return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) - } - - triggerRepo, err := s.buildTriggerRepo(ctx, repo) - if err != nil { - return "", fmt.Errorf("building trigger repo: %w", err) - } - - trigger := tangled.Pipeline_TriggerMetadata{Repo: triggerRepo} - if pull.IsPullRequest { - var pullAt *string - if pull.Pull != "" { - pullAtStr := pull.Pull.String() - pullAt = &pullAtStr - } - trigger.Kind = string(workflow.TriggerKindPullRequest) - trigger.PullRequest = &tangled.Pipeline_PullRequestTriggerData{ - SourceBranch: pull.SourceBranch, - TargetBranch: pull.TargetBranch, - SourceSha: sha, - Pull: pullAt, - } - } else { - var refPtr *string - if ref != "" { - refPtr = &ref - } - trigger.Kind = string(workflow.TriggerKindManual) - trigger.Manual = &tangled.Pipeline_ManualTriggerData{ - Sha: sha, - Ref: refPtr, - Inputs: inputs, - } - } - - repoCloneUri := s.newRepoCloneUrl(repo.Knot, repoDid) - repoPath := s.newRepoPath(repoDid) - sourceInfo := triggerRepo // default: code comes from the repo itself - if sourceRepo != "" && sourceRepo != repoDid { - sourceInfo, err = s.resolveSourceRepoInfo(ctx, sourceRepo) - if err != nil { - return "", err - } - sourceRepoStr := sourceRepo.String() - trigger.SourceRepo = &sourceRepoStr - repoCloneUri = models.BuildRepoURL(sourceInfo) - repoPath = s.newRepoPath(sourceRepo) - } - - pipelineId, err := s.runPipeline(ctx, repoDid, trigger, nil, repoCloneUri, repoPath, sha, workflows, sourceInfo) - if err != nil { - return "", err - } - if pipelineId.Rkey == "" { - return "", xrpc.ErrNoMatchingWorkflows - } - return pipelineId.AtUri(), nil -} - -func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, error) { - if err := git.SparseSyncGitRepo(ctx, repoUri, repoPath, rev); err != nil { - return nil, fmt.Errorf("syncing git repo: %w", err) - } - gr, err := kgit.Open(repoPath, rev) - if err != nil { - return nil, fmt.Errorf("opening git repo: %w", err) - } - - workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) - if errors.Is(err, object.ErrDirectoryNotFound) { - // return empty RawPipeline when directory doesn't exist - return nil, nil - } else if err != nil { - return nil, fmt.Errorf("loading file tree: %w", err) - } - - var rawPipeline workflow.RawPipeline - for _, e := range workflowDir { - if !e.IsFile() { - continue - } - - fpath := filepath.Join(workflow.WorkflowDir, e.Name) - contents, err := gr.RawContent(fpath) - if err != nil { - return nil, fmt.Errorf("reading raw content of '%s': %w", fpath, err) - } - - rawPipeline = append(rawPipeline, workflow.RawWorkflow{ - Name: e.Name, - Contents: contents, - }) - } - - return rawPipeline, nil -} - -func (s *Spindle) StartJobWorkers(ctx context.Context) { - for range s.cfg.Server.MaxJobCount { - go func() { - for { - job, err := s.db.DequeueJob(ctx) - if err != nil { - s.l.Error("failed to dequeue job", "error", err) - } - if job == nil { - // sleep until a new job wakes us - select { - case <-ctx.Done(): - return - case <-s.jobWake: - } - continue - } - s.runJob(ctx, job) - } - }() - } -} - -func (s *Spindle) runJob(ctx context.Context, job *db.JobRow) { - pipelineId := models.PipelineId{ - Knot: job.PipelineIdKnot, - Rkey: job.PipelineIdRkey, - } - - pipelineEnv := models.PipelineEnvVarsForSource(job.Tpl.TriggerMetadata, pipelineId, job.SourceRepo) - trustedSource := true - if tm := job.Tpl.TriggerMetadata; tm != nil && tm.SourceRepo != nil && - *tm.SourceRepo != "" && *tm.SourceRepo != job.RepoDid { - trustedSource = false - } - - initTpl := job.Tpl - if job.SourceRepo != nil && job.Tpl.TriggerMetadata != nil { - tm := *job.Tpl.TriggerMetadata - tm.Repo = job.SourceRepo - initTpl.TriggerMetadata = &tm - } - - workflows := make(map[models.Engine][]models.Workflow) - for _, w := range job.Tpl.Workflows { - if w == nil { - continue - } - eng, ok := s.engs[w.Engine] - if !ok { - _ = s.db.StatusFailed(models.WorkflowId{ - PipelineId: pipelineId, - Name: w.Name, - }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) - continue - } - - ewf, err := eng.InitWorkflow(*w, initTpl) - if err != nil { - _ = s.db.StatusFailed(models.WorkflowId{ - PipelineId: pipelineId, - Name: w.Name, - }, fmt.Sprintf("init workflow: %s", err), -1, s.n) - continue - } - - if ewf.Environment == nil { - ewf.Environment = make(map[string]string) - } - maps.Copy(ewf.Environment, pipelineEnv) - workflows[eng] = append(workflows[eng], *ewf) - } - - engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.rootCtx, &models.Pipeline{ - RepoDid: syntax.DID(job.RepoDid), - Workflows: workflows, - TrustedSource: trustedSource, - }, pipelineId) -} - -// enqueues the workflows in tpl. -func (s *Spindle) processPipeline(repoDid syntax.DID, tpl tangled.Pipeline, pipelineId models.PipelineId, sourceRepo *tangled.Pipeline_TriggerRepo) error { - err := s.db.EnqueueJob(s.rootCtx, repoDid.String(), pipelineId, sourceRepo, tpl) - if err != nil { - return fmt.Errorf("failed to enqueue durable job: %w", err) - } - s.l.Info("pipeline enqueued successfully to db", "id", pipelineId) - - // wake up an idle worker to pick up more jobs if any - select { - case s.jobWake <- struct{}{}: - default: - } - - // pipelines visible from now on, they are sitting in queue - for _, w := range tpl.Workflows { - if w == nil { - continue - } - if err := s.db.StatusPending(models.WorkflowId{ - PipelineId: pipelineId, - Name: w.Name, - }, s.n); err != nil { - return fmt.Errorf("db.StatusPending: %w", err) - } - } - 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(repo syntax.DID) string { - return filepath.Join(s.cfg.Server.RepoDir, repo.String()) -} - -func (s *Spindle) newRepoCloneUrl(knot string, did syntax.DID) string { - scheme := "https://" - if s.cfg.Server.Dev { - scheme = "http://" - } - return fmt.Sprintf("%s%s/%s", scheme, knot, did) -} - -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 -} - -func (s *Spindle) resolvePipelineRepoDid(repo *tangled.Pipeline_TriggerRepo) (syntax.DID, error) { - if repo.RepoDid == nil || *repo.RepoDid == "" { - return "", fmt.Errorf("pipeline trigger missing repoDid") - } - repoDid, err := syntax.ParseDID(*repo.RepoDid) - if err != nil { - return "", fmt.Errorf("parse repoDid %s: %w", *repo.RepoDid, err) - } - if _, err := s.db.GetRepoByDid(repoDid); err != nil { - return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) - } - return repoDid, nil -} - -func (s *Spindle) configureOwner() error { - cfgOwner := s.cfg.Server.Owner - - existing, err := s.e.GetSpindleUsersByRole("server:owner", rbacDomain) - if err != nil { - return err - } - - switch len(existing) { - case 0: - // no owner configured, continue - case 1: - // find existing owner - existingOwner := existing[0] - - // no ownership change, this is okay - if existingOwner == s.cfg.Server.Owner { - break - } - - // remove existing owner - err = s.e.RemoveSpindleOwner(rbacDomain, existingOwner) - if err != nil { - return nil - } - default: - return fmt.Errorf("more than one owner in DB, try deleting %q and starting over", s.cfg.Server.DBPath) - } - - return s.e.AddSpindleOwner(rbacDomain, cfgOwner) -} diff --git a/spindle/server.go.mill b/spindle/server.go.mill deleted file mode 100644 index 65c0bafd8..000000000 --- a/spindle/server.go.mill +++ /dev/null @@ -1,1028 +0,0 @@ -package spindle - -import ( - "context" - "database/sql" - _ "embed" - "encoding/json" - "errors" - "fmt" - "log/slog" - "maps" - "net/http" - "path/filepath" - "sync" - "time" - - "github.com/bluesky-social/indigo/atproto/syntax" - indigoxrpc "github.com/bluesky-social/indigo/xrpc" - "github.com/go-chi/chi/v5" - "github.com/go-git/go-git/v5/plumbing/object" - "github.com/hashicorp/go-version" - "tangled.org/core/api/tangled" - "tangled.org/core/eventconsumer" - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/eventstream" - "tangled.org/core/idresolver" - "tangled.org/core/jetstream" - knotdb "tangled.org/core/knotserver/db" - kgit "tangled.org/core/knotserver/git" - "tangled.org/core/log" - "tangled.org/core/notifier" - "tangled.org/core/rbac" - "tangled.org/core/repoident" - "tangled.org/core/repoverify" - "tangled.org/core/spindle/artifactstore" - "tangled.org/core/spindle/config" - "tangled.org/core/spindle/db" - "tangled.org/core/spindle/engine" - "tangled.org/core/spindle/engines/dummy" - "tangled.org/core/spindle/engines/nixery" - "tangled.org/core/spindle/git" - "tangled.org/core/spindle/mill" - "tangled.org/core/spindle/mill/executor" - "tangled.org/core/spindle/models" - "tangled.org/core/spindle/queue" - "tangled.org/core/spindle/secrets" - "tangled.org/core/spindle/xrpc" - "tangled.org/core/tid" - "tangled.org/core/workflow" - "tangled.org/core/xrpc/serviceauth" -) - -//go:embed motd -var defaultMotd []byte - -const ( - rbacDomain = "thisserver" -) - -type Spindle struct { - jc *jetstream.JetstreamClient - tap *Tap - embedTap *embeddedTap - db *db.DB - e *rbac.Enforcer - l *slog.Logger - n *notifier.Notifier - engs map[string]models.Engine - jq *queue.Queue - cfg *config.Config - ks *eventconsumer.Consumer - res *idresolver.Resolver - verify repoverify.Verifier - vault secrets.Manager - motd []byte - motdMu sync.RWMutex - rootCtx context.Context - store artifactstore.Store - stores *artifactstore.Stores - reader artifactstore.Reader - // set only when this spindle hosts the mill or joins one as an executor - mill *mill.Mill - exec *executor.Executor -} - -// New creates a new Spindle server with the provided configuration and engines. -func New(ctx context.Context, cfg *config.Config, d *db.DB, engines map[string]models.Engine) (*Spindle, error) { - logger := log.FromContext(ctx) - n := notifier.New() - - if cfg.Role == config.RoleExecutor { - if err := cleanupOrphanRepos(ctx, d, logger); err != nil { - return nil, fmt.Errorf("failed to run startup cleanup: %w", err) - } - } else if err := runStartupMigrations(ctx, d, cfg.Server.Tap.Embed, cfg.Server.Tap.DBPath, logger); err != nil { - return nil, fmt.Errorf("failed to run startup migrations: %w", err) - } - - spindle := &Spindle{ - db: d, - l: logger, - n: &n, - engs: engines, - cfg: cfg, - motd: defaultMotd, - rootCtx: ctx, - } - diskFallback := "" - if cfg.Role == config.RoleStandalone { - diskFallback = cfg.Server.LogDir - if cfg.ArtifactStores.Disk.Dir == "" { - logger.Warn("using SPINDLE_SERVER_LOG_DIR as the implicit disk artifact store; configure SPINDLE_ARTIFACT_STORES_DISK_DIR explicitly") - } - } - stores, err := artifactstore.NewStores(cfg.ArtifactStores, diskFallback, cfg.LegacyS3.LogBucket) - if err != nil { - return nil, fmt.Errorf("failed to setup artifact stores: %w", err) - } - spindle.stores = stores - if cfg.LegacyS3.LogBucket != "" { - logger.Warn("SPINDLE_S3_LOG_BUCKET is deprecated; use SPINDLE_ARTIFACT_STORES_S3_BUCKET") - } - if cfg.Role == config.RoleStandalone { - spindle.reader = stores - } else { - name := cfg.Mill.ArtifactStore - if name == "" { - names := stores.Names() - if len(names) != 1 { - return nil, fmt.Errorf("%s requires SPINDLE_MILL_ARTIFACT_STORE when %d artifact stores are configured", cfg.Role, len(names)) - } - name = names[0] - logger.Warn("SPINDLE_MILL_ARTIFACT_STORE is not set; inferred the only configured store", "store", name) - } - store, ok := stores.Store(name) - if !ok { - return nil, fmt.Errorf("SPINDLE_MILL_ARTIFACT_STORE=%q is not configured", name) - } - spindle.store = store - spindle.reader = store - } - if cfg.Role == config.RoleExecutor { - return spindle, nil - } - - e, err := rbac.NewEnforcer(cfg.Server.DBPath) - if err != nil { - return nil, fmt.Errorf("failed to setup rbac enforcer: %w", err) - } - e.E.EnableAutoSave(true) - spindle.e = e - - switch cfg.Server.Secrets.Provider { - case "openbao": - if cfg.Server.Secrets.OpenBao.ProxyAddr == "" { - return nil, fmt.Errorf("openbao proxy address is required when using openbao secrets provider") - } - spindle.vault, err = secrets.NewOpenBaoManager( - cfg.Server.Secrets.OpenBao.ProxyAddr, - logger, - secrets.WithMountPath(cfg.Server.Secrets.OpenBao.Mount), - ) - if err != nil { - return nil, fmt.Errorf("failed to setup openbao secrets provider: %w", err) - } - logger.Info("using openbao secrets provider", "proxy_address", cfg.Server.Secrets.OpenBao.ProxyAddr, "mount", cfg.Server.Secrets.OpenBao.Mount) - case "sqlite", "": - spindle.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) - default: - return nil, fmt.Errorf("unknown secrets provider: %s", cfg.Server.Secrets.Provider) - } - - if cfg.Role == config.RoleStandalone { - spindle.jq = queue.NewQueue(cfg.Server.QueueSize, cfg.Server.MaxJobCount) - logger.Info("initialized queue", "queueSize", cfg.Server.QueueSize, "numWorkers", cfg.Server.MaxJobCount) - } - - collections := []string{ - tangled.SpindleMemberNSID, - tangled.RepoNSID, - tangled.RepoCollaboratorNSID, - tangled.RepoPullNSID, - } - jc, err := jetstream.NewJetstreamClient(cfg.Server.JetstreamEndpoint, "spindle", collections, nil, log.SubLogger(logger, "jetstream"), d, true, true) - if err != nil { - return nil, fmt.Errorf("failed to setup jetstream client: %w", err) - } - spindle.jc = jc - jc.AddDid(cfg.Server.Owner) - // pull records are created by arbitrary users too, same hack as in tap - jc.ExemptCollection(tangled.RepoPullNSID) - - // Check if the spindle knows about any Dids; - dids, err := d.GetAllDids() - if err != nil { - return nil, fmt.Errorf("failed to get all dids: %w", err) - } - for _, d := range dids { - jc.AddDid(d) - } - - knownRepos, err := d.AllRepos() - if err != nil { - return nil, fmt.Errorf("failed to get known repos: %w", err) - } - for _, r := range knownRepos { - if r.Owner != "" { - jc.AddDid(r.Owner.String()) - } - } - - spindle.res = idresolver.DefaultResolver(cfg.Server.PlcUrl) - spindle.verify = repoverify.New(spindle.res, cfg.Server.Dev) - - err = e.AddSpindle(rbacDomain) - if err != nil { - return nil, fmt.Errorf("failed to set rbac domain: %w", err) - } - err = spindle.configureOwner() - if err != nil { - return nil, err - } - logger.Info("owner set", "did", cfg.Server.Owner) - - cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath) - if err != nil { - return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err) - } - - err = jc.StartJetstream(ctx, spindle.ingest()) - if err != nil { - return nil, fmt.Errorf("failed to start jetstream consumer: %w", err) - } - - // spindle listen to knot stream for sh.tangled.git.refUpdate - // which will sync the local workflow files in spindle and enqueues the - // pipeline job for on-push workflows - ccfg := eventconsumer.NewConsumerConfig() - ccfg.Logger = log.SubLogger(logger, "eventconsumer") - ccfg.ProcessFunc = spindle.processKnotStream - ccfg.CursorStore = cursorStore - if cfg.Server.Dev { - ccfg.RetryInterval = 5 * time.Second - ccfg.MaxRetryInterval = 10 * time.Second - } else { - ccfg.RetryInterval = 1 * time.Minute - ccfg.MaxRetryInterval = 10 * time.Minute - } - knownKnots, err := d.Knots() - if err != nil { - return nil, err - } - for _, knot := range knownKnots { - logger.Info("adding source start", "knot", knot) - src := eventconsumer.NewKnotSource(knot) - eventconsumer.MigrateLegacyCursor(cursorStore, src) - ccfg.Sources[src] = struct{}{} - } - spindle.ks = eventconsumer.NewConsumer(*ccfg) - if cfg.Server.Tap.Embed { - pw, err := randomAdminPassword() - if err != nil { - return nil, err - } - cfg.Server.Tap.AdminPassword = pw - logger.Info("embedded tap: using random admin password") - } - spindle.tap = NewTapClient(spindle) - - return spindle, nil -} -func (s *Spindle) DB() *db.DB { - return s.db -} -func (s *Spindle) Queue() *queue.Queue { - return s.jq -} - -// Engines returns the map of available engines. -func (s *Spindle) Engines() map[string]models.Engine { - return s.engs -} - -// Vault returns the secrets manager instance. -func (s *Spindle) Vault() secrets.Manager { - return s.vault -} - -// Notifier returns the notifier instance. -func (s *Spindle) Notifier() *notifier.Notifier { - return s.n -} - -// Enforcer returns the RBAC enforcer instance. -func (s *Spindle) Enforcer() *rbac.Enforcer { - return s.e -} - -// SetMotdContent sets custom MOTD content, replacing the embedded default. -func (s *Spindle) SetMotdContent(content []byte) { - s.motdMu.Lock() - defer s.motdMu.Unlock() - s.motd = content -} - -// GetMotdContent returns the current MOTD content. -func (s *Spindle) GetMotdContent() []byte { - s.motdMu.RLock() - defer s.motdMu.RUnlock() - return s.motd -} - -// runs the server. blocks -func (s *Spindle) Start(ctx context.Context) error { - // only standalone runs the local queue. mill hosts place directly onto - // executors, and executors only run jobs explicitly assigned by a mill - if s.cfg.Role == config.RoleStandalone { - if s.jq != nil { - s.jq.Start() - defer s.jq.Stop() - } - } - - // an executor dials out to its mill and takes work from it - if s.exec != nil { - go s.exec.Connect(ctx) - } - - if stopper, ok := s.vault.(secrets.Stopper); ok { - defer stopper.Stop() - } - - if s.cfg.Role != config.RoleExecutor { - tapCtx, tapCancel := context.WithCancel(ctx) - - if s.cfg.Server.Tap.Embed { - emb, err := startEmbeddedTap(tapCtx, s.cfg, log.SubLogger(s.l, "embedtap")) - if err != nil { - tapCancel() - return fmt.Errorf("starting embedded tap: %w", err) - } - s.embedTap = emb - defer func() { - tapCancel() - s.embedTap.Shutdown() - }() - - go s.watchTapDrain(tapCtx, tapCancel) - } else { - defer tapCancel() - } - - go func() { - s.l.Info("starting knot event consumer") - s.ks.Start(ctx) - }() - - s.l.Info("starting tap client", "url", s.cfg.Server.Tap.Url) - s.tap.Start(tapCtx) - } - - s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr) - return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router()) -} - -func (s *Spindle) declareTapInterest(ctx context.Context) { - repos, err := s.db.AllRepos() - if err != nil { - s.l.Warn("tap declare: failed to load known repos", "err", err) - return - } - seen := make(map[syntax.DID]struct{}, len(repos)) - dids := make([]syntax.DID, 0, len(repos)) - for _, r := range repos { - if r.Owner == "" { - continue - } - if _, ok := seen[r.Owner]; ok { - continue - } - seen[r.Owner] = struct{}{} - dids = append(dids, r.Owner) - } - if err := s.tap.AddOwnerDIDs(ctx, dids); err != nil { - s.l.Warn("tap declare: AddRepos rejected", "count", len(dids), "err", err) - return - } - s.l.Info("tap declare: known owner DIDs registered", "count", len(dids)) -} - -func Run(ctx context.Context) error { - cfg, err := config.Load(ctx) - if err != nil { - return fmt.Errorf("failed to load config: %w", err) - } - - if err := ensureGitVersion(); err != nil { - return fmt.Errorf("ensuring git version: %w", err) - } - - d, err := db.Make(ctx, cfg.Server.DBPath) - if err != nil { - return fmt.Errorf("failed to setup db: %w", err) - } - - logger := log.FromContext(ctx) - - var engines map[string]models.Engine - var m *mill.Mill - - if cfg.Role == config.RoleMill { - // mill host: register engines that place jobs on executors instead of - // running them. all names share one Mill - m = mill.New(log.SubLogger(logger, "mill"), mill.Config{ - LogDir: cfg.Server.LogDir, - MaxPending: cfg.Mill.MaxPending, - ReconnectGrace: cfg.Mill.ReconnectGrace, - }) - engines = map[string]models.Engine{ - "nixery": mill.NewEngine("nixery", m), - "microvm": mill.NewEngine("microvm", m), - "dummy": mill.NewEngine("dummy", m), - } - } else { - // standalone and executor both run real engines locally. - nixeryEng, err := nixery.New(ctx, cfg) - if err != nil { - return err - } - microvmEng, err := newMicrovmEngine(ctx, cfg, d) - if err != nil { - return err - } - engines = map[string]models.Engine{ - "nixery": nixeryEng, - "microvm": microvmEng, - "dummy": dummy.New(logger), - } - } - - s, err := New(ctx, cfg, d, engines) - if err != nil { - return err - } - - if m != nil { - // resolve the chicken-and-egg: the engines (built above) hold the mill, - // but the mill's db/notifier are created inside New - m.Attach(s.DB(), s.Notifier()) - s.mill = m - if err := m.RestoreState(); err != nil { - return fmt.Errorf("restoring mill state: %w", err) - } - } - if cfg.Role == config.RoleExecutor { - s.exec, err = executor.New(cfg, engines, s.DB(), s.Notifier(), log.SubLogger(logger, "executor"), s.store) - if err != nil { - return err - } - } - - return s.Start(ctx) -} - -func (s *Spindle) Router() http.Handler { - mux := chi.NewRouter() - - mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { - w.Write(s.GetMotdContent()) - }) - if s.cfg.Role == config.RoleExecutor { - return mux - } - - mux.HandleFunc("/events", s.Events) - mux.HandleFunc("/logs/{knot}/{rkey}/{name}", s.Logs) - - // mill host: executors dial in here (plain ws, shared-secret auth) - if s.mill != nil { - mux.HandleFunc("/mill", s.mill.HandleExecutorConn) - } - - mux.Mount("/xrpc", s.XrpcRouter()) - return mux -} - -func (s *Spindle) XrpcRouter() http.Handler { - serviceAuth := serviceauth.NewServiceAuth(s.l, s.res.Directory(), s.cfg.Server.Did().String()) - - l := log.SubLogger(s.l, "xrpc") - - x := xrpc.Xrpc{ - Logger: l, - Db: s.db, - Enforcer: s.e, - Engines: s.engs, - Config: s.cfg, - ArtifactReader: s.reader, - Resolver: s.res, - Vault: s.vault, - Notifier: s.Notifier(), - ServiceAuth: serviceAuth, - Trigger: s, - } - - return x.Router() -} - -func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) 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 == knotdb.RepoCollaboratorUpdateNSID { - return s.ingestKnotCollaborator(ctx, l, src, msg) - } - 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("repo", event.Repo, "ref", event.Ref, "newSha", event.NewSha) - l.Debug("debug") - - repoDid := syntax.DID(event.Repo) - repo, err := s.db.GetRepoByDid(repoDid) - if err != nil { - return fmt.Errorf("unknown repoDid %s: %w", repoDid, err) - } - - if src.Host != repo.Knot { - return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, repo.Knot) - } - - if kgit.HasSkipCIPushOption(event.PushOptions) { - l.Info("push event requested ci skip, skipping the event") - return nil - } - - // NOTE: we are blindly trusting the knot that it will return only repos it own - repoCloneUri := s.newRepoCloneUrl(src.Host, repoDid) - repoPath := s.newRepoPath(repoDid) - if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, event.NewSha); err != nil { - return fmt.Errorf("sync git repo: %w", err) - } - l.Info("synced git repo") - - triggerRepo, err := s.buildTriggerRepo(ctx, repo) - if err != nil { - return fmt.Errorf("building trigger repo: %w", err) - } - - trigger := tangled.Pipeline_TriggerMetadata{ - Kind: string(workflow.TriggerKindPush), - Push: &tangled.Pipeline_PushTriggerData{ - Ref: event.Ref, - OldSha: event.OldSha, - NewSha: event.NewSha, - }, - Repo: triggerRepo, - } - - pipelineId, err := s.runPipeline(ctx, repoDid, trigger, event.ChangedFiles, repoCloneUri, repoPath, event.NewSha, nil, triggerRepo) - if err != nil { - return err - } - if pipelineId.Rkey == "" { - l.Info("no workflow matched 'push' trigger, skipping the event") - return nil - } - l.Info("pipeline triggered", "pipeline", pipelineId.AtUri()) - } - - return nil -} - -func (s *Spindle) ingestKnotCollaborator(ctx context.Context, l *slog.Logger, src eventconsumer.Source, msg eventstream.Event) error { - var rec knotdb.RepoCollaboratorUpdate - if err := json.Unmarshal(msg.EventJson, &rec); err != nil { - l.Error("error unmarshalling collaboratorUpdate", "err", err) - return err - } - - subject, err := syntax.ParseDID(rec.Subject) - if err != nil { - l.Info("skipping collaboratorUpdate with malformed subject", "subject", rec.Subject, "err", err) - return nil - } - repoDid, err := syntax.ParseDID(rec.Repo) - if err != nil { - l.Info("skipping collaboratorUpdate with malformed repo", "repo", rec.Repo, "err", err) - return nil - } - - repo, err := s.db.GetRepoByDid(repoDid) - if errors.Is(err, sql.ErrNoRows) { - l.Info("skipping collaboratorUpdate for unknown repo", "repo", repoDid) - return nil - } - if err != nil { - return fmt.Errorf("lookup repo %s: %w", repoDid, err) - } - if src.Host != repo.Knot { - l.Warn("dropping collaboratorUpdate from non-owning knot", "src", src.Host, "repoKnot", repo.Knot) - return nil - } - - switch rec.Op { - case knotdb.AclOpAdd: - if err := s.e.AddCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { - return fmt.Errorf("add collaborator policy: %w", err) - } - if err := s.db.AddKnotCollaborator(repoDid, subject); err != nil { - return fmt.Errorf("track collaborator: %w", err) - } - l.Info("added knot-managed collaborator", "subject", subject, "repo", repoDid) - case knotdb.AclOpRemove: - if err := s.e.RemoveCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { - return fmt.Errorf("remove collaborator policy: %w", err) - } - if err := s.db.DeleteRepoCollaboratorBySubjectRepo(subject, repoDid); err != nil { - return fmt.Errorf("delete collaborator row: %w", err) - } - l.Info("removed knot-managed collaborator", "subject", subject, "repo", repoDid) - default: - return fmt.Errorf("collaboratorUpdate unknown op %q", rec.Op) - } - return nil -} - -// buildTriggerRepo gathers trigger metadata, resolving default branch from the knot -func (s *Spindle) buildTriggerRepo(ctx context.Context, repo *db.Repo) (*tangled.Pipeline_TriggerRepo, error) { - rkey := string(repo.Rkey) - repoDid := repo.RepoDid.String() - return s.buildTriggerRepoFrom(ctx, repo.Knot, repo.Owner.String(), rkey, repoDid), nil -} - -func (s *Spindle) buildTriggerRepoFrom(ctx context.Context, knot, did, rkey, repoDid string) *tangled.Pipeline_TriggerRepo { - scheme := "https" - if s.cfg.Server.Dev { - scheme = "http" - } - client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, knot)} - - // this should maybe (?) be in the refUpdate event itself to save a roundtrip - defaultBranch := "" - if out, err := tangled.RepoGetDefaultBranch(ctx, client, repoDid); err == nil { - defaultBranch = out.Name - } - - var rkeyPtr *string - if rkey != "" { - rkeyPtr = &rkey - } - return &tangled.Pipeline_TriggerRepo{ - Did: did, - Knot: knot, - Repo: rkeyPtr, - RepoDid: &repoDid, - DefaultBranch: defaultBranch, - } -} - -func (s *Spindle) resolvePipelineSourceRepo(ctx context.Context, trigger *tangled.Pipeline_TriggerMetadata) (*tangled.Pipeline_TriggerRepo, error) { - if trigger == nil { - return nil, nil - } - if trigger.SourceRepo == nil || *trigger.SourceRepo == "" { - return trigger.Repo, nil - } - repoDid, err := syntax.ParseDID(*trigger.SourceRepo) - if err != nil { - return nil, fmt.Errorf("parse sourceRepo %s: %w", *trigger.SourceRepo, err) - } - return s.resolveSourceRepoInfo(ctx, repoDid) -} - -// resolveSourceRepoInfo resolves trigger-repo metadata for a source repo DID. -func (s *Spindle) resolveSourceRepoInfo(ctx context.Context, repoDid syntax.DID) (*tangled.Pipeline_TriggerRepo, error) { - repo, err := s.db.GetRepoByDid(repoDid) - if err == nil { - return s.buildTriggerRepo(ctx, repo) - } - - // verify repo, we don't want git sync to point to arbitrary endpoints - res, err := s.verify(ctx, repoident.RepoDid(repoDid)) - if err != nil { - return nil, fmt.Errorf("verify sourceRepo %s: %w", repoDid, err) - } - return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host, res.OwnerDid.String(), res.Rkey, repoDid.String()), nil -} - -// runPipeline compiles and enqueues the pipeline for the given revision. -// sourceRepo is the resolved repo the code was checked out from, forwarded to -// processPipeline for env vars. -func (s *Spindle) runPipeline(ctx context.Context, repoDid syntax.DID, trigger tangled.Pipeline_TriggerMetadata, changedFiles []string, repoCloneUri, repoPath, rev string, only []string, sourceRepo *tangled.Pipeline_TriggerRepo) (models.PipelineId, error) { - l := log.FromContext(ctx) - - compiler := workflow.Compiler{ - ChangedFiles: changedFiles, - Trigger: trigger, - } - - rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev) - if err != nil { - return models.PipelineId{}, fmt.Errorf("loading pipeline: %w", err) - } - if len(rawPipeline) == 0 { - return models.PipelineId{}, nil - } - - tpl := compiler.Compile(compiler.Parse(rawPipeline)) - // todo(dawn): pass compile error to workflow log - for _, w := range compiler.Diagnostics.Errors { - l.Error(w.String()) - } - for _, w := range compiler.Diagnostics.Warnings { - l.Warn(w.String()) - } - - if len(only) > 0 { - tpl.Workflows = filterWorkflows(tpl.Workflows, only) - } - if len(tpl.Workflows) == 0 { - return models.PipelineId{}, nil - } - - pipelineId := models.PipelineId{ - Knot: trigger.Repo.Knot, - Rkey: tid.TID(), - } - if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { - return models.PipelineId{}, fmt.Errorf("creating pipeline event: %w", err) - } - err = s.processPipeline(repoDid, tpl, pipelineId, sourceRepo) - return pipelineId, err -} - -// filterWorkflows filters workflows to the requested names -func filterWorkflows(workflows []*tangled.Pipeline_Workflow, only []string) []*tangled.Pipeline_Workflow { - allowed := make(map[string]struct{}, len(only)) - for _, n := range only { - allowed[n] = struct{}{} - } - var filtered []*tangled.Pipeline_Workflow - for _, w := range workflows { - if w == nil { - continue - } - if _, ok := allowed[w.Name]; ok { - filtered = append(filtered, w) - } - } - return filtered -} - -// TriggerManual dispatches a pipeline at sha, authorized against and recorded -// under repoDid. sourceRepo, pull, and inputs are optional trigger payload. -func (s *Spindle) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull xrpc.PullContext, inputs []*tangled.Pipeline_Pair) (syntax.ATURI, error) { - repo, err := s.db.GetRepoByDid(repoDid) - if err != nil { - return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) - } - - triggerRepo, err := s.buildTriggerRepo(ctx, repo) - if err != nil { - return "", fmt.Errorf("building trigger repo: %w", err) - } - - trigger := tangled.Pipeline_TriggerMetadata{Repo: triggerRepo} - if pull.IsPullRequest { - var pullAt *string - if pull.Pull != "" { - pullAtStr := pull.Pull.String() - pullAt = &pullAtStr - } - trigger.Kind = string(workflow.TriggerKindPullRequest) - trigger.PullRequest = &tangled.Pipeline_PullRequestTriggerData{ - SourceBranch: pull.SourceBranch, - TargetBranch: pull.TargetBranch, - SourceSha: sha, - Pull: pullAt, - } - } else { - var refPtr *string - if ref != "" { - refPtr = &ref - } - trigger.Kind = string(workflow.TriggerKindManual) - trigger.Manual = &tangled.Pipeline_ManualTriggerData{ - Sha: sha, - Ref: refPtr, - Inputs: inputs, - } - } - - repoCloneUri := s.newRepoCloneUrl(repo.Knot, repoDid) - repoPath := s.newRepoPath(repoDid) - sourceInfo := triggerRepo // default: code comes from the repo itself - if sourceRepo != "" && sourceRepo != repoDid { - sourceInfo, err = s.resolveSourceRepoInfo(ctx, sourceRepo) - if err != nil { - return "", err - } - sourceRepoStr := sourceRepo.String() - trigger.SourceRepo = &sourceRepoStr - repoCloneUri = models.BuildRepoURL(sourceInfo) - repoPath = s.newRepoPath(sourceRepo) - } - - pipelineId, err := s.runPipeline(ctx, repoDid, trigger, nil, repoCloneUri, repoPath, sha, workflows, sourceInfo) - if err != nil { - return "", err - } - if pipelineId.Rkey == "" { - return "", xrpc.ErrNoMatchingWorkflows - } - return pipelineId.AtUri(), nil -} - -func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, error) { - if err := git.SparseSyncGitRepo(ctx, repoUri, repoPath, rev); err != nil { - return nil, fmt.Errorf("syncing git repo: %w", err) - } - gr, err := kgit.Open(repoPath, rev) - if err != nil { - return nil, fmt.Errorf("opening git repo: %w", err) - } - - workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) - if errors.Is(err, object.ErrDirectoryNotFound) { - // return empty RawPipeline when directory doesn't exist - return nil, nil - } else if err != nil { - return nil, fmt.Errorf("loading file tree: %w", err) - } - - var rawPipeline workflow.RawPipeline - for _, e := range workflowDir { - if !e.IsFile() { - continue - } - - fpath := filepath.Join(workflow.WorkflowDir, e.Name) - contents, err := gr.RawContent(fpath) - if err != nil { - return nil, fmt.Errorf("reading raw content of '%s': %w", fpath, err) - } - - rawPipeline = append(rawPipeline, workflow.RawWorkflow{ - Name: e.Name, - Contents: contents, - }) - } - - return rawPipeline, nil -} - -// processPipeline enqueues the workflows in tpl. -func (s *Spindle) processPipeline(repoDid syntax.DID, tpl tangled.Pipeline, pipelineId models.PipelineId, sourceRepo *tangled.Pipeline_TriggerRepo) error { - // derive security-relevant things like whether this run is trusted and can be passed - // secrets to from the original metadata. - pipelineEnv := models.PipelineEnvVarsForSource(tpl.TriggerMetadata, pipelineId, sourceRepo) - trustedSource := true - if tm := tpl.TriggerMetadata; tm != nil && tm.SourceRepo != nil && - *tm.SourceRepo != "" && *tm.SourceRepo != repoDid.String() { - trustedSource = false - } - - // swap the repo with our sourceRepo if we are running a pipeline on a fork. - // the metadata stays the same. we check whether the repo is trusted above, - // so this only affects the clone URL. - initTpl := tpl - if sourceRepo != nil && tpl.TriggerMetadata != nil { - tm := *tpl.TriggerMetadata - tm.Repo = sourceRepo - initTpl.TriggerMetadata = &tm - } - - // filter & init workflows - workflows := make(map[models.Engine][]models.Workflow) - for _, w := range tpl.Workflows { - if w == nil { - continue - } - eng, ok := s.engs[w.Engine] - if !ok { - err := s.db.StatusFailed(models.WorkflowId{ - PipelineId: pipelineId, - Name: w.Name, - }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) - if err != nil { - return fmt.Errorf("db.StatusFailed: %w", err) - } - - continue - } - - ewf, err := eng.InitWorkflow(*w, initTpl) - if err != nil { - err = s.db.StatusFailed(models.WorkflowId{ - PipelineId: pipelineId, - Name: w.Name, - }, fmt.Sprintf("init workflow: %s", err), -1, s.n) - if err != nil { - return fmt.Errorf("db.StatusFailed: %w", err) - } - - continue - } - - // inject TANGLED_* env vars after InitWorkflow - // This prevents user-defined env vars from overriding them - if ewf.Environment == nil { - ewf.Environment = make(map[string]string) - } - maps.Copy(ewf.Environment, pipelineEnv) - - workflows[eng] = append(workflows[eng], *ewf) - } - - pipeline := &models.Pipeline{ - RepoDid: repoDid, - Workflows: workflows, - TrustedSource: trustedSource, - } - - if s.mill != nil { - // mill host: no bounded pool. each job blocks in placement - // (AcquireWorkflowSlot) which the user sees as pending. the only bound - // is the mill's maxPending. rootCtx is the long-lived consumer context, - // so the goroutine safely outlives this call - go engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.stores, s.db, s.n, s.rootCtx, pipeline, pipelineId) - s.l.Info("pipeline handed to mill placement", "id", pipelineId) - } else if s.jq != nil { - ok := s.jq.Enqueue(repoDid, queue.Job{ - Run: func() error { - engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.stores, s.db, s.n, s.rootCtx, pipeline, pipelineId) - return nil - }, - OnFail: func(jobError error) { - s.l.Error("pipeline run failed", "error", jobError) - }, - }) - if !ok { - return fmt.Errorf("failed to enqueue pipeline: queue is full") - } - s.l.Info("pipeline enqueued successfully", "id", pipelineId) - } else { - return fmt.Errorf("no queue or mill available to process pipeline") - } - - // after successful enqueue, emit StatusPending for all workflows - for _, ewfs := range workflows { - for _, ewf := range ewfs { - err := s.db.StatusPending(models.WorkflowId{ - PipelineId: pipelineId, - Name: ewf.Name, - }, s.n) - if err != nil { - return fmt.Errorf("db.StatusPending: %w", err) - } - } - } - 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(repo syntax.DID) string { - return filepath.Join(s.cfg.Server.RepoDir, repo.String()) -} - -func (s *Spindle) newRepoCloneUrl(knot string, did syntax.DID) string { - scheme := "https://" - if s.cfg.Server.Dev { - scheme = "http://" - } - return fmt.Sprintf("%s%s/%s", scheme, knot, did) -} - -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 -} - -func (s *Spindle) configureOwner() error { - cfgOwner := s.cfg.Server.Owner - - existing, err := s.e.GetSpindleUsersByRole("server:owner", rbacDomain) - if err != nil { - return err - } - - switch len(existing) { - case 0: - // no owner configured, continue - case 1: - // find existing owner - existingOwner := existing[0] - - // no ownership change, this is okay - if existingOwner == s.cfg.Server.Owner { - break - } - - // remove existing owner - err = s.e.RemoveSpindleOwner(rbacDomain, existingOwner) - if err != nil { - return nil - } - default: - return fmt.Errorf("more than one owner in DB, try deleting %q and starting over", s.cfg.Server.DBPath) - } - - return s.e.AddSpindleOwner(rbacDomain, cfgOwner) -} diff --git a/spindle/server_test.go b/spindle/server_test.go index f9fbd8324..5e4be236d 100644 --- a/spindle/server_test.go +++ b/spindle/server_test.go @@ -86,7 +86,7 @@ func TestExecutorRoleBuildsMinimalSpindle(t *testing.T) { t.Fatalf("New() error = %v", err) } - if s.jc != nil || s.tap != nil || s.e != nil || s.ks != nil || s.res != nil || s.vault != nil { + if s.jc != nil || s.tap != nil || s.e != nil || s.feed != nil || s.res != nil || s.vault != nil { t.Fatal("executor role built coordinator-only spindle dependencies") } diff --git a/spindle/tapclient.go b/spindle/tapclient.go index 81eb06dd5..ad1aa8bb8 100644 --- a/spindle/tapclient.go +++ b/spindle/tapclient.go @@ -18,7 +18,6 @@ import ( indigoxrpc "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" avmodels "tangled.org/core/appview/models" - "tangled.org/core/eventconsumer" "tangled.org/core/gitutil" "tangled.org/core/log" "tangled.org/core/rbac" @@ -212,8 +211,9 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error return fmt.Errorf("add repo policy: %w", err) } - src := eventconsumer.NewKnotSource(record.Knot) - t.spindle.ks.AddSource(t.spindle.rootCtx, src) + if t.spindle.feed != nil { + t.spindle.feed.Subscribe(t.spindle.rootCtx, record.Knot) + } repo := db.Repo{ Knot: record.Knot, @@ -275,10 +275,24 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error func (t *Tap) teardownRepo(l *slog.Logger, repo *db.Repo, ownerDid syntax.DID, rkey syntax.RecordKey) error { if repo.RepoDid != "" { - return t.spindle.WipeRepo(context.Background(), repo.RepoDid, "repo record removed") + if err := t.spindle.WipeRepo(context.Background(), repo.RepoDid, "repo record removed"); err != nil { + return err + } + } else { + // rows without a repo did never got past registration + if err := t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey); err != nil { + return err + } } - // rows without a repo did never got past registration - return t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey) + if t.spindle.feed != nil { + if left, err := t.spindle.db.CountReposByKnot(repo.Knot); err != nil { + l.Warn("counting repos left on knot", "knot", repo.Knot, "err", err) + } else if left == 0 { + t.spindle.feed.Unsubscribe(repo.Knot) + } + } + // TODO: clear sparse-synced git repo + return nil } func (t *Tap) processCollaborator(ctx context.Context, evt *tapc.RecordEventData) error { diff --git a/spindle/tapclient_test.go b/spindle/tapclient_test.go index a197a49bc..2c0027400 100644 --- a/spindle/tapclient_test.go +++ b/spindle/tapclient_test.go @@ -15,7 +15,6 @@ import ( "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" - "tangled.org/core/eventconsumer" "tangled.org/core/idresolver" "tangled.org/core/jetstream" "tangled.org/core/notifier" @@ -80,10 +79,6 @@ func TestProcessRepo_MembershipCheck(t *testing.T) { cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" - ccfg := eventconsumer.NewConsumerConfig() - ccfg.Logger = slog.Default() - ks := eventconsumer.NewConsumer(*ccfg) - jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) @@ -93,7 +88,6 @@ func TestProcessRepo_MembershipCheck(t *testing.T) { e: e, l: slog.Default(), cfg: cfg, - ks: ks, jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), @@ -261,10 +255,6 @@ func TestProcessRepo_HijackRepoDidCheck(t *testing.T) { cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" - ccfg := eventconsumer.NewConsumerConfig() - ccfg.Logger = slog.Default() - ks := eventconsumer.NewConsumer(*ccfg) - jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) @@ -274,7 +264,6 @@ func TestProcessRepo_HijackRepoDidCheck(t *testing.T) { e: e, l: slog.Default(), cfg: cfg, - ks: ks, jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false),