package spindle import ( "context" "encoding/json" "log/slog" "path/filepath" "strings" "testing" "time" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" "tangled.org/core/idresolver" "tangled.org/core/jetstream" "tangled.org/core/notifier" "tangled.org/core/rbac" "tangled.org/core/spindle/artifactstore" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/webhook" "tangled.org/core/tapc" "tangled.org/core/workflow" ) func TestBackfillOwnersBootstrapsConfiguredOwner(t *testing.T) { d, _ := newTestSpindleDB(t) owner := syntax.DID("did:plc:spindleowner") cfg := &config.Config{} cfg.Server.Owner = owner.String() s := &Spindle{ db: d, l: slog.Default(), cfg: cfg, wh: webhook.New(d, false), } if dids, err := s.backfillOwners(); err != nil || len(dids) != 1 || dids[0] != owner { t.Fatalf("backfill owners = %v, %v, want [%s]", dids, err, owner) } } func TestBackfillOwnersFailsWithoutTheWholeList(t *testing.T) { d, _ := newTestSpindleDB(t) s := &Spindle{db: d, l: slog.Default(), cfg: &config.Config{}} d.Close() if dids, err := s.backfillOwners(); err == nil { t.Fatalf("backfill owners = %v from a closed db, want an error", dids) } } func TestProcessRepo_MembershipCheck(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) } s := &Spindle{ db: d, e: e, vault: newTestVault(t), l: slog.Default(), cfg: cfg, jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), } tap := &Tap{ spindle: s, logger: slog.Default(), } ownerDid := syntax.DID("did:plc:memberowner") nonMemberDid := syntax.DID("did:plc:nonmemberowner") repoDid := "did:plc:testrepo123" err := db.AddToAllowlist(d, ownerDid) if err != nil { t.Fatalf("AddToAllowlist: %v", err) } recNonMember := tangled.Repo{ Knot: "knot.test", RepoDid: &repoDid, Spindle: &cfg.Server.Hostname, CreatedAt: time.Now().Format(time.RFC3339), } recNonMemberJson, _ := json.Marshal(recNonMember) err = tap.processRepo(context.Background(), &tapc.RecordEventData{ Live: true, Did: nonMemberDid, Rkey: "test-repo-rkey", Collection: syntax.NSID(tangled.RepoNSID), Action: tapc.RecordCreateAction, Record: recNonMemberJson, }) if err != nil { t.Fatalf("processRepo returned error for non-member: %v", err) } _, err = d.GetRepoByOwnerRkey(nonMemberDid, "test-repo-rkey") if err == nil { t.Fatal("repo for non-member was registered in DB, expected rejection") } recMember := tangled.Repo{ Knot: "knot.test", RepoDid: &repoDid, Spindle: &cfg.Server.Hostname, CreatedAt: time.Now().Format(time.RFC3339), } recMemberJson, _ := json.Marshal(recMember) err = tap.processRepo(context.Background(), &tapc.RecordEventData{ Live: true, Did: ownerDid, Rkey: "test-repo-rkey", Collection: syntax.NSID(tangled.RepoNSID), Action: tapc.RecordCreateAction, Record: recMemberJson, }) if err != nil { t.Fatalf("processRepo returned unexpected error for member: %v", err) } _, err = d.GetRepoByOwnerRkey(ownerDid, "test-repo-rkey") if err != nil { t.Fatalf("repo for member was not registered in DB: %v", err) } } func TestProcessPull_PushAllowedCheck(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) } s := &Spindle{ db: d, e: e, l: slog.Default(), cfg: cfg, res: idresolver.DefaultResolver("https://plc.test"), jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), } repoOwnerDid := syntax.DID("did:plc:repoowner") nonPusherDid := syntax.DID("did:plc:nonpusher") pusherDid := syntax.DID("did:plc:pusher") repoDid := syntax.DID("did:plc:testrepo123") err := d.AddRepo(db.Repo{ Knot: "knot.test", Owner: repoOwnerDid, Rkey: "test-repo-rkey", RepoDid: repoDid, CreatedAt: time.Now().Format(time.RFC3339), }) if err != nil { t.Fatalf("AddRepo: %v", err) } err = e.AddRepo(repoOwnerDid.String(), rbac.ThisServer, repoDid.String()) if err != nil { t.Fatalf("AddRepo permissions: %v", err) } err = e.AddCollaborator(pusherDid.String(), rbac.ThisServer, repoDid.String()) if err != nil { t.Fatalf("AddCollaborator: %v", err) } pullRecord := tangled.RepoPull{ Target: &tangled.RepoPull_Target{ Branch: "main", Repo: repoDid.String(), }, Source: &tangled.RepoPull_Source{ Branch: "feature", Repo: nil, // branch-based PR (source repo is nil) }, } pullRecordJson, _ := json.Marshal(pullRecord) err = s.processPull(context.Background(), &tapc.RecordEventData{ Live: true, Did: nonPusherDid, Rkey: "pull-rkey-1", Collection: syntax.NSID(tangled.RepoPullNSID), Action: tapc.RecordCreateAction, Record: pullRecordJson, }) if err != nil { t.Fatalf("processPull returned error for non-pusher: %v", err) } err = s.processPull(context.Background(), &tapc.RecordEventData{ Live: true, Did: pusherDid, Rkey: "pull-rkey-2", Collection: syntax.NSID(tangled.RepoPullNSID), Action: tapc.RecordCreateAction, Record: pullRecordJson, }) if err != nil { t.Fatalf("processPull returned error for pusher: %v", err) } // legacy round-based records fall back to fetching the patch blob; the // fetch fails here (plc/pds are not real) which proves we did not skip legacyRecord := pullRecord legacyRecord.Rounds = []*tangled.RepoPull_Round{{CreatedAt: time.Now().Format(time.RFC3339)}} legacyRecordJson, _ := json.Marshal(legacyRecord) err = s.processPull(context.Background(), &tapc.RecordEventData{ Live: true, Did: pusherDid, Rkey: "pull-rkey-3", Collection: syntax.NSID(tangled.RepoPullNSID), Action: tapc.RecordCreateAction, Record: legacyRecordJson, }) if err == nil || !strings.Contains(err.Error(), "resolve PR owner") { t.Fatalf("expected legacy round-based PR to fetch the patch blob, got: %v", err) } // versions take precedence over rounds bothRecord := legacyRecord bothRecord.Versions = []*tangled.RepoPull_Version{{Head: "deadbeef", Base: "cafe", CreatedAt: time.Now().Format(time.RFC3339)}} bothRecordJson, _ := json.Marshal(bothRecord) err = s.processPull(context.Background(), &tapc.RecordEventData{ Live: true, Did: pusherDid, Rkey: "pull-rkey-4", Collection: syntax.NSID(tangled.RepoPullNSID), Action: tapc.RecordCreateAction, Record: bothRecordJson, }) if err != nil { t.Fatalf("processPull returned error for versioned PR: %v", err) } // the pipeline itself swallows downstream failures, so assert the push // access check directly allowed, err := s.isPullTriggerAuthorized(nonPusherDid.String(), nonPusherDid.String(), repoDid.String()) if err != nil { t.Fatalf("isPullTriggerAuthorized(non-pusher): %v", err) } if allowed { t.Fatal("non-pusher was authorized to trigger a pull pipeline") } allowed, err = s.isPullTriggerAuthorized(pusherDid.String(), pusherDid.String(), repoDid.String()) if err != nil { t.Fatalf("isPullTriggerAuthorized(pusher): %v", err) } if !allowed { t.Fatal("pusher was not authorized to trigger a pull pipeline") } } func TestProcessRepo_HijackRepoDidCheck(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) } s := &Spindle{ db: d, e: e, l: slog.Default(), cfg: cfg, jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), } tap := &Tap{ spindle: s, logger: slog.Default(), } aliceDid := syntax.DID("did:plc:alice") bobDid := syntax.DID("did:plc:bob") repoDid := "did:plc:sharedrepo" err := db.AddToAllowlist(d, aliceDid) if err != nil { t.Fatalf("AddToAllowlist %s: %v", aliceDid, err) } err = db.AddToAllowlist(d, bobDid) if err != nil { t.Fatalf("AddToAllowlist %s: %v", bobDid, err) } err = d.AddRepo(db.Repo{ Knot: "knot.test", Owner: aliceDid, Rkey: "alice-repo", RepoDid: syntax.DID(repoDid), CreatedAt: time.Now().Format(time.RFC3339), }) if err != nil { t.Fatalf("d.AddRepo: %v", err) } // bob tries to register alice's repo did, must reject the hijack recBob := tangled.Repo{ Knot: "knot.test", RepoDid: &repoDid, Spindle: &cfg.Server.Hostname, CreatedAt: time.Now().Format(time.RFC3339), } recBobJson, _ := json.Marshal(recBob) err = tap.processRepo(context.Background(), &tapc.RecordEventData{ Live: true, Did: bobDid, Rkey: "bob-repo", Collection: syntax.NSID(tangled.RepoNSID), Action: tapc.RecordCreateAction, Record: recBobJson, }) if err != nil { t.Fatalf("processRepo returned error on duplicate repoDid hijack attempt: %v", err) } _, err = d.GetRepoByOwnerRkey(bobDid, "bob-repo") if err == nil { t.Fatal("bob successfully hijacked alice's repoDid in DB, expected rejection") } } func TestProcessCollaborator_RBAC(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" ownerDid := syntax.DID("did:plc:repoowner") otherDid := syntax.DID("did:plc:otheractor") subjectDid := syntax.DID("did:plc:collabsubject") repoDid := syntax.DID("did:plc:testrepo123") h, err := syntax.ParseHandle("collabsubject.test") if err != nil { t.Fatalf("syntax.ParseHandle: %v", err) } mockIdent := &identity.Identity{ DID: subjectDid, Handle: h, } resolver := idresolver.NewMockResolver(idresolver.MockDirectory{Ident: mockIdent}) jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) } s := &Spindle{ db: d, e: e, l: slog.Default(), cfg: cfg, res: resolver, jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), } tap := &Tap{ spindle: s, logger: slog.Default(), } err = d.AddRepo(db.Repo{ Knot: "knot.test", Owner: ownerDid, Rkey: "test-repo-rkey", RepoDid: repoDid, CreatedAt: time.Now().Format(time.RFC3339), }) if err != nil { t.Fatalf("AddRepo: %v", err) } collabRecord := tangled.RepoCollaborator{ Subject: subjectDid.String(), Repo: repoDid.String(), } collabRecordJson, _ := json.Marshal(collabRecord) err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ Live: true, Did: otherDid, Rkey: "collab-rkey-1", Collection: syntax.NSID(tangled.RepoCollaboratorNSID), Action: tapc.RecordCreateAction, Record: collabRecordJson, }) if err != nil { t.Fatalf("processCollaborator returned error: %v", err) } _, err = d.GetRepoCollaborator(otherDid, "collab-rkey-1") if err == nil { t.Fatal("collaborator from non-owner was registered in DB") } err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ Live: true, Did: ownerDid, Rkey: "collab-rkey-2", Collection: syntax.NSID(tangled.RepoCollaboratorNSID), Action: tapc.RecordCreateAction, Record: collabRecordJson, }) if err != nil { t.Fatalf("processCollaborator returned error: %v", err) } _, err = d.GetRepoCollaborator(ownerDid, "collab-rkey-2") if err == nil { t.Fatal("collaborator registered despite missing Casbin invite permission") } err = e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()) if err != nil { t.Fatalf("AddRepo permissions: %v", err) } err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ Live: true, Did: ownerDid, Rkey: "collab-rkey-3", Collection: syntax.NSID(tangled.RepoCollaboratorNSID), Action: tapc.RecordCreateAction, Record: collabRecordJson, }) if err != nil { t.Fatalf("processCollaborator failed for authorized owner: %v", err) } c, err := d.GetRepoCollaborator(ownerDid, "collab-rkey-3") if err != nil { t.Fatalf("GetRepoCollaborator error: %v", err) } if c.Subject != subjectDid || c.RepoDid != repoDid { t.Fatalf("unexpected collaborator: %+v", c) } ok, err := e.IsRepoCollaborator(subjectDid.String(), rbac.ThisServer, repoDid.String()) if err != nil || !ok { t.Fatalf("Casbin policy for collaborator missing or err: %v", err) } err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ Live: true, Did: ownerDid, Rkey: "collab-rkey-3", Collection: syntax.NSID(tangled.RepoCollaboratorNSID), Action: tapc.RecordDeleteAction, }) if err != nil { t.Fatalf("delete collaborator process returned error: %v", err) } _, err = d.GetRepoCollaborator(ownerDid, "collab-rkey-3") if err == nil { t.Fatal("collaborator DB row remained after deletion") } ok, err = e.IsRepoCollaborator(subjectDid.String(), rbac.ThisServer, repoDid.String()) if err != nil || ok { t.Fatal("Casbin policy for collaborator remained after deletion") } } func TestTeardownRepo_RBAC(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" cfg.Server.RepoDir = t.TempDir() cfg.Server.LogDir = t.TempDir() jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) } vault, verr := secrets.NewSQLiteManager(filepath.Join(t.TempDir(), "secrets.db")) if verr != nil { t.Fatalf("NewSQLiteManager: %v", verr) } stores, serr := artifactstore.NewStores(config.ArtifactStores{ Disk: config.ArtifactStoreDisk{Dir: t.TempDir()}, }, "", "") if serr != nil { t.Fatalf("NewStores: %v", serr) } s := &Spindle{ db: d, e: e, l: slog.Default(), cfg: cfg, jc: jc, vault: vault, stores: stores, n: ptr(notifier.New()), engs: map[string]models.Engine{}, rootCtx: context.Background(), wh: webhook.New(d, false), } tap := &Tap{ spindle: s, logger: slog.Default(), } ownerDid := syntax.DID("did:plc:repoowner") repoDid := syntax.DID("did:plc:testrepo123") collabDid := syntax.DID("did:plc:collab") err := d.AddRepo(db.Repo{ Knot: "knot.test", Owner: ownerDid, Rkey: "test-repo-rkey", RepoDid: repoDid, CreatedAt: time.Now().Format(time.RFC3339), }) if err != nil { t.Fatalf("AddRepo DB: %v", err) } err = e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()) if err != nil { t.Fatalf("AddRepo policy: %v", err) } err = d.AddRepoCollaborator(db.RepoCollaborator{ OwnerDid: ownerDid, Rkey: "collab-rkey", Subject: collabDid, RepoDid: repoDid, }) if err != nil { t.Fatalf("AddCollaborator DB: %v", err) } err = e.AddCollaborator(collabDid.String(), rbac.ThisServer, repoDid.String()) if err != nil { t.Fatalf("AddCollaborator policy: %v", err) } // entities the delete path only reaches through DeleteAllData if err := vault.AddSecret(context.Background(), secrets.UnlockedSecret{Key: "api_key", Value: "v", Repo: secrets.RepoIdentifier(repoDid.String()), CreatedBy: ownerDid}); err != nil { t.Fatalf("AddSecret: %v", err) } if err := d.EnqueueJob(context.Background(), repoDid.String(), models.PipelineId("p1"), nil, tangled.Pipeline{}, "", ""); err != nil { t.Fatalf("EnqueueJob: %v", err) } err = tap.processRepo(context.Background(), &tapc.RecordEventData{ Live: true, Did: ownerDid, Rkey: "test-repo-rkey", Collection: syntax.NSID(tangled.RepoNSID), Action: tapc.RecordDeleteAction, }) if err != nil { t.Fatalf("processRepo delete returned error: %v", err) } _, err = d.GetRepoByOwnerRkey(ownerDid, "test-repo-rkey") if err == nil { t.Fatal("repo remained in DB after delete") } collabs, err := d.ListCollaboratorsByRepoDid(repoDid) if err != nil { t.Fatalf("ListCollaboratorsByRepoDid: %v", err) } if len(collabs) > 0 { t.Fatal("collaborators remained in DB after delete") } ok, err := e.IsRepoOwner(ownerDid.String(), rbac.ThisServer, repoDid.String()) if err != nil || ok { t.Fatal("repo owner policy remained in Casbin after delete") } ok, err = e.IsRepoCollaborator(collabDid.String(), rbac.ThisServer, repoDid.String()) if err != nil || ok { t.Fatal("collaborator policy remained in Casbin after delete") } if got, kerr := vault.GetSecretsUnlocked(context.Background(), secrets.RepoIdentifier(repoDid.String())); kerr != nil || len(got) > 0 { t.Fatalf("secrets remained in vault after delete: %v %v", got, kerr) } var njobs int if err := d.QueryRow(`select count(*) from jobs where repo_did = ?`, repoDid.String()).Scan(&njobs); err != nil || njobs != 0 { t.Fatalf("jobs remained in DB after delete: %d %v", njobs, err) } } func TestProcessRepo_ForgeDeleteRejection(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) } s := &Spindle{ db: d, e: e, l: slog.Default(), cfg: cfg, jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), } tap := &Tap{ spindle: s, logger: slog.Default(), } aliceDid := syntax.DID("did:plc:alice") bobDid := syntax.DID("did:plc:bob") repoDid := syntax.DID("did:plc:sharedrepo") err := d.AddRepo(db.Repo{ Knot: "knot.test", Owner: aliceDid, Rkey: "test-repo-rkey", RepoDid: repoDid, CreatedAt: time.Now().Format(time.RFC3339), }) if err != nil { t.Fatalf("AddRepo DB: %v", err) } err = e.AddRepo(aliceDid.String(), rbac.ThisServer, repoDid.String()) if err != nil { t.Fatalf("AddRepo policy: %v", err) } // bob tries to delete alice's repo, must reject forged delete err = tap.processRepo(context.Background(), &tapc.RecordEventData{ Live: true, Did: bobDid, Rkey: "test-repo-rkey", Collection: syntax.NSID(tangled.RepoNSID), Action: tapc.RecordDeleteAction, }) if err != nil { t.Fatalf("processRepo returned error on delete: %v", err) } _, err = d.GetRepoByOwnerRkey(aliceDid, "test-repo-rkey") if err != nil { t.Fatalf("Alice's repo was deleted or error: %v", err) } ok, err := e.IsRepoOwner(aliceDid.String(), rbac.ThisServer, repoDid.String()) if err != nil || !ok { t.Fatal("Alice's owner policy was removed from Casbin by forged delete") } } func TestProcessCollaborator_ForgeDeleteRejection(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) } s := &Spindle{ db: d, e: e, l: slog.Default(), cfg: cfg, res: idresolver.DefaultResolver("https://plc.test"), jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), } tap := &Tap{ spindle: s, logger: slog.Default(), } ownerDid := syntax.DID("did:plc:repoowner") bobDid := syntax.DID("did:plc:bob") collabDid := syntax.DID("did:plc:collab") repoDid := syntax.DID("did:plc:testrepo123") err := d.AddRepo(db.Repo{ Knot: "knot.test", Owner: ownerDid, Rkey: "test-repo-rkey", RepoDid: repoDid, CreatedAt: time.Now().Format(time.RFC3339), }) if err != nil { t.Fatalf("AddRepo: %v", err) } err = e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()) if err != nil { t.Fatalf("AddRepo permissions: %v", err) } err = d.AddRepoCollaborator(db.RepoCollaborator{ OwnerDid: ownerDid, Rkey: "collab-rkey", Subject: collabDid, RepoDid: repoDid, }) if err != nil { t.Fatalf("AddRepoCollaborator: %v", err) } err = e.AddCollaborator(collabDid.String(), rbac.ThisServer, repoDid.String()) if err != nil { t.Fatalf("AddCollaborator policy: %v", err) } // bob tries to delete alice's collaborator, must reject forged delete err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ Live: true, Did: bobDid, Rkey: "collab-rkey", Collection: syntax.NSID(tangled.RepoCollaboratorNSID), Action: tapc.RecordDeleteAction, }) if err != nil { t.Fatalf("processCollaborator delete returned error: %v", err) } _, err = d.GetRepoCollaborator(ownerDid, "collab-rkey") if err != nil { t.Fatalf("collaborator was deleted from DB: %v", err) } ok, err := e.IsRepoCollaborator(collabDid.String(), rbac.ThisServer, repoDid.String()) if err != nil || !ok { t.Fatal("collaborator policy was removed from Casbin by forged delete") } } func TestPullStatusAction(t *testing.T) { tests := []struct { name string status string wantAction string wantOK bool }{ {"open maps to reopened", tangled.RepoPullStatusOpen, workflow.PullRequestActionReopened, true}, {"closed maps to closed", tangled.RepoPullStatusClosed, workflow.PullRequestActionClosed, true}, {"merged maps to merged", tangled.RepoPullStatusMerged, workflow.PullRequestActionMerged, true}, {"unknown status is rejected", "sh.tangled.repo.pull.status.bogus", "", false}, {"empty status is rejected", "", "", false}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { action, ok := pullStatusAction(tt.status) if ok != tt.wantOK { t.Fatalf("pullStatusAction(%q) ok = %v, want %v", tt.status, ok, tt.wantOK) } if action != tt.wantAction { t.Fatalf("pullStatusAction(%q) action = %q, want %q", tt.status, action, tt.wantAction) } }) } } func TestProcessPullStatus(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if jcerr != nil { t.Fatalf("NewJetstreamClient: %v", jcerr) } s := &Spindle{ db: d, e: e, l: slog.Default(), cfg: cfg, res: idresolver.DefaultResolver("https://plc.test"), jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), } statusRecord := func(status, pullUri string) []byte { rec := tangled.RepoPullStatus{ Status: status, Pull: pullUri, CreatedAt: time.Now().Format(time.RFC3339), } b, _ := json.Marshal(rec) return b } validPullUri := "at://did:plc:pullowner/sh.tangled.repo.pull/pull-rkey-1" // non-create actions are ignored if err := s.processPullStatus(context.Background(), &tapc.RecordEventData{ Live: true, Did: syntax.DID("did:plc:actor"), Rkey: "status-rkey", Collection: syntax.NSID(tangled.RepoPullStatusNSID), Action: tapc.RecordUpdateAction, Record: statusRecord(tangled.RepoPullStatusClosed, validPullUri), }); err != nil { t.Fatalf("update action should be a no-op, got: %v", err) } // unknown status variant is skipped without error if err := s.processPullStatus(context.Background(), &tapc.RecordEventData{ Live: true, Did: syntax.DID("did:plc:actor"), Rkey: "status-rkey", Collection: syntax.NSID(tangled.RepoPullStatusNSID), Action: tapc.RecordCreateAction, Record: statusRecord("sh.tangled.repo.pull.status.bogus", validPullUri), }); err != nil { t.Fatalf("unknown status should be skipped, got: %v", err) } // a malformed pull at-uri is skipped without error if err := s.processPullStatus(context.Background(), &tapc.RecordEventData{ Live: true, Did: syntax.DID("did:plc:actor"), Rkey: "status-rkey", Collection: syntax.NSID(tangled.RepoPullStatusNSID), Action: tapc.RecordCreateAction, Record: statusRecord(tangled.RepoPullStatusClosed, "not-an-at-uri"), }); err != nil { t.Fatalf("invalid pull at-uri should be skipped, got: %v", err) } // a status pointing at a non-pull subject is skipped without error if err := s.processPullStatus(context.Background(), &tapc.RecordEventData{ Live: true, Did: syntax.DID("did:plc:actor"), Rkey: "status-rkey", Collection: syntax.NSID(tangled.RepoPullStatusNSID), Action: tapc.RecordCreateAction, Record: statusRecord(tangled.RepoPullStatusClosed, "at://did:plc:x/sh.tangled.repo.issue/y"), }); err != nil { t.Fatalf("non-pull subject should be skipped, got: %v", err) } // a valid close event resolves the pull record; fetch fails because plc/pds // are not real, confirming we reached the fetch stage with a mapped action. err := s.processPullStatus(context.Background(), &tapc.RecordEventData{ Live: true, Did: syntax.DID("did:plc:actor"), Rkey: "status-rkey", Collection: syntax.NSID(tangled.RepoPullStatusNSID), Action: tapc.RecordCreateAction, Record: statusRecord(tangled.RepoPullStatusClosed, validPullUri), }) if err == nil { t.Fatal("expected error fetching pull record against fake pds, got nil") } if !strings.Contains(err.Error(), "fetch pull record") { t.Fatalf("expected fetch pull record error, got: %v", err) } } func TestIsPullTriggerAuthorized(t *testing.T) { d, e := newTestSpindleDB(t) s := &Spindle{ db: d, e: e, l: slog.Default(), wh: webhook.New(d, false), } repoDid := syntax.DID("did:plc:targetrepo") ownerDid := syntax.DID("did:plc:owner") // has push (repo owner) collaboratorDid := "did:plc:collaborator" // has push noPushAuthorDid := "did:plc:nopushauthor" // pull author without push strangerDid := "did:plc:stranger" // no push, not the author if err := e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()); err != nil { t.Fatalf("AddRepo permissions: %v", err) } if err := e.AddCollaborator(collaboratorDid, rbac.ThisServer, repoDid.String()); err != nil { t.Fatalf("AddCollaborator: %v", err) } cases := []struct { name string eventDid string pullDid string allowed bool }{ // direct pr path: event author == pull author {"direct pr by pushing author", ownerDid.String(), ownerDid.String(), true}, {"direct pr by non-pushing author", noPushAuthorDid, noPushAuthorDid, false}, // status path: actor differs from pull author, pull author has push {"actor with push acts on authorized pr", collaboratorDid, ownerDid.String(), true}, {"pull author acts on own authorized pr", ownerDid.String(), ownerDid.String(), true}, {"stranger without push acts on authorized pr", strangerDid, ownerDid.String(), false}, // pull author must always have push, even if the actor does {"pushing actor on unauthorized pull author", ownerDid.String(), noPushAuthorDid, false}, {"non-pushing actor on unauthorized pull author", strangerDid, noPushAuthorDid, false}, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { allowed, err := s.isPullTriggerAuthorized(tc.eventDid, tc.pullDid, repoDid.String()) if err != nil { t.Fatalf("isPullTriggerAuthorized: %v", err) } if allowed != tc.allowed { t.Fatalf("event=%q pull=%q allowed=%v, want %v", tc.eventDid, tc.pullDid, allowed, tc.allowed) } }) } } func TestProcessRepo_SeveralRecordsForOneRepoDid(t *testing.T) { type rec struct { rkey string name string createdAt string live bool } cases := []struct { name string registered []rec events []rec want string }{ { name: "live rename", registered: []rec{{rkey: "old", createdAt: "2026-07-01T10:00:00Z"}}, events: []rec{{rkey: "new", createdAt: "2026-07-02T10:00:00+03:00", live: true}}, want: "new", }, { name: "live rename onto a record whose createdAt an edit reset", registered: []rec{{rkey: "frytg-web", createdAt: "2026-08-21T18:09:19+03:00"}}, events: []rec{{rkey: "frytg-digital", createdAt: "2026-07-28T15:29:15Z", live: true}}, want: "frytg-digital", }, { name: "live rename keeping createdAt", registered: []rec{{rkey: "old", createdAt: "2026-07-01T10:00:00Z"}}, events: []rec{{rkey: "new", createdAt: "2026-07-01T10:00:00Z", live: true}}, want: "new", }, { name: "backfilled alias with the same createdAt", registered: []rec{{rkey: "atproto-deeplink", createdAt: "2026-05-30T02:22:08Z"}}, events: []rec{{rkey: "deeplink", createdAt: "2026-05-30T02:22:08Z"}}, want: "atproto-deeplink", }, { name: "backfilled alias older than the registration", registered: []rec{{rkey: "hefter", createdAt: "2026-08-13T23:55:57+03:00"}}, events: []rec{{rkey: "ledger-bridge", createdAt: "2026-07-30T16:42:22Z"}}, want: "hefter", }, { name: "rename missed while down, found by backfill", registered: []rec{{rkey: "old", createdAt: "2026-09-02T01:00:00+03:00"}}, events: []rec{{rkey: "new", createdAt: "2026-09-01T23:00:00Z"}}, want: "new", }, { name: "backfill drops a legacy row without createdAt", registered: []rec{{rkey: "stale-bogus-rkey"}}, events: []rec{{rkey: "fresh-pds-rkey", createdAt: "2024-06-01T00:00:00Z"}}, want: "fresh-pds-rkey", }, { name: "fresh backfill keeps the newest record", events: []rec{ {rkey: "3lun4rlswkc22", name: "rockbox-zig", createdAt: "2025-07-23T13:28:28Z"}, {rkey: "rockboxd", createdAt: "2026-07-31T09:23:46+03:00"}, {rkey: "rockboxdd", createdAt: "2026-07-31T09:23:30+03:00"}, }, want: "rockboxd", }, { name: "fresh backfill prefers a name rkey over a tid rkey", events: []rec{ {rkey: "3mcnqv7zdo222", name: "tranquil-pds", createdAt: "2026-01-17T23:21:07Z"}, {rkey: "tranquil-pds", createdAt: "2026-01-17T23:21:07Z"}, }, want: "tranquil-pds", }, { name: "fresh backfill keeps the first of a tie", events: []rec{ {rkey: "atproto-deeplink", createdAt: "2026-05-30T02:22:08Z"}, {rkey: "deeplink", createdAt: "2026-05-30T02:22:08Z"}, }, want: "atproto-deeplink", }, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { d, e := newTestSpindleDB(t) cfg := &config.Config{} cfg.Server.Hostname = "spindle.test" cfg.Server.RepoDir = t.TempDir() jc, err := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) if err != nil { t.Fatalf("NewJetstreamClient: %v", err) } tap := &Tap{ spindle: &Spindle{ db: d, e: e, vault: newTestVault(t), l: slog.Default(), cfg: cfg, jc: jc, rootCtx: context.Background(), wh: webhook.New(d, false), }, logger: slog.Default(), } owner := syntax.DID("did:plc:owner") repoDid := "did:plc:samerepo" if err := e.AddSpindle(rbac.ThisServer); err != nil { t.Fatalf("AddSpindle: %v", err) } if err := e.AddSpindleMember(rbac.ThisServer, owner.String()); err != nil { t.Fatalf("AddSpindleMember: %v", err) } for _, r := range tc.registered { if err := d.AddRepo(db.Repo{ Knot: "knot.test", Owner: owner, Rkey: syntax.RecordKey(r.rkey), RepoDid: syntax.DID(repoDid), Name: r.rkey, CreatedAt: r.createdAt, }); err != nil { t.Fatalf("AddRepo: %v", err) } } for _, r := range tc.events { record := tangled.Repo{ Knot: "knot.test", RepoDid: &repoDid, Spindle: &cfg.Server.Hostname, CreatedAt: r.createdAt, } if r.name != "" { record.Name = &r.name } raw, _ := json.Marshal(record) if err := tap.processRepo(context.Background(), &tapc.RecordEventData{ Live: r.live, Did: owner, Rkey: syntax.RecordKey(r.rkey), Collection: syntax.NSID(tangled.RepoNSID), Action: tapc.RecordCreateAction, Record: raw, }); err != nil { t.Fatalf("processRepo %s: %v", r.rkey, err) } } rows, err := d.ReposByDid(syntax.DID(repoDid)) if err != nil { t.Fatalf("ReposByDid: %v", err) } var got []string for _, r := range rows { got = append(got, r.Rkey.String()) } if len(got) != 1 || got[0] != tc.want { t.Fatalf("rows for repo did = %v, want [%s]", got, tc.want) } }) } } func TestOlderRecordVersionsDontUndoNewerOnes(t *testing.T) { d, e := newTestSpindleDB(t) owner := syntax.DID("did:plc:owner") subject := syntax.DID("did:plc:subject") repo := syntax.DID("did:plc:repo") if err := d.AddRepo(db.Repo{Knot: "knot.test", Owner: owner, Rkey: "repo", RepoDid: repo}); err != nil { t.Fatal(err) } if err := e.AddRepo(owner.String(), rbac.ThisServer, repo.String()); err != nil { t.Fatal(err) } resolver := idresolver.NewMockResolver(idresolver.MockDirectory{Ident: &identity.Identity{DID: subject}}) tap := &Tap{ spindle: &Spindle{db: d, e: e, l: slog.Default(), cfg: &config.Config{}, res: resolver}, logger: slog.Default(), } raw, _ := json.Marshal(tangled.RepoCollaborator{Subject: subject.String(), Repo: repo.String()}) send := func(live bool, action tapc.RecordAction, rev string) { t.Helper() evt := tapc.Event{Type: tapc.EvtRecord, Record: &tapc.RecordEventData{ Live: live, Did: owner, Rev: rev, Collection: syntax.NSID(tangled.RepoCollaboratorNSID), Rkey: "collab", Action: action, }} if action != tapc.RecordDeleteAction { evt.Record.Record = raw } if err := tap.processEvent(context.Background(), evt); err != nil { t.Fatal(err) } } isCollaborator := func() bool { t.Helper() _, err := d.GetRepoCollaborator(owner, "collab") return err == nil } send(true, tapc.RecordCreateAction, "3jzfcijpj2z2a") send(true, tapc.RecordDeleteAction, "3jzfcijpj2z2c") send(false, tapc.RecordCreateAction, "3jzfcijpj2z2b") if isCollaborator() { t.Fatal("an older backfilled create brought back a collaborator deleted live") } send(false, tapc.RecordCreateAction, "3jzfcijpj2z2e") send(true, tapc.RecordDeleteAction, "3jzfcijpj2z2d") if !isCollaborator() { t.Fatal("a live delete older than the backfilled record removed it") } send(false, tapc.RecordCreateAction, "3jzfcijpj2z2e") if !isCollaborator() { t.Fatal("replaying the version spindle has dropped it") } if rev, err := d.RecordRev("at://did:plc:owner/sh.tangled.repo.collaborator/collab"); err != nil || rev != "3jzfcijpj2z2e" { t.Fatalf("record rev = %q, %v, want the newest one handled", rev, err) } }