From 1b2cefa523bd8c1022c09748925ecc77899881a4 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Wed, 22 Jul 2026 19:40:14 +0000 Subject: [PATCH] spindle: remove collaborator PDS record ingest Signed-off-by: Seongmin Lee --- spindle/casbin_copy.go | 87 --------------------------------------------------------------------------------------- spindle/ingester.go | 2 +- spindle/server.go | 32 +++++++++++++------------------- spindle/startup_migrations.go | 23 ----------------------- spindle/startup_migrations_test.go | 468 ++---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- spindle/tapclient.go | 407 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ spindle/db/collaborators.go | 115 ------------------------------------------------------------------------------------------------------------------- spindle/db/collaborators_test.go | 94 ---------------------------------------------------------------------------------------------- spindle/db/db.go | 42 ++++++++++++++++++++++++++++++++++++++++++ spindle/db/repos.go | 27 ++++++++++++++++++++++----- spindle/xrpc/add_secret.go | 9 ++++++--- spindle/xrpc/ci_pipeline_trigger_pipeline.go | 3 +-- spindle/xrpc/list_secrets.go | 9 ++++++--- spindle/xrpc/remove_secret.go | 9 ++++++--- spindle/xrpc/xrpc.go | 2 +- 15 file(s) changed, 261 insertion(s)(+), 1068 deletion(s)(-) diff --git a/spindle/casbin_copy.go b/spindle/casbin_copy.go deleted file mode 100644 --- a/spindle/casbin_copy.go +++ /dev/null @@ -1,87 +0,0 @@ -package spindle - -import ( - "context" - "fmt" - "log/slog" - - "github.com/bluesky-social/indigo/atproto/syntax" - "tangled.org/core/rbac" - "tangled.org/core/spindle/db" -) - -func migrateLegacyRepoCasbin(ctx context.Context, d *db.DB, e *rbac.Enforcer, logger *slog.Logger, owner syntax.DID, name string, rkey syntax.RecordKey, repoDid syntax.DID) { - candidates := legacyKeyCandidates(owner, name, rkey) - if siblings, err := d.SiblingRkeysForRepoDid(owner, repoDid, rkey); err == nil { - var fold func(rest []string, acc []string) []string - fold = func(rest []string, acc []string) []string { - if len(rest) == 0 { - return acc - } - return fold(rest[1:], append(acc, owner.String()+"/"+rest[0])) - } - candidates = fold(siblings, candidates) - } else { - logger.Warn("legacy casbin rekey: sibling lookup failed", "err", err) - } - if len(candidates) == 0 { - return - } - flag := "legacy-casbin-rekey:" + repoDid.String() + ":" + rkey.String() - var exists bool - if err := d.QueryRowContext(ctx, `select exists (select 1 from migrations where name = ?)`, flag).Scan(&exists); err != nil { - logger.Warn("legacy casbin rekey: check migration flag", "err", err) - return - } - if exists { - return - } - - if err := e.AddRepo(owner.String(), rbac.ThisServer, repoDid.String()); err != nil { - logger.Warn("legacy casbin rekey: owner add new key failed", "err", err) - return - } - - collabs, err := d.ListCollaboratorsByRepoDid(repoDid) - if err != nil { - logger.Warn("legacy casbin rekey: list collaborators failed", "err", err) - return - } - - var addCollabs func(remaining []db.RepoCollaborator) error - addCollabs = func(remaining []db.RepoCollaborator) error { - if len(remaining) == 0 { - return nil - } - c := remaining[0] - if err := e.AddCollaborator(c.Subject.String(), rbac.ThisServer, repoDid.String()); err != nil { - return fmt.Errorf("AddCollaborator %s -> %s: %w", c.Subject, repoDid, err) - } - return addCollabs(remaining[1:]) - } - if err := addCollabs(collabs); err != nil { - logger.Warn("legacy casbin rekey: collaborator add failed", "err", err) - return - } - - var wipeCandidates func(remaining []string) error - wipeCandidates = func(remaining []string) error { - if len(remaining) == 0 { - return nil - } - if err := e.WipeRepoPolicies(rbac.ThisServer, remaining[0]); err != nil { - return fmt.Errorf("WipeRepoPolicies %s: %w", remaining[0], err) - } - return wipeCandidates(remaining[1:]) - } - if err := wipeCandidates(candidates); err != nil { - logger.Warn("legacy casbin rekey: wipe failed", "err", err) - return - } - - if _, err := d.ExecContext(ctx, `insert or ignore into migrations (name) values (?)`, flag); err != nil { - logger.Warn("legacy casbin rekey: mark flag failed", "err", err) - return - } - logger.Info("legacy casbin rekeyed", "owner", owner, "name", name, "rkey", rkey, "repoDid", repoDid, "candidates", candidates, "collabs", len(collabs)) -} diff --git a/spindle/ingester.go b/spindle/ingester.go --- a/spindle/ingester.go +++ b/spindle/ingester.go @@ -20,7 +20,7 @@ var err error switch e.Commit.Collection { - case tangled.RepoNSID, tangled.RepoCollaboratorNSID: + case tangled.RepoNSID: if evt, ok := jetstreamToTapEvent(e); ok { err = s.tap.processEvent(ctx, evt) } diff --git a/spindle/server.go b/spindle/server.go --- a/spindle/server.go +++ b/spindle/server.go @@ -29,7 +29,7 @@ kgit "tangled.org/core/knotserver/git" "tangled.org/core/log" "tangled.org/core/notifier" - "tangled.org/core/rbac" + "tangled.org/core/rbac/v2" "tangled.org/core/repoident" "tangled.org/core/repoverify" "tangled.org/core/spindle/config" @@ -46,10 +46,6 @@ //go:embed motd var defaultMotd []byte - -const ( - rbacDomain = "thisserver" -) type Spindle struct { jc *jetstream.JetstreamClient @@ -79,7 +75,7 @@ if err != nil { return nil, fmt.Errorf("failed to setup rbac enforcer: %w", err) } - e.E.EnableAutoSave(true) + e.EnableAutoSave(true) n := notifier.New() @@ -114,7 +110,6 @@ collections := []string{ tangled.RepoNSID, - tangled.RepoCollaboratorNSID, tangled.RepoPullNSID, } jc, err := jetstream.NewJetstreamClient(cfg.Server.JetstreamEndpoint, "spindle", collections, nil, log.SubLogger(logger, "jetstream"), d, true, true) @@ -235,6 +230,10 @@ // Enforcer returns the RBAC enforcer instance. func (s *Spindle) Enforcer() *rbac.Enforcer { return s.e +} + +func (s *Spindle) VerifyRepo(ctx context.Context, repo syntax.DID) (repoverify.Result, error) { + return s.verify(ctx, repoident.RepoDid(repo)) } // SetMotdContent sets custom MOTD content, replacing the embedded default. @@ -382,10 +381,11 @@ 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 { + switch msg.Nsid { + case knotdb.RepoCollaboratorUpdateNSID: return s.ingestKnotCollaborator(ctx, l, src, msg) - } - if msg.Nsid == tangled.GitRefUpdateNSID { + + case tangled.GitRefUpdateNSID: event := tangled.GitRefUpdate{} if err := json.Unmarshal(msg.EventJson, &event); err != nil { l.Error("error unmarshalling", "err", err) @@ -475,19 +475,13 @@ switch rec.Op { case knotdb.AclOpAdd: - if err := s.e.AddCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { + if err := s.e.AddRepoCollaborator(subject, repoDid); 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 { + if err := s.e.RemoveRepoCollaborator(subject, repoDid); 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: @@ -747,7 +741,7 @@ } } -func (s *Spindle) runJob(ctx context.Context, job *db.JobRow) { +func (s *Spindle) runJob(_ context.Context, job *db.JobRow) { pipelineId := models.PipelineId{ Knot: job.PipelineIdKnot, Rkey: job.PipelineIdRkey, diff --git a/spindle/startup_migrations.go b/spindle/startup_migrations.go --- a/spindle/startup_migrations.go +++ b/spindle/startup_migrations.go @@ -15,35 +15,12 @@ const forceTapResyncFlag = "force-tap-repo-resync-v1" func runStartupMigrations(ctx context.Context, d *db.DB, tapEmbed bool, tapDBPath string, logger *slog.Logger) error { - if err := cleanupOrphanRepos(ctx, d, logger); err != nil { - return fmt.Errorf("cleanup orphan repos: %w", err) - } if !tapEmbed { logger.Warn("tap not embedded: legacy repos won't auto-resync; trigger external tap resync to migrate secrets/casbin") return nil } if err := nudgeTapForResync(ctx, d, tapDBPath, logger); err != nil { return fmt.Errorf("nudge tap for resync: %w", err) - } - return nil -} - -func cleanupOrphanRepos(ctx context.Context, d *db.DB, logger *slog.Logger) error { - res, err := d.ExecContext(ctx, ` - delete from repos - where coalesce(repo_did, '') = '' - and exists ( - select 1 from repos r2 - where r2.owner = repos.owner - and coalesce(r2.repo_did, '') <> '' - ) - `) - if err != nil { - return fmt.Errorf("delete orphan repos: %w", err) - } - n, _ := res.RowsAffected() - if n > 0 { - logger.Info("cleaned up orphan repos missing repo_did", "deleted", n) } return nil } diff --git a/spindle/startup_migrations_test.go b/spindle/startup_migrations_test.go --- a/spindle/startup_migrations_test.go +++ b/spindle/startup_migrations_test.go @@ -3,7 +3,6 @@ import ( "context" "database/sql" - "fmt" "io" "log/slog" "path/filepath" @@ -12,7 +11,7 @@ "github.com/bluesky-social/indigo/atproto/syntax" - "tangled.org/core/rbac" + "tangled.org/core/rbac/v2" "tangled.org/core/spindle/db" "tangled.org/core/spindle/secrets" ) @@ -88,7 +87,7 @@ if err != nil { t.Fatalf("rbac.NewEnforcer: %v", err) } - e.E.EnableAutoSave(true) + e.EnableAutoSave(true) return d, e } @@ -112,18 +111,6 @@ }) if err != nil { t.Fatalf("AddSecret(%s/%s): %v", repo, key, err) - } -} - -func mustAddCollab(t *testing.T, d *db.DB, owner, rkey, subject, repoDid string) { - t.Helper() - if err := d.AddRepoCollaborator(db.RepoCollaborator{ - OwnerDid: syntax.DID(owner), - Rkey: syntax.RecordKey(rkey), - Subject: syntax.DID(subject), - RepoDid: syntax.DID(repoDid), - }); err != nil { - t.Fatalf("AddRepoCollaborator(%s): %v", rkey, err) } } @@ -367,183 +354,6 @@ } } -func TestMigrateLegacyRepoCasbin_NameCandidate(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, e := newTestSpindleDB(t) - - if err := e.AddSpindle(rbacDomain); err != nil { - t.Fatalf("AddSpindle: %v", err) - } - - owner := "did:plc:akshay" - repoDid := "did:plc:boltless" - displayName := "myrepo" - rkey := "3kspindlerkey00a" - collab := "did:plc:limpet" - oldNameKey := owner + "/" + displayName - oldRkeyKey := owner + "/" + rkey - - mustAddCollab(t, d, owner, "3kcollabrkey0001", collab, repoDid) - - if err := e.AddRepo(owner, rbacDomain, oldNameKey); err != nil { - t.Fatalf("seed AddRepo at Name key: %v", err) - } - if err := e.AddCollaborator(collab, rbacDomain, oldNameKey); err != nil { - t.Fatalf("seed AddCollaborator at Name key: %v", err) - } - - migrateLegacyRepoCasbin(ctx, d, e, logger, syntax.DID(owner), displayName, syntax.RecordKey(rkey), syntax.DID(repoDid)) - - if got, err := e.IsSettingsAllowed(owner, rbacDomain, repoDid); err != nil || !got { - t.Errorf("owner should have settings at new repoDid key, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(collab, rbacDomain, repoDid); err != nil || !got { - t.Errorf("collab should have settings at new repoDid key, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(owner, rbacDomain, oldNameKey); err != nil || got { - t.Errorf("owner Name-keyed policy should be removed, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(collab, rbacDomain, oldNameKey); err != nil || got { - t.Errorf("collab Name-keyed policy should be removed, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(owner, rbacDomain, oldRkeyKey); err != nil || got { - t.Errorf("owner rkey-keyed policy should be absent (never added), allowed=%v err=%v", got, err) - } - - migrateLegacyRepoCasbin(ctx, d, e, logger, syntax.DID(owner), displayName, syntax.RecordKey(rkey), syntax.DID(repoDid)) - - if got, err := e.IsSettingsAllowed(collab, rbacDomain, repoDid); err != nil || !got { - t.Errorf("collab settings still expected after idempotent re-run, allowed=%v err=%v", got, err) - } - - var marked int - if err := d.QueryRow( - `select count(*) from migrations where name = ?`, - "legacy-casbin-rekey:"+repoDid+":"+rkey, - ).Scan(&marked); err != nil { - t.Fatalf("query migrations: %v", err) - } - if marked != 1 { - t.Errorf("expected per-repo flag recorded exactly once, got %d", marked) - } -} - -func TestMigrateLegacyRepoCasbin_RkeyCandidate(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, e := newTestSpindleDB(t) - - if err := e.AddSpindle(rbacDomain); err != nil { - t.Fatalf("AddSpindle: %v", err) - } - - owner := "did:plc:akshay" - repoDid := "did:plc:boltless" - displayName := "myrepo" - rkey := "3kspindlerkey00a" - collab := "did:plc:limpet" - oldRkeyKey := owner + "/" + rkey - - mustAddCollab(t, d, owner, "3kcollabrkey0001", collab, repoDid) - - if err := e.AddRepo(owner, rbacDomain, oldRkeyKey); err != nil { - t.Fatalf("seed AddRepo at rkey: %v", err) - } - if err := e.AddCollaborator(collab, rbacDomain, oldRkeyKey); err != nil { - t.Fatalf("seed AddCollaborator at rkey: %v", err) - } - - migrateLegacyRepoCasbin(ctx, d, e, logger, syntax.DID(owner), displayName, syntax.RecordKey(rkey), syntax.DID(repoDid)) - - if got, err := e.IsSettingsAllowed(owner, rbacDomain, repoDid); err != nil || !got { - t.Errorf("owner should have settings at new repoDid key, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(collab, rbacDomain, repoDid); err != nil || !got { - t.Errorf("collab should have settings at new repoDid key, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(owner, rbacDomain, oldRkeyKey); err != nil || got { - t.Errorf("owner rkey-keyed policy should be removed, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(collab, rbacDomain, oldRkeyKey); err != nil || got { - t.Errorf("collab rkey-keyed policy should be removed, allowed=%v err=%v", got, err) - } -} - -func TestMigrateLegacyRepoCasbin_BothCandidates(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, e := newTestSpindleDB(t) - - if err := e.AddSpindle(rbacDomain); err != nil { - t.Fatalf("AddSpindle: %v", err) - } - - owner := "did:plc:akshay" - repoDid := "did:plc:boltless" - displayName := "myrepo" - rkey := "3kspindlerkey00a" - collab := "did:plc:limpet" - oldNameKey := owner + "/" + displayName - oldRkeyKey := owner + "/" + rkey - - mustAddCollab(t, d, owner, "3kcollabrkey0001", collab, repoDid) - - if err := e.AddRepo(owner, rbacDomain, oldNameKey); err != nil { - t.Fatalf("seed AddRepo at Name key: %v", err) - } - if err := e.AddRepo(owner, rbacDomain, oldRkeyKey); err != nil { - t.Fatalf("seed AddRepo at rkey: %v", err) - } - if err := e.AddCollaborator(collab, rbacDomain, oldNameKey); err != nil { - t.Fatalf("seed AddCollaborator at Name key: %v", err) - } - if err := e.AddCollaborator(collab, rbacDomain, oldRkeyKey); err != nil { - t.Fatalf("seed AddCollaborator at rkey: %v", err) - } - - migrateLegacyRepoCasbin(ctx, d, e, logger, syntax.DID(owner), displayName, syntax.RecordKey(rkey), syntax.DID(repoDid)) - - if got, err := e.IsSettingsAllowed(owner, rbacDomain, oldNameKey); err != nil || got { - t.Errorf("owner Name-keyed policy should be removed, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(owner, rbacDomain, oldRkeyKey); err != nil || got { - t.Errorf("owner rkey-keyed policy should be removed, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(collab, rbacDomain, oldNameKey); err != nil || got { - t.Errorf("collab Name-keyed policy should be removed, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(collab, rbacDomain, oldRkeyKey); err != nil || got { - t.Errorf("collab rkey-keyed policy should be removed, allowed=%v err=%v", got, err) - } -} - -func TestMigrateLegacyRepoCasbin_BothEmpty(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, e := newTestSpindleDB(t) - - if err := e.AddSpindle(rbacDomain); err != nil { - t.Fatalf("AddSpindle: %v", err) - } - - owner := syntax.DID("did:plc:akshay") - repoDid := syntax.DID("did:plc:boltless") - - migrateLegacyRepoCasbin(ctx, d, e, logger, owner, "", "", repoDid) - - var marked int - if err := d.QueryRow( - `select count(*) from migrations where name like ?`, - "legacy-casbin-rekey:"+repoDid.String()+":%", - ).Scan(&marked); err != nil { - t.Fatalf("query migrations: %v", err) - } - if marked != 0 { - t.Errorf("empty inputs should not record flag, got %d", marked) - } -} - func TestNudgeTapForResync(t *testing.T) { ctx := context.Background() logger := slog.New(slog.NewTextHandler(io.Discard, nil)) @@ -678,279 +488,5 @@ } if marked != 0 { t.Errorf("non-embed mode should skip tap nudge flag, got %d", marked) - } -} - -func TestCleanupOrphanRepos_DeletesWhenSiblingExists(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, _ := newTestSpindleDB(t) - - owner := "did:plc:akshay" - if _, err := d.Exec(`insert into repos (knot, owner, rkey, repo_did, created_at) values - ('k', ?, 'legacy_name', null, null), - ('k', ?, '3kspindlerkey00a', 'did:plc:boltless', '2024-01-01T00:00:00Z')`, - owner, owner); err != nil { - t.Fatalf("seed: %v", err) - } - - if err := cleanupOrphanRepos(ctx, d, logger); err != nil { - t.Fatalf("cleanupOrphanRepos: %v", err) - } - - var nullCount int - if err := d.QueryRow(`select count(*) from repos where repo_did is null`).Scan(&nullCount); err != nil { - t.Fatalf("null count: %v", err) - } - if nullCount != 0 { - t.Errorf("orphan should be deleted when sibling exists, got %d remaining", nullCount) - } - - var sibCount int - if err := d.QueryRow(`select count(*) from repos where repo_did is not null`).Scan(&sibCount); err != nil { - t.Fatalf("sibling count: %v", err) - } - if sibCount != 1 { - t.Errorf("sibling row should be preserved, got %d", sibCount) - } -} - -func TestCleanupOrphanRepos_KeepsWhenAlone(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, _ := newTestSpindleDB(t) - - owner := "did:plc:akshay" - if _, err := d.Exec(`insert into repos (knot, owner, rkey, repo_did, created_at) values - ('k', ?, 'legacy_name', null, null)`, owner); err != nil { - t.Fatalf("seed: %v", err) - } - - if err := cleanupOrphanRepos(ctx, d, logger); err != nil { - t.Fatalf("cleanupOrphanRepos: %v", err) - } - - var remaining int - if err := d.QueryRow(`select count(*) from repos where owner = ?`, owner).Scan(&remaining); err != nil { - t.Fatalf("count: %v", err) - } - if remaining != 1 { - t.Errorf("orphan with no sibling should be kept (preserves owner registration), got %d", remaining) - } -} - -func TestCleanupOrphanRepos_PerOwnerScope(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, _ := newTestSpindleDB(t) - - ownerA := "did:plc:akshay" - ownerB := "did:plc:limpet" - if _, err := d.Exec(`insert into repos (knot, owner, rkey, repo_did, created_at) values - ('k', ?, 'legacy_a', null, null), - ('k', ?, '3krealkey', 'did:plc:boltless', '2024-01-01T00:00:00Z'), - ('k', ?, 'legacy_b', null, null)`, - ownerA, ownerA, ownerB); err != nil { - t.Fatalf("seed: %v", err) - } - - if err := cleanupOrphanRepos(ctx, d, logger); err != nil { - t.Fatalf("cleanupOrphanRepos: %v", err) - } - - var ownerARows, ownerBRows int - if err := d.QueryRow(`select count(*) from repos where owner = ?`, ownerA).Scan(&ownerARows); err != nil { - t.Fatalf("count A: %v", err) - } - if ownerARows != 1 { - t.Errorf("ownerA: orphan should be deleted (sibling exists), expected 1 row, got %d", ownerARows) - } - if err := d.QueryRow(`select count(*) from repos where owner = ?`, ownerB).Scan(&ownerBRows); err != nil { - t.Fatalf("count B: %v", err) - } - if ownerBRows != 1 { - t.Errorf("ownerB: orphan should be kept (no sibling), expected 1 row, got %d", ownerBRows) - } -} - -func TestCleanupOrphanRepos_EmptyStringRepoDid(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, _ := newTestSpindleDB(t) - - owner := "did:plc:akshay" - if _, err := d.Exec(`insert into repos (knot, owner, rkey, repo_did, created_at) values - ('k', ?, 'legacy_empty', '', null), - ('k', ?, '3krealkey', 'did:plc:boltless', '2024-01-01T00:00:00Z')`, - owner, owner); err != nil { - t.Fatalf("seed: %v", err) - } - - if err := cleanupOrphanRepos(ctx, d, logger); err != nil { - t.Fatalf("cleanupOrphanRepos: %v", err) - } - - var emptyCount int - if err := d.QueryRow(`select count(*) from repos where coalesce(repo_did, '') = ''`).Scan(&emptyCount); err != nil { - t.Fatalf("empty count: %v", err) - } - if emptyCount != 0 { - t.Errorf("empty-string repo_did orphan should be deleted when sibling exists, got %d remaining", emptyCount) - } -} - -func TestMigrateLegacyRepoCasbin_MultipleCollabsAllRekeyed(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, e := newTestSpindleDB(t) - - if err := e.AddSpindle(rbacDomain); err != nil { - t.Fatalf("AddSpindle: %v", err) - } - - owner := "did:plc:akshay" - repoDid := "did:plc:boltless" - displayName := "myrepo" - rkey := "3kspindlerkey00a" - oldNameKey := owner + "/" + displayName - collabs := []string{"did:plc:limpet", "did:plc:nautilus", "did:plc:whelk", "did:plc:cuttle"} - - var addCollabRows func(rest []string, idx int) - addCollabRows = func(rest []string, idx int) { - if len(rest) == 0 { - return - } - mustAddCollab(t, d, owner, fmt.Sprintf("3kcollabrkey%04d", idx), rest[0], repoDid) - addCollabRows(rest[1:], idx+1) - } - addCollabRows(collabs, 0) - - if err := e.AddRepo(owner, rbacDomain, oldNameKey); err != nil { - t.Fatalf("seed owner: %v", err) - } - var seedAll func(rest []string) error - seedAll = func(rest []string) error { - if len(rest) == 0 { - return nil - } - if err := e.AddCollaborator(rest[0], rbacDomain, oldNameKey); err != nil { - return err - } - return seedAll(rest[1:]) - } - if err := seedAll(collabs); err != nil { - t.Fatalf("seed collab policies: %v", err) - } - - migrateLegacyRepoCasbin(ctx, d, e, logger, syntax.DID(owner), displayName, syntax.RecordKey(rkey), syntax.DID(repoDid)) - - var assertEach func(rest []string) - assertEach = func(rest []string) { - if len(rest) == 0 { - return - } - c := rest[0] - if got, err := e.IsSettingsAllowed(c, rbacDomain, repoDid); err != nil || !got { - t.Errorf("collab %s should have settings at repoDid, allowed=%v err=%v", c, got, err) - } - if got, err := e.IsPushAllowed(c, rbacDomain, repoDid); err != nil || !got { - t.Errorf("collab %s should have push at repoDid, allowed=%v err=%v", c, got, err) - } - if got, err := e.IsSettingsAllowed(c, rbacDomain, oldNameKey); err != nil || got { - t.Errorf("collab %s old policy should be wiped, allowed=%v err=%v", c, got, err) - } - assertEach(rest[1:]) - } - assertEach(collabs) -} - -func TestMigrateLegacyRepoCasbin_RenameSiblingsEachWiped(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, e := newTestSpindleDB(t) - - if err := e.AddSpindle(rbacDomain); err != nil { - t.Fatalf("AddSpindle: %v", err) - } - - owner := syntax.DID("did:plc:akshay") - repoDid := syntax.DID("did:plc:di4gol2smljyj6gjnjdu5qrg") - siblings := []string{"pre-rename-life", "i-renamed-this", "post-rename-rename", "post-rename-renamed-again"} - - var seedAll func(rest []string) error - seedAll = func(rest []string) error { - if len(rest) == 0 { - return nil - } - if err := e.AddRepo(owner.String(), rbacDomain, owner.String()+"/"+rest[0]); err != nil { - return err - } - return seedAll(rest[1:]) - } - if err := seedAll(siblings); err != nil { - t.Fatalf("seed siblings: %v", err) - } - - var run func(rest []string) - run = func(rest []string) { - if len(rest) == 0 { - return - } - migrateLegacyRepoCasbin(ctx, d, e, logger, owner, "", syntax.RecordKey(rest[0]), repoDid) - run(rest[1:]) - } - run(siblings) - - var assertWiped func(rest []string) - assertWiped = func(rest []string) { - if len(rest) == 0 { - return - } - key := owner.String() + "/" + rest[0] - if got, err := e.IsSettingsAllowed(owner.String(), rbacDomain, key); err != nil || got { - t.Errorf("rename sibling %s should be wiped, allowed=%v err=%v", rest[0], got, err) - } - assertWiped(rest[1:]) - } - assertWiped(siblings) - - if got, err := e.IsSettingsAllowed(owner.String(), rbacDomain, repoDid.String()); err != nil || !got { - t.Errorf("owner should retain settings at repoDid, allowed=%v err=%v", got, err) - } -} - -func TestMigrateLegacyRepoCasbin_StrandedCollabWiped(t *testing.T) { - ctx := context.Background() - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - d, e := newTestSpindleDB(t) - - if err := e.AddSpindle(rbacDomain); err != nil { - t.Fatalf("AddSpindle: %v", err) - } - - owner := "did:plc:akshay" - repoDid := "did:plc:boltless" - displayName := "myrepo" - rkey := "3kspindlerkey00a" - strandedCollab := "did:plc:nautilus" - oldNameKey := owner + "/" + displayName - - if err := e.AddRepo(owner, rbacDomain, oldNameKey); err != nil { - t.Fatalf("seed AddRepo at Name key: %v", err) - } - if err := e.AddCollaborator(strandedCollab, rbacDomain, oldNameKey); err != nil { - t.Fatalf("seed stranded collab at Name key: %v", err) - } - - migrateLegacyRepoCasbin(ctx, d, e, logger, syntax.DID(owner), displayName, syntax.RecordKey(rkey), syntax.DID(repoDid)) - - if got, err := e.IsSettingsAllowed(strandedCollab, rbacDomain, oldNameKey); err != nil || got { - t.Errorf("stranded collab should be wiped from old key, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(owner, rbacDomain, oldNameKey); err != nil || got { - t.Errorf("owner old policy should be wiped, allowed=%v err=%v", got, err) - } - if got, err := e.IsSettingsAllowed(owner, rbacDomain, repoDid); err != nil || !got { - t.Errorf("owner should have settings at new repoDid key, allowed=%v err=%v", got, err) } } diff --git a/spindle/tapclient.go b/spindle/tapclient.go --- a/spindle/tapclient.go +++ b/spindle/tapclient.go @@ -2,14 +2,13 @@ import ( "context" - "database/sql" "encoding/json" - "errors" "fmt" "log/slog" "net/http" "net/url" - "sync" + "slices" + "strings" "time" "github.com/bluesky-social/indigo/atproto/syntax" @@ -18,7 +17,6 @@ avmodels "tangled.org/core/appview/models" "tangled.org/core/eventconsumer" "tangled.org/core/log" - "tangled.org/core/rbac" "tangled.org/core/spindle/db" "tangled.org/core/spindle/git" "tangled.org/core/spindle/models" @@ -28,29 +26,21 @@ ) const ( - maxPendingPerRepo = 64 - pendingCollabTTL = 10 * time.Minute + collaboratorPageLimit = 100 + maxCollaboratorPages = 50 ) -type pendingCollabEvent struct { - evt *tapc.RecordEventData - at time.Time -} - type Tap struct { - logger *slog.Logger - spindle *Spindle - tap tapc.Client - pendingMu sync.Mutex - pendingCollabs map[syntax.DID][]pendingCollabEvent + logger *slog.Logger + spindle *Spindle + tap tapc.Client } func NewTapClient(s *Spindle) *Tap { return &Tap{ - logger: log.SubLogger(s.l, "tapclient"), - spindle: s, - tap: tapc.NewClient(s.cfg.Server.Tap.Url, s.cfg.Server.Tap.AdminPassword), - pendingCollabs: make(map[syntax.DID][]pendingCollabEvent), + logger: log.SubLogger(s.l, "tapclient"), + spindle: s, + tap: tapc.NewClient(s.cfg.Server.Tap.Url, s.cfg.Server.Tap.AdminPassword), } } @@ -66,7 +56,6 @@ EventHandler: t.processEvent, ConnectHandler: t.onConnect, }) - go t.purgePendingCollabsLoop(t.spindle.rootCtx) } func (t *Tap) onConnect(ctx context.Context) { @@ -80,12 +69,11 @@ switch evt.Record.Collection.String() { case tangled.RepoNSID: return t.processRepo(ctx, evt.Record) - case tangled.RepoCollaboratorNSID: - return t.processCollaborator(ctx, evt.Record) } return nil } +// processRepo ingests repo declaration record: `sh.tangled.repo`. It skips alias records. func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error { l := t.logger.With("collection", tangled.RepoNSID, "did", evt.Did, "rkey", evt.Rkey) @@ -97,18 +85,6 @@ record := tangled.Repo{} if err := json.Unmarshal(evt.Record, &record); err != nil { l.Warn("skipping invalid repo record", "err", err) - return nil - } - - hostname := t.spindle.cfg.Server.Hostname - prior, priorErr := t.spindle.db.GetRepoByOwnerRkey(ownerDid, rkey) - knownRepo := priorErr == nil - - if record.Spindle == nil || *record.Spindle != hostname { - if knownRepo { - l.Info("tearing down repo reassigned from this spindle", "newSpindle", record.Spindle) - return t.teardownRepo(l, prior, ownerDid, rkey) - } return nil } @@ -131,24 +107,40 @@ return nil } - // check if this repo DID is already owned by someone else - existingRepo, err := t.spindle.db.GetRepoByDid(repoDid) - if err == nil { - if existingRepo.Owner != ownerDid { - l.Warn("rejecting repo record: repoDid already registered by another owner", "repoDid", repoDid, "existingOwner", existingRepo.Owner, "newOwner", ownerDid) + // ignore repos not pointing this spindle. + hostname := t.spindle.cfg.Server.Hostname + if record.Spindle == nil || *record.Spindle != hostname { + // teardown existing repo + prior, err := t.spindle.db.GetRepoByDid(repoDid) + if err != nil { return nil } - } else if !errors.Is(err, sql.ErrNoRows) { - return fmt.Errorf("lookup existing repo by DID: %w", err) + if prior.Owner == ownerDid && prior.Rkey == rkey { + l.Info("tearing down repo reassigned from this spindle", "newSpindle", record.Spindle) + return t.teardownRepo(l, prior.RepoDid) + } + l.Warn("ignoring reassignment from non-registering record", "owner", prior.Owner, "rkey", prior.Rkey) + return nil } - if err := t.spindle.e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()); err != nil { + // verify repo declaration + // NOTE: we are ignoring repo declaration records that *might* become correct pointer in + // future with repository rename or ownership transfer. On repo rename, new pointer record + // should always be recreated even when it already exists. + verified, err := t.verifyRepoDeclaration(ctx, l, repoDid, ownerDid, rkey.String(), record.Knot) + if err != nil { + l.Warn("failed to verify repo declaration", "err", err) + return nil + } + if !verified { + // ignore alias records + return nil + } + + if err := t.spindle.e.SetRepoOwner(ownerDid, repoDid); err != nil { l.Error("failed to add repo policy", "err", err) return fmt.Errorf("add repo policy: %w", err) } - - src := eventconsumer.NewKnotSource(record.Knot) - t.spindle.ks.AddSource(t.spindle.rootCtx, src) repo := db.Repo{ Knot: record.Knot, @@ -163,6 +155,9 @@ return fmt.Errorf("add repo: %w", err) } + t.reconcileCollaborators(ctx, l, record.Knot, repoDid, ownerDid) + t.spindle.ks.AddSource(t.spindle.rootCtx, eventconsumer.NewKnotSource(record.Knot)) + // setup sparse sync repoCloneUri := t.spindle.newRepoCloneUrl(repo.Knot, repo.RepoDid) repoPath := t.spindle.newRepoPath(repo.RepoDid) @@ -175,7 +170,6 @@ legacyName = *record.Name } migrateLegacyRepoSecrets(ctx, t.spindle.db, t.spindle.vault, l, ownerDid, legacyName, rkey, repoDid) - migrateLegacyRepoCasbin(ctx, t.spindle.db, t.spindle.e, l, ownerDid, legacyName, rkey, repoDid) if removed, err := t.spindle.db.CollapseRepoSiblings(ownerDid, repoDid); err != nil { l.Warn("collapse rename siblings failed", "err", err) @@ -190,146 +184,135 @@ } t.spindle.jc.AddDid(ownerDid.String()) - t.drainPendingCollabs(ctx, repoDid) - case tapc.RecordDeleteAction: - repo, err := t.spindle.db.GetRepoByOwnerRkey(ownerDid, rkey) + repoDid, err := t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey) if err != nil { - l.Info("skipping delete for unknown repo") + return fmt.Errorf("deleting repo record: %w", err) + } + + if repoDid == "" { + // no record is deleted. record was not pointing this spindle or was an alias record return nil } - return t.teardownRepo(l, repo, ownerDid, rkey) + return t.teardownRepo(l, repoDid) } return nil } -func (t *Tap) teardownRepo(l *slog.Logger, repo *db.Repo, ownerDid syntax.DID, rkey syntax.RecordKey) error { - if repo.RepoDid != "" { - collabs, err := t.spindle.db.ListCollaboratorsByRepoDid(repo.RepoDid) - if err != nil { - l.Error("failed to list collaborators for cleanup", "err", err) - return fmt.Errorf("list collaborators: %w", err) - } - for _, c := range collabs { - if err := t.spindle.e.RemoveCollaborator(c.Subject.String(), rbac.ThisServer, repo.RepoDid.String()); err != nil { - l.Error("failed to remove collaborator policy", "subject", c.Subject, "err", err) - return fmt.Errorf("remove collaborator policy: %w", err) - } - } - if err := t.spindle.db.DeleteRepoCollaboratorsByRepoDid(repo.RepoDid); err != nil { - l.Error("failed to clear collaborator rows", "err", err) - return err - } - if err := t.spindle.e.RemoveRepo(ownerDid.String(), rbac.ThisServer, repo.RepoDid.String()); err != nil { - l.Error("failed to remove repo policy", "err", err) - return fmt.Errorf("remove repo policy: %w", err) - } +func (t *Tap) verifyRepoDeclaration(ctx context.Context, l *slog.Logger, repo, owner syntax.DID, rkey, knot string) (bool, error) { + result, err := t.spindle.VerifyRepo(ctx, repo) + if err != nil { + return false, err } - if err := t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey); err != nil { - l.Error("failed to delete repo row", "err", err) - return fmt.Errorf("delete repo row: %w", err) + l = l.With( + "repoDid", result.RepoDid, + "repoOwner", result.OwnerDid, + "repoKnot", result.KnotURL.Host, + "claimedOwner", owner, + "claimedRkey", rkey, + "claimedKnot", knot, + ) + if syntax.DID(result.OwnerDid) != owner { + l.Warn("rejecting repo event: owner mismatch") + return false, nil + } + if result.Rkey != rkey { + l.Warn("rejecting repo event: rkey mismatch") + return false, nil + } + if !strings.EqualFold(knot, result.KnotURL.Host) { + l.Warn("rejecting repo event: record knot does not match DID-doc endpoint") + return false, nil + } + return true, nil +} + +func (t *Tap) reconcileCollaborators(ctx context.Context, l *slog.Logger, knot string, repo, owner syntax.DID) { + wanted, err := t.fetchKnotCollaborators(ctx, knot, repo) + if err != nil { + l.Warn("collaborator reconcile: failed to fetch roster from knot", "knot", knot, "err", err) + return + } + + have, err := t.spindle.e.GetRepoCollaborators(repo) + if err != nil { + l.Warn("collaborator reconcile: failed to read current grants", "err", err) + return + } + + // the owner is a collaborator by role inheritance, never by an explicit grant + delete(wanted, owner) + + for did := range wanted { + if slices.Contains(have, did) { + continue + } + if err := t.spindle.grantCollaborator(did, repo); err != nil { + l.Error("collaborator reconcile: failed to add", "subject", did, "err", err) + return + } + l.Info("collaborator reconcile: added", "subject", did) + } + for _, did := range have { + if did == owner { + continue + } + if _, ok := wanted[did]; ok { + continue + } + if err := t.spindle.e.RemoveRepoCollaborator(did, repo); err != nil { + l.Error("collaborator reconcile: failed to remove", "subject", did, "err", err) + return + } + l.Info("collaborator reconcile: removed", "subject", did) + } +} + +func (t *Tap) fetchKnotCollaborators(ctx context.Context, knot string, repo syntax.DID) (map[syntax.DID]struct{}, error) { + scheme := "https" + if t.spindle.cfg.Server.Dev { + scheme = "http" + } + xc := &indigoxrpc.Client{ + Host: fmt.Sprintf("%s://%s", scheme, knot), + Client: &http.Client{Timeout: 30 * time.Second}, + } + + subjects := make(map[syntax.DID]struct{}) + cursor := "" + for range maxCollaboratorPages { + out, err := tangled.RepoListCollaborators(ctx, xc, cursor, collaboratorPageLimit, "", repo.String()) + if err != nil { + return nil, err + } + for _, item := range out.Items { + did, err := syntax.ParseDID(item.Subject) + if err != nil { + continue + } + subjects[did] = struct{}{} + } + if out.Cursor == nil || *out.Cursor == "" { + return subjects, nil + } + cursor = *out.Cursor + } + return nil, fmt.Errorf("collaborator roster exceeded %d pages", maxCollaboratorPages) +} + +func (t *Tap) teardownRepo(l *slog.Logger, repo syntax.DID) error { + if repo == "" { + return nil + } + if err := t.spindle.db.DeleteRepo(repo); err != nil { + l.Error("failed to remove repo", "err", err) + return fmt.Errorf("remove repo: %w", err) + } + if err := t.spindle.e.DeleteRepo(repo); err != nil { + l.Error("failed to remove repo policy", "err", err) + return fmt.Errorf("remove repo policy: %w", err) } // TODO: clear sparse-synced git repo - return nil -} - -func (t *Tap) processCollaborator(ctx context.Context, evt *tapc.RecordEventData) error { - l := t.logger.With("collection", tangled.RepoCollaboratorNSID, "did", evt.Did, "rkey", evt.Rkey) - - switch evt.Action { - case tapc.RecordCreateAction, tapc.RecordUpdateAction: - record := tangled.RepoCollaborator{} - if err := json.Unmarshal(evt.Record, &record); err != nil { - l.Warn("skipping invalid collaborator record", "err", err) - return nil - } - - actor := evt.Did - rkey := evt.Rkey - - subjectDid, err := syntax.ParseDID(record.Subject) - if err != nil { - l.Info("skipping collaborator with malformed subject DID", "subject", record.Subject, "err", err) - return nil - } - if _, err := t.spindle.res.ResolveIdent(ctx, subjectDid.String()); err != nil { - l.Info("skipping unresolvable collaborator subject", "subject", subjectDid, "err", err) - return nil - } - - repoRefDid, err := syntax.ParseDID(record.Repo) - if err != nil { - l.Info("skipping collaborator with non-DID repo ref", "repo", record.Repo, "err", err) - return nil - } - repo, lookupErr := t.spindle.db.GetRepoByDid(repoRefDid) - if errors.Is(lookupErr, sql.ErrNoRows) { - t.bufferCollab(repoRefDid, evt) - l.Info("buffering collaborator until repo arrives", "repo", repoRefDid) - return nil - } - if lookupErr != nil { - return fmt.Errorf("lookup repo %s: %w", repoRefDid, lookupErr) - } - repoDid := repo.RepoDid - ownerDid := repo.Owner - - if actor != ownerDid { - l.Info("rejecting collaborator with non-owner actor", "actor", actor, "owner", ownerDid) - return nil - } - - ok, err := t.spindle.e.IsCollaboratorInviteAllowed(ownerDid.String(), rbac.ThisServer, repoDid.String()) - if err != nil { - l.Error("invite permission check failed", "err", err) - return fmt.Errorf("invite check: %w", err) - } - if !ok { - l.Info("rejecting collaborator invite", "owner", ownerDid, "repo", repoDid) - return nil - } - - prior, priorErr := t.spindle.db.GetRepoCollaborator(actor, rkey) - staleSubject := priorErr == nil && (prior.Subject != subjectDid || prior.RepoDid != repoDid) - - if err := t.spindle.e.AddCollaborator(subjectDid.String(), rbac.ThisServer, repoDid.String()); err != nil { - l.Error("failed to add collaborator policy", "err", err) - return fmt.Errorf("add collaborator policy: %w", err) - } - if staleSubject { - if err := t.spindle.e.RemoveCollaborator(prior.Subject.String(), rbac.ThisServer, prior.RepoDid.String()); err != nil { - l.Error("failed to remove stale collaborator policy", "err", err) - return fmt.Errorf("remove stale collaborator: %w", err) - } - } - if err := t.spindle.db.AddRepoCollaborator(db.RepoCollaborator{ - OwnerDid: actor, - Rkey: rkey, - Subject: subjectDid, - RepoDid: repoDid, - }); err != nil { - l.Error("failed to persist collaborator row", "err", err) - return fmt.Errorf("track collaborator: %w", err) - } - - case tapc.RecordDeleteAction: - actor := evt.Did - rkey := evt.Rkey - - tracked, err := t.spindle.db.GetRepoCollaborator(actor, rkey) - if err != nil { - l.Info("skipping delete for unknown collaborator record") - return nil - } - if err := t.spindle.e.RemoveCollaborator(tracked.Subject.String(), rbac.ThisServer, tracked.RepoDid.String()); err != nil { - l.Error("failed to remove collaborator policy", "err", err) - return fmt.Errorf("remove collaborator policy: %w", err) - } - if err := t.spindle.db.DeleteRepoCollaborator(actor, rkey); err != nil { - l.Error("failed to delete collaborator row", "err", err) - return fmt.Errorf("delete collaborator row: %w", err) - } - } return nil } @@ -375,8 +358,8 @@ return nil } - // check if pull record author has push access to target repo - allowed, err := s.e.IsPushAllowed(evt.Did.String(), rbac.ThisServer, repo.RepoDid.String()) + // check if pull record author can trigger CI in target repo + allowed, err := s.e.IsRepoCiTriggerAllowed(evt.Did, repo.RepoDid) if err != nil { return fmt.Errorf("checking push access for pull record author: %w", err) } @@ -475,74 +458,6 @@ // no-op } return nil -} - -func (t *Tap) bufferCollab(repoDid syntax.DID, evt *tapc.RecordEventData) { - t.pendingMu.Lock() - defer t.pendingMu.Unlock() - list := t.pendingCollabs[repoDid] - list = append(list, pendingCollabEvent{evt: evt, at: time.Now()}) - if len(list) > maxPendingPerRepo { - list = list[len(list)-maxPendingPerRepo:] - } - t.pendingCollabs[repoDid] = list -} - -func (t *Tap) drainPendingCollabs(ctx context.Context, repoDid syntax.DID) { - t.pendingMu.Lock() - list := t.pendingCollabs[repoDid] - delete(t.pendingCollabs, repoDid) - t.pendingMu.Unlock() - if len(list) == 0 { - return - } - cutoff := time.Now().Add(-pendingCollabTTL) - for _, p := range list { - if p.at.Before(cutoff) { - continue - } - if err := t.processCollaborator(ctx, p.evt); err != nil { - t.logger.Warn("replaying buffered collaborator failed", "repo", repoDid, "rkey", p.evt.Rkey, "err", err) - } - } -} - -func (t *Tap) purgePendingCollabsLoop(ctx context.Context) { - ticker := time.NewTicker(pendingCollabTTL / 2) - defer ticker.Stop() - for { - select { - case <-ctx.Done(): - return - case <-ticker.C: - t.purgeStalePendingCollabs() - } - } -} - -func (t *Tap) purgeStalePendingCollabs() { - cutoff := time.Now().Add(-pendingCollabTTL) - t.pendingMu.Lock() - defer t.pendingMu.Unlock() - expired := 0 - for did, list := range t.pendingCollabs { - kept := list[:0] - for _, p := range list { - if !p.at.Before(cutoff) { - kept = append(kept, p) - } else { - expired++ - } - } - if len(kept) == 0 { - delete(t.pendingCollabs, did) - } else { - t.pendingCollabs[did] = kept - } - } - if expired > 0 { - t.logger.Warn("expired buffered collaborator events without matching repo arrival", "count", expired, "ttl", pendingCollabTTL) - } } func (s *Spindle) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*avmodels.PullSubmission, error) { diff --git a/spindle/db/collaborators.go b/spindle/db/collaborators.go deleted file mode 100644 --- a/spindle/db/collaborators.go +++ /dev/null @@ -1,115 +0,0 @@ -package db - -import ( - "database/sql" - "fmt" - - "github.com/bluesky-social/indigo/atproto/syntax" -) - -type RepoCollaborator struct { - OwnerDid syntax.DID - Rkey syntax.RecordKey - Subject syntax.DID - RepoDid syntax.DID -} - -func (d *DB) AddRepoCollaborator(c RepoCollaborator) error { - _, err := d.Exec( - `insert into repo_collaborators (owner_did, rkey, subject, repo_did) - values (?, ?, ?, ?) - on conflict(owner_did, rkey) do update set - subject = excluded.subject, - repo_did = excluded.repo_did`, - c.OwnerDid.String(), c.Rkey.String(), c.Subject.String(), c.RepoDid.String(), - ) - 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 { - return nil, err - } - return &RepoCollaborator{ - OwnerDid: syntax.DID(owner), - Rkey: syntax.RecordKey(rkey), - Subject: syntax.DID(subject), - RepoDid: syntax.DID(repoDid), - }, nil -} - -func (d *DB) GetRepoCollaborator(ownerDid syntax.DID, rkey syntax.RecordKey) (*RepoCollaborator, error) { - return scanCollab(d.QueryRow( - `select owner_did, rkey, subject, repo_did from repo_collaborators where owner_did = ? and rkey = ?`, - ownerDid.String(), rkey.String(), - )) -} - -func (d *DB) DeleteRepoCollaborator(ownerDid syntax.DID, rkey syntax.RecordKey) error { - res, err := d.Exec(`delete from repo_collaborators where owner_did = ? and rkey = ?`, ownerDid.String(), rkey.String()) - if err != nil { - return err - } - n, err := res.RowsAffected() - if err != nil { - return err - } - if n == 0 { - return sql.ErrNoRows - } - 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 { - return fmt.Errorf("delete collaborators for %s: %w", repoDid, err) - } - return nil -} - -func (d *DB) ListCollaboratorsByRepoDid(repoDid syntax.DID) ([]RepoCollaborator, error) { - rows, err := d.Query( - `select owner_did, rkey, subject, repo_did from repo_collaborators where repo_did = ?`, - repoDid.String(), - ) - if err != nil { - return nil, fmt.Errorf("list 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/collaborators_test.go b/spindle/db/collaborators_test.go deleted file mode 100644 --- a/spindle/db/collaborators_test.go +++ /dev/null @@ -1,94 +0,0 @@ -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/db/db.go b/spindle/db/db.go --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -287,6 +287,48 @@ return err } + // we use casbin instead + if err := orm.RunMigration(conn, logger, "drop-legacy-acl-tables", func(tx *sql.Tx) error { + _, err := tx.Exec(` + drop table if exists repo_collaborators; + drop table if exists known_dids; + `) + return err + }); err != nil { + return err + } + + // repo_did is required + if err := orm.RunMigration(conn, logger, "enforce-repo_did", func(tx *sql.Tx) error { + _, err := tx.Exec(` + create table repos_new ( + repo_did text primary key, + knot text not null, + owner text not null, + rkey text not null, + created_at text, + addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + + unique(owner, rkey) + ); + insert into repos_new (repo_did, knot, owner, rkey, created_at, addedAt) + select repo_did, knot, owner, rkey, created_at, addedAt from repos + where coalesce(repo_did, '') <> '' + and rowid = ( + select r2.rowid from repos r2 + where r2.repo_did = repos.repo_did + order by r2.created_at is null, r2.created_at desc, r2.rowid desc + limit 1 + ); + + drop table repos; + alter table repos_new rename to repos; + `) + return err + }); err != nil { + return err + } + return nil } diff --git a/spindle/db/repos.go b/spindle/db/repos.go --- a/spindle/db/repos.go +++ b/spindle/db/repos.go @@ -2,6 +2,7 @@ import ( "database/sql" + "errors" "github.com/bluesky-social/indigo/atproto/syntax" ) @@ -122,20 +123,20 @@ func (d *DB) GetRepoByDid(repoDid syntax.DID) (*Repo, error) { return scanRepo(d.QueryRow( - `select knot, owner, rkey, coalesce(repo_did, '') from repos where repo_did = ?`, + `select knot, owner, rkey, repo_did from repos where repo_did = ?`, repoDid.String(), )) } func (d *DB) GetRepoByOwnerRkey(owner syntax.DID, rkey syntax.RecordKey) (*Repo, error) { return scanRepo(d.QueryRow( - `select knot, owner, rkey, coalesce(repo_did, '') from repos where owner = ? and rkey = ?`, + `select knot, owner, rkey, repo_did from repos where owner = ? and rkey = ?`, owner.String(), rkey.String(), )) } func (d *DB) AllRepos() ([]Repo, error) { - rows, err := d.Query(`select knot, owner, rkey, coalesce(repo_did, '') from repos`) + rows, err := d.Query(`select knot, owner, rkey, repo_did from repos`) if err != nil { return nil, err } @@ -157,7 +158,23 @@ return repos, nil } -func (d *DB) DeleteRepoByOwnerRkey(owner syntax.DID, rkey syntax.RecordKey) error { - _, err := d.Exec(`delete from repos where owner = ? and rkey = ?`, owner.String(), rkey.String()) +func (d *DB) DeleteRepo(repoDid syntax.DID) error { + _, err := d.Exec(`delete from repos where repo_did = ?`, repoDid) return err +} + +// DeleteRepoByOwnerRkey deletes a repo by (owner,rkey) pair and returns deleted repos DID. +func (d *DB) DeleteRepoByOwnerRkey(owner syntax.DID, rkey syntax.RecordKey) (syntax.DID, error) { + var repoDid string + err := d.QueryRow( + `delete from repos where owner = ? and rkey = ? returning repo_did`, + owner.String(), rkey.String(), + ).Scan(&repoDid) + if errors.Is(err, sql.ErrNoRows) { + return "", nil + } + if err != nil { + return "", err + } + return syntax.DID(repoDid), nil } diff --git a/spindle/xrpc/add_secret.go b/spindle/xrpc/add_secret.go --- a/spindle/xrpc/add_secret.go +++ b/spindle/xrpc/add_secret.go @@ -10,7 +10,6 @@ "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" - "tangled.org/core/rbac" "tangled.org/core/spindle/secrets" xrpcerr "tangled.org/core/xrpc/errors" ) @@ -69,9 +68,13 @@ fail(xrpcerr.GenericError(fmt.Errorf("repo record %s has no repoDid", repoAt))) return } - repoDid := *repoRec.RepoDid + repoDid, err := syntax.ParseDID(*repoRec.RepoDid) + if err != nil { + fail(xrpcerr.GenericError(fmt.Errorf("repo record %q has invalid repoDid: %q", repoAt, *repoRec.RepoDid))) + return + } - if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, repoDid); !ok || err != nil { + if ok, err := x.Enforcer.IsRepoSecretsAllowed(actorDid, repoDid); !ok || err != nil { l.Error("insufficient permissions", "did", actorDid.String()) writeError(w, xrpcerr.AccessControlError(actorDid.String()), http.StatusUnauthorized) return diff --git a/spindle/xrpc/ci_pipeline_trigger_pipeline.go b/spindle/xrpc/ci_pipeline_trigger_pipeline.go --- a/spindle/xrpc/ci_pipeline_trigger_pipeline.go +++ b/spindle/xrpc/ci_pipeline_trigger_pipeline.go @@ -10,7 +10,6 @@ "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" - "tangled.org/core/rbac" xrpcerr "tangled.org/core/xrpc/errors" ) @@ -162,7 +161,7 @@ return "", xerr, false } - isPushAllowed, err := x.Enforcer.IsPushAllowed(actorDid.String(), rbac.ThisServer, repoDid.String()) + isPushAllowed, err := x.Enforcer.IsRepoCiTriggerAllowed(actorDid, repoDid) if err != nil || !isPushAllowed { return "", xrpcerr.AccessControlError(actorDid.String()), false } diff --git a/spindle/xrpc/list_secrets.go b/spindle/xrpc/list_secrets.go --- a/spindle/xrpc/list_secrets.go +++ b/spindle/xrpc/list_secrets.go @@ -10,7 +10,6 @@ "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" - "tangled.org/core/rbac" "tangled.org/core/spindle/secrets" xrpcerr "tangled.org/core/xrpc/errors" ) @@ -64,9 +63,13 @@ fail(xrpcerr.GenericError(fmt.Errorf("repo record %s has no repoDid", repoAt))) return } - repoDid := *repoRec.RepoDid + repoDid, err := syntax.ParseDID(*repoRec.RepoDid) + if err != nil { + fail(xrpcerr.GenericError(fmt.Errorf("repo record %q has invalid repoDid: %q", repoAt, *repoRec.RepoDid))) + return + } - if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, repoDid); !ok || err != nil { + if ok, err := x.Enforcer.IsRepoSecretsAllowed(actorDid, repoDid); !ok || err != nil { l.Error("insufficient permissions", "did", actorDid.String()) writeError(w, xrpcerr.AccessControlError(actorDid.String()), http.StatusUnauthorized) return diff --git a/spindle/xrpc/remove_secret.go b/spindle/xrpc/remove_secret.go --- a/spindle/xrpc/remove_secret.go +++ b/spindle/xrpc/remove_secret.go @@ -9,7 +9,6 @@ "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" - "tangled.org/core/rbac" "tangled.org/core/spindle/secrets" xrpcerr "tangled.org/core/xrpc/errors" ) @@ -63,9 +62,13 @@ fail(xrpcerr.GenericError(fmt.Errorf("repo record %s has no repoDid", repoAt))) return } - repoDid := *repoRec.RepoDid + repoDid, err := syntax.ParseDID(*repoRec.RepoDid) + if err != nil { + fail(xrpcerr.GenericError(fmt.Errorf("repo record %q has invalid repoDid: %q", repoAt, *repoRec.RepoDid))) + return + } - if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, repoDid); !ok || err != nil { + if ok, err := x.Enforcer.IsRepoSecretsAllowed(actorDid, repoDid); !ok || err != nil { l.Error("insufficient permissions", "did", actorDid.String()) writeError(w, xrpcerr.AccessControlError(actorDid.String()), http.StatusUnauthorized) return diff --git a/spindle/xrpc/xrpc.go b/spindle/xrpc/xrpc.go --- a/spindle/xrpc/xrpc.go +++ b/spindle/xrpc/xrpc.go @@ -14,7 +14,7 @@ "tangled.org/core/api/tangled" "tangled.org/core/idresolver" "tangled.org/core/notifier" - "tangled.org/core/rbac" + "tangled.org/core/rbac/v2" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" -- tangled.sh