From 8bb90de1c562c4f3b7b232f48200009215dcb05d Mon Sep 17 00:00:00 2001 From: dawn Date: Mon, 13 Jul 2026 23:55:46 +0300 Subject: [PATCH] spindle: sync knot-managed collaborators from the acl firehose Signed-off-by: dawn --- spindle/db/collaborators.go | 21 +++++++ spindle/db/collaborators_test.go | 94 ++++++++++++++++++++++++++++++++ spindle/server.go | 59 ++++++++++++++++++++ 3 files changed, 174 insertions(+) create mode 100644 spindle/db/collaborators_test.go diff --git a/spindle/db/collaborators.go b/spindle/db/collaborators.go index 3356f946..2a5344c0 100644 --- a/spindle/db/collaborators.go +++ b/spindle/db/collaborators.go @@ -26,6 +26,16 @@ func (d *DB) AddRepoCollaborator(c RepoCollaborator) error { return err } +func (d *DB) AddKnotCollaborator(repoDid, subject syntax.DID) error { + _, err := d.Exec( + `insert into repo_collaborators (owner_did, rkey, subject, repo_did) + values (?, ?, ?, ?) + on conflict(owner_did, rkey) do nothing`, + repoDid.String(), subject.String(), subject.String(), repoDid.String(), + ) + return err +} + func scanCollab(row interface{ Scan(...any) error }) (*RepoCollaborator, error) { var owner, rkey, subject, repoDid string if err := row.Scan(&owner, &rkey, &subject, &repoDid); err != nil { @@ -61,6 +71,17 @@ func (d *DB) DeleteRepoCollaborator(ownerDid syntax.DID, rkey syntax.RecordKey) return nil } +func (d *DB) DeleteRepoCollaboratorBySubjectRepo(subject, repoDid syntax.DID) error { + _, err := d.Exec( + `delete from repo_collaborators where repo_did = ? and subject = ?`, + repoDid.String(), subject.String(), + ) + if err != nil { + return fmt.Errorf("delete collaborator %s on %s: %w", subject, repoDid, err) + } + return nil +} + func (d *DB) DeleteRepoCollaboratorsByRepoDid(repoDid syntax.DID) error { _, err := d.Exec(`delete from repo_collaborators where repo_did = ?`, repoDid.String()) if err != nil { diff --git a/spindle/db/collaborators_test.go b/spindle/db/collaborators_test.go new file mode 100644 index 00000000..8547811b --- /dev/null +++ b/spindle/db/collaborators_test.go @@ -0,0 +1,94 @@ +package db + +import ( + "testing" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +func subjectsOf(t *testing.T, d *DB, repoDid syntax.DID) []syntax.DID { + t.Helper() + rows, err := d.ListCollaboratorsByRepoDid(repoDid) + if err != nil { + t.Fatalf("ListCollaboratorsByRepoDid: %v", err) + } + out := make([]syntax.DID, 0, len(rows)) + for _, r := range rows { + out = append(out, r.Subject) + } + return out +} + +func TestAddKnotCollaborator_PersistsAndIsIdempotent(t *testing.T) { + d := newTestDB(t) + repo := syntax.DID("did:plc:repo") + bob := syntax.DID("did:plc:bob") + + if err := d.AddKnotCollaborator(repo, bob); err != nil { + t.Fatalf("add: %v", err) + } + if err := d.AddKnotCollaborator(repo, bob); err != nil { + t.Fatalf("re-add: %v", err) + } + + got := subjectsOf(t, d, repo) + if len(got) != 1 || got[0] != bob { + t.Fatalf("collaborators = %v, want exactly [bob]", got) + } +} + +func TestDeleteRepoCollaboratorBySubjectRepo(t *testing.T) { + d := newTestDB(t) + repo := syntax.DID("did:plc:repo") + bob := syntax.DID("did:plc:bob") + carol := syntax.DID("did:plc:carol") + + if err := d.AddKnotCollaborator(repo, bob); err != nil { + t.Fatalf("add bob: %v", err) + } + if err := d.AddKnotCollaborator(repo, carol); err != nil { + t.Fatalf("add carol: %v", err) + } + + if err := d.DeleteRepoCollaboratorBySubjectRepo(bob, repo); err != nil { + t.Fatalf("delete bob: %v", err) + } + got := subjectsOf(t, d, repo) + if len(got) != 1 || got[0] != carol { + t.Fatalf("after removing bob, collaborators = %v, want [carol]", got) + } + + if err := d.DeleteRepoCollaboratorBySubjectRepo(bob, repo); err != nil { + t.Fatalf("idempotent delete: %v", err) + } +} + +func TestKnotCollaborator_NoCollisionAcrossReposAndSubjects(t *testing.T) { + d := newTestDB(t) + repoA := syntax.DID("did:plc:repoA") + repoB := syntax.DID("did:plc:repoB") + bob := syntax.DID("did:plc:bob") + carol := syntax.DID("did:plc:carol") + + for _, c := range []struct{ repo, subj syntax.DID }{ + {repoA, bob}, {repoB, bob}, {repoA, carol}, + } { + if err := d.AddKnotCollaborator(c.repo, c.subj); err != nil { + t.Fatalf("add %s/%s: %v", c.repo, c.subj, err) + } + } + + if got := subjectsOf(t, d, repoA); len(got) != 2 { + t.Errorf("repoA collaborators = %v, want bob+carol", got) + } + if got := subjectsOf(t, d, repoB); len(got) != 1 || got[0] != bob { + t.Errorf("repoB collaborators = %v, want [bob]", got) + } + + if err := d.DeleteRepoCollaboratorBySubjectRepo(bob, repoA); err != nil { + t.Fatalf("delete bob@repoA: %v", err) + } + if got := subjectsOf(t, d, repoB); len(got) != 1 || got[0] != bob { + t.Errorf("repoB after removing bob@repoA = %v, want still [bob]", got) + } +} diff --git a/spindle/server.go b/spindle/server.go index 2bbc4ac3..a00461c4 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -2,6 +2,7 @@ package spindle import ( "context" + "database/sql" _ "embed" "encoding/json" "errors" @@ -24,6 +25,7 @@ import ( "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" @@ -408,6 +410,9 @@ func (s *Spindle) XrpcRouter() http.Handler { 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 { @@ -469,6 +474,60 @@ func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Sourc 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) -- 2.51.2