diff --git a/idresolver/mock.go b/idresolver/mock.go new file mode 100644 index 00000000..3b6a8c3c --- /dev/null +++ b/idresolver/mock.go @@ -0,0 +1,10 @@ +package idresolver + +import "github.com/bluesky-social/indigo/atproto/identity" + +func NewMockResolver(dir identity.Directory) *Resolver { + return &Resolver{ + directory: dir, + base: &identity.BaseDirectory{}, + } +} diff --git a/spindle/ingester_test.go b/spindle/ingester_test.go index 2a3890bd..99a309ec 100644 --- a/spindle/ingester_test.go +++ b/spindle/ingester_test.go @@ -3,12 +3,16 @@ package spindle import ( "context" "encoding/json" + "log/slog" + "strings" + "tangled.org/core/jetstream" "testing" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/jetstream/pkg/models" "tangled.org/core/api/tangled" + "tangled.org/core/rbac" "tangled.org/core/spindle/config" "tangled.org/core/tapc" ) @@ -85,3 +89,205 @@ func TestEmbeddedTapDoesNotSubscribeToPullRecords(t *testing.T) { } } } + +func TestIngestMember_RBAC(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(), + } + + actorDid := "did:plc:adminactor" + subjectDid := "did:plc:newmember" + rbacDomain := rbac.ThisServer + + memberRecord := tangled.SpindleMember{ + Instance: "spindle.test", + Subject: subjectDid, + } + memberRecordJson, _ := json.Marshal(memberRecord) + + evt := &models.Event{ + Did: actorDid, + Kind: models.EventKindCommit, + Commit: &models.Commit{ + Operation: models.CommitOperationCreate, + Collection: tangled.SpindleMemberNSID, + RKey: "member-rkey-1", + Record: memberRecordJson, + }, + } + + err := s.ingestMember(context.Background(), evt) + if err == nil { + t.Fatal("expected permission denied error, got nil") + } + if !strings.Contains(err.Error(), "permission denied") { + t.Fatalf("expected permission denied, got error: %v", err) + } + + var dbCount int + err = d.QueryRow(`select count(*) from spindle_members where subject = ?`, subjectDid).Scan(&dbCount) + if err != nil { + t.Fatalf("DB query error: %v", err) + } + if dbCount > 0 { + t.Fatal("spindle member was registered in DB on failed auth") + } + + err = e.AddSpindle(rbacDomain) + if err != nil { + t.Fatalf("AddSpindle: %v", err) + } + err = e.AddSpindleOwner(rbacDomain, actorDid) + if err != nil { + t.Fatalf("AddSpindleOwner: %v", err) + } + + err = s.ingestMember(context.Background(), evt) + if err != nil { + t.Fatalf("ingestMember failed for authorized actor: %v", err) + } + + err = d.QueryRow(`select count(*) from spindle_members where subject = ?`, subjectDid).Scan(&dbCount) + if err != nil || dbCount != 1 { + t.Fatalf("expected exactly 1 member in DB, got: %d (err: %v)", dbCount, err) + } + + isMember, err := e.IsSpindleMember(subjectDid, rbacDomain) + if err != nil || !isMember { + t.Fatalf("expected subject to be spindle member in Casbin, got: %t (err: %v)", isMember, err) + } + + deleteEvt := &models.Event{ + Did: actorDid, + Kind: models.EventKindCommit, + Commit: &models.Commit{ + Operation: models.CommitOperationDelete, + Collection: tangled.SpindleMemberNSID, + RKey: "member-rkey-1", + }, + } + + err = s.ingestMember(context.Background(), deleteEvt) + if err != nil { + t.Fatalf("ingestMember delete failed: %v", err) + } + + err = d.QueryRow(`select count(*) from spindle_members where subject = ?`, subjectDid).Scan(&dbCount) + if err != nil || dbCount != 0 { + t.Fatalf("expected 0 members in DB after delete, got: %d (err: %v)", dbCount, err) + } + + isMember, err = e.IsSpindleMember(subjectDid, rbacDomain) + if err != nil || isMember { + t.Fatalf("expected subject to NOT be spindle member in Casbin, got: %t (err: %v)", isMember, err) + } +} + +func TestIngestMember_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(), + } + + adminDid := "did:plc:adminactor" + bobDid := "did:plc:bobactor" + subjectDid := "did:plc:newmember" + rbacDomain := rbac.ThisServer + + err := e.AddSpindle(rbacDomain) + if err != nil { + t.Fatalf("AddSpindle: %v", err) + } + err = e.AddSpindleOwner(rbacDomain, adminDid) + if err != nil { + t.Fatalf("AddSpindleOwner: %v", err) + } + + memberRecord := tangled.SpindleMember{ + Instance: "spindle.test", + Subject: subjectDid, + } + memberRecordJson, _ := json.Marshal(memberRecord) + + evt := &models.Event{ + Did: adminDid, + Kind: models.EventKindCommit, + Commit: &models.Commit{ + Operation: models.CommitOperationCreate, + Collection: tangled.SpindleMemberNSID, + RKey: "member-rkey-1", + Record: memberRecordJson, + }, + } + + err = s.ingestMember(context.Background(), evt) + if err != nil { + t.Fatalf("ingestMember failed for admin: %v", err) + } + + var dbCount int + err = d.QueryRow(`select count(*) from spindle_members where subject = ?`, subjectDid).Scan(&dbCount) + if err != nil || dbCount != 1 { + t.Fatalf("expected member in DB, got: %d (err: %v)", dbCount, err) + } + + isMember, err := e.IsSpindleMember(subjectDid, rbacDomain) + if err != nil || !isMember { + t.Fatalf("expected subject to be spindle member, got %t (err: %v)", isMember, err) + } + + // bob tries to delete alice's spindle member record, must reject forged delete + deleteEvt := &models.Event{ + Did: bobDid, // Bob is the actor + Kind: models.EventKindCommit, + Commit: &models.Commit{ + Operation: models.CommitOperationDelete, + Collection: tangled.SpindleMemberNSID, + RKey: "member-rkey-1", + }, + } + + err = s.ingestMember(context.Background(), deleteEvt) + if err != nil { + t.Fatalf("ingestMember delete returned error: %v", err) + } + + err = d.QueryRow(`select count(*) from spindle_members where subject = ?`, subjectDid).Scan(&dbCount) + if err != nil || dbCount != 1 { + t.Fatalf("member was deleted from DB, expected remaining, count: %d (err: %v)", dbCount, err) + } + + isMember, err = e.IsSpindleMember(subjectDid, rbacDomain) + if err != nil || !isMember { + t.Fatal("member policy was removed from Casbin by forged delete") + } +} diff --git a/spindle/tapclient_test.go b/spindle/tapclient_test.go new file mode 100644 index 00000000..a46a9d3c --- /dev/null +++ b/spindle/tapclient_test.go @@ -0,0 +1,694 @@ +package spindle + +import ( + "context" + "encoding/json" + "log/slog" + "strings" + "tangled.org/core/jetstream" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/eventconsumer" + "tangled.org/core/idresolver" + "tangled.org/core/rbac" + "tangled.org/core/spindle/config" + "tangled.org/core/spindle/db" + + "tangled.org/core/tapc" +) + +type mockDirectory struct { + ident *identity.Identity +} + +func (m *mockDirectory) LookupDID(ctx context.Context, did syntax.DID) (*identity.Identity, error) { + return m.ident, nil +} + +func (m *mockDirectory) LookupHandle(ctx context.Context, handle syntax.Handle) (*identity.Identity, error) { + return m.ident, nil +} + +func (m *mockDirectory) Lookup(ctx context.Context, id syntax.AtIdentifier) (*identity.Identity, error) { + return m.ident, nil +} + +func (m *mockDirectory) Purge(ctx context.Context, id syntax.AtIdentifier) error { + return nil +} + +func TestProcessRepo_MembershipCheck(t *testing.T) { + d, e := newTestSpindleDB(t) + + cfg := &config.Config{} + cfg.Server.Hostname = "spindle.test" + + ccfg := eventconsumer.NewConsumerConfig() + ccfg.Logger = slog.Default() + ks := eventconsumer.NewConsumer(*ccfg) + + jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) + if jcerr != nil { + t.Fatalf("NewJetstreamClient: %v", jcerr) + } + s := &Spindle{ + db: d, + e: e, + l: slog.Default(), + cfg: cfg, + ks: ks, + jc: jc, + rootCtx: context.Background(), + } + + tap := &Tap{ + spindle: s, + logger: slog.Default(), + } + + ownerDid := syntax.DID("did:plc:memberowner") + nonMemberDid := syntax.DID("did:plc:nonmemberowner") + repoDid := "did:plc:testrepo123" + + err := e.AddSpindle(rbac.ThisServer) + if err != nil { + t.Fatalf("AddSpindle: %v", err) + } + err = e.AddSpindleMember(rbac.ThisServer, ownerDid.String()) + if err != nil { + t.Fatalf("AddSpindleMember: %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.Fatal("expected git clone error for valid member, but got nil") + } + + if !strings.Contains(err.Error(), "setting up sparse-clone git repo") { + t.Fatalf("expected sparse-clone error, got: %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(), + } + + 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) + } + + // fetch fails because plc/pds are not real + 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.Fatal("expected error from fetchLatestSubmission for valid pusher, but got nil") + } + + if !strings.Contains(err.Error(), "checking push access") && !strings.Contains(err.Error(), "resolve PR owner") && !strings.Contains(err.Error(), "invalid memory address") { + t.Fatalf("expected failed identity resolution or connection error, got: %v", err) + } +} + +func TestProcessRepo_HijackRepoDidCheck(t *testing.T) { + d, e := newTestSpindleDB(t) + + cfg := &config.Config{} + cfg.Server.Hostname = "spindle.test" + + ccfg := eventconsumer.NewConsumerConfig() + ccfg.Logger = slog.Default() + ks := eventconsumer.NewConsumer(*ccfg) + + jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) + if jcerr != nil { + t.Fatalf("NewJetstreamClient: %v", jcerr) + } + s := &Spindle{ + db: d, + e: e, + l: slog.Default(), + cfg: cfg, + ks: ks, + jc: jc, + rootCtx: context.Background(), + } + + tap := &Tap{ + spindle: s, + logger: slog.Default(), + } + + aliceDid := syntax.DID("did:plc:alice") + bobDid := syntax.DID("did:plc:bob") + repoDid := "did:plc:sharedrepo" + + err := e.AddSpindle(rbac.ThisServer) + if err != nil { + t.Fatalf("AddSpindle: %v", err) + } + err = e.AddSpindleMember(rbac.ThisServer, aliceDid.String()) + if err != nil { + t.Fatalf("AddSpindleMember alice: %v", err) + } + err = e.AddSpindleMember(rbac.ThisServer, bobDid.String()) + if err != nil { + t.Fatalf("AddSpindleMember bob: %v", 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(&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(), + } + + 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" + + 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(), + } + + 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) + } + + 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") + } +} + +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(), + } + + 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(), + } + + 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") + } +} diff --git a/spindle/xrpc/xrpc_test.go b/spindle/xrpc/xrpc_test.go new file mode 100644 index 00000000..9aec07e5 --- /dev/null +++ b/spindle/xrpc/xrpc_test.go @@ -0,0 +1,419 @@ +package xrpc + +import ( + "bytes" + "context" + "encoding/json" + "github.com/bluesky-social/indigo/atproto/identity" + "log/slog" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/idresolver" + "tangled.org/core/rbac" + "tangled.org/core/spindle/config" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" + "tangled.org/core/spindle/secrets" +) + +type mockTrigger struct { + triggered bool +} + +func (m *mockTrigger) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull PullContext, inputs []*tangled.Pipeline_Pair) (syntax.ATURI, error) { + m.triggered = true + return syntax.ParseATURI("at://did:plc:repoowner/sh.tangled.ci.pipeline/testrkey") +} + +func newTestXrpcDB(t *testing.T) (*db.DB, *rbac.Enforcer) { + t.Helper() + p := filepath.Join(t.TempDir(), "spindle_xrpc.db") + d, err := db.Make(context.Background(), p) + if err != nil { + t.Fatalf("db.Make: %v", err) + } + t.Cleanup(func() { d.Close() }) + e, err := rbac.NewEnforcer(p) + if err != nil { + t.Fatalf("rbac.NewEnforcer: %v", err) + } + e.E.EnableAutoSave(true) + return d, e +} + +func TestTriggerPipeline_RBAC(t *testing.T) { + d, e := newTestXrpcDB(t) + + 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) + } + + trigger := &mockTrigger{} + x := &Xrpc{ + Logger: slog.Default(), + Db: d, + Enforcer: e, + Config: &config.Config{}, + Trigger: trigger, + } + + sendReq := func(actor syntax.DID, input tangled.CiTriggerPipeline_Input) (*httptest.ResponseRecorder, int) { + body, _ := json.Marshal(input) + req := httptest.NewRequest(http.MethodPost, "/com.atproto.repo.createRecord", bytes.NewReader(body)) + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + + w := httptest.NewRecorder() + x.TriggerPipeline(w, req) + return w, w.Code + } + + sha := "0123456789abcdef0123456789abcdef01234567" + ref := "refs/heads/main" + + input := tangled.CiTriggerPipeline_Input{ + Repo: repoDid.String(), + Trigger: &tangled.CiTriggerPipeline_Input_Trigger{ + CiTrigger_Manual: &tangled.CiTrigger_Manual{ + Sha: sha, + Ref: &ref, + }, + }, + } + + w, code := sendReq(pusherDid, input) + if code != http.StatusOK { + t.Fatalf("expected 200 for pusher, got %d (body: %s)", code, w.Body.String()) + } + if !trigger.triggered { + t.Fatal("expected pipeline trigger to be called") + } + + trigger.triggered = false + + w, code = sendReq(nonPusherDid, input) + if code != http.StatusBadRequest { + t.Fatalf("expected 400 for non-pusher, got %d", code) + } + if !strings.Contains(w.Body.String(), "AccessControl") { + t.Fatalf("expected AccessControl, got: %s", w.Body.String()) + } + if trigger.triggered { + t.Fatal("expected pipeline trigger not to be called for non-pusher") + } + + badInput := input + badInput.Repo = "did:plc:unknownrepo" + w, code = sendReq(pusherDid, badInput) + if code != http.StatusBadRequest { + t.Fatalf("expected 400 for unknown repo, got %d", code) + } + if !strings.Contains(w.Body.String(), "RepoNotFound") { + t.Fatalf("expected RepoNotFound, got: %s", w.Body.String()) + } +} + +func TestCancelPipeline_RBAC(t *testing.T) { + d, e := newTestXrpcDB(t) + + 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) + } + + pipelineTid := "3mrkp6iz6os2o" + repoDidStr := repoDid.String() + tpl := tangled.Pipeline{ + TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ + Kind: "manual", + Repo: &tangled.Pipeline_TriggerRepo{ + RepoDid: &repoDidStr, + Knot: "knot.test", + Did: repoOwnerDid.String(), + }, + }, + Workflows: []*tangled.Pipeline_Workflow{ + {Name: "test-workflow"}, + }, + } + err = d.CreatePipelineEvent(pipelineTid, tpl, nil) + if err != nil { + t.Fatalf("CreatePipelineEvent: %v", err) + } + + _, err = d.Exec(`UPDATE pipelines SET repo_did = ? WHERE id = ?`, repoDid.String(), pipelineTid) + if err != nil { + t.Fatalf("Update pipeline repo association: %v", err) + } + + x := &Xrpc{ + Logger: slog.Default(), + Db: d, + Enforcer: e, + Config: &config.Config{}, + Engines: make(map[string]models.Engine), + } + + sendReq := func(actor syntax.DID, input tangled.CiCancelPipeline_Input) (*httptest.ResponseRecorder, int) { + body, _ := json.Marshal(input) + req := httptest.NewRequest(http.MethodPost, "/com.atproto.repo.createRecord", bytes.NewReader(body)) + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + + w := httptest.NewRecorder() + x.CancelPipeline(w, req) + return w, w.Code + } + + input := tangled.CiCancelPipeline_Input{ + Repo: repoDid.String(), + Pipeline: pipelineTid, + } + + w, code := sendReq(pusherDid, input) + if code != http.StatusOK { + t.Fatalf("expected 200 for pusher, got %d (body: %s)", code, w.Body.String()) + } + + w, code = sendReq(nonPusherDid, input) + if code != http.StatusBadRequest { + t.Fatalf("expected 400 for non-pusher, got %d", code) + } + if !strings.Contains(w.Body.String(), "AccessControl") { + t.Fatalf("expected AccessControl, got: %s", w.Body.String()) + } +} + +type mockDirectory struct { + ident *identity.Identity +} + +func (m *mockDirectory) LookupDID(ctx context.Context, did syntax.DID) (*identity.Identity, error) { + return m.ident, nil +} + +func (m *mockDirectory) LookupHandle(ctx context.Context, handle syntax.Handle) (*identity.Identity, error) { + return m.ident, nil +} + +func (m *mockDirectory) Lookup(ctx context.Context, id syntax.AtIdentifier) (*identity.Identity, error) { + return m.ident, nil +} + +func (m *mockDirectory) Purge(ctx context.Context, id syntax.AtIdentifier) error { + return nil +} + +func TestSecrets_RBAC(t *testing.T) { + d, e := newTestXrpcDB(t) + + 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) + } + + vault, err := secrets.NewSQLiteManager(":memory:") + if err != nil { + t.Fatalf("secrets.NewSQLiteManager: %v", err) + } + + var ts *httptest.Server + ts = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.HasPrefix(r.URL.Path, "/xrpc/com.atproto.repo.getRecord") { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{ + "uri": "at://did:plc:repoowner/sh.tangled.repo/test-repo-rkey", + "cid": "bafybeigdyrzt5s2nuxwos7552", + "value": { + "$type": "sh.tangled.repo", + "knot": "knot.test", + "repoDid": "did:plc:testrepo123", + "spindle": "spindle.test", + "createdAt": "2026-07-26T12:00:00Z" + } + }`)) + return + } + w.WriteHeader(http.StatusNotFound) + })) + defer ts.Close() + + h, err := syntax.ParseHandle("repoowner.test") + if err != nil { + t.Fatalf("syntax.ParseHandle: %v", err) + } + + mockIdent := &identity.Identity{ + DID: repoOwnerDid, + Handle: h, + Services: map[string]identity.ServiceEndpoint{ + "atproto_pds": { + Type: "AtprotoPersonalDataServer", + URL: ts.URL, + }, + }, + } + + resolver := idresolver.NewMockResolver(&mockDirectory{ident: mockIdent}) + + x := &Xrpc{ + Logger: slog.Default(), + Db: d, + Enforcer: e, + Config: &config.Config{}, + Resolver: resolver, + Vault: vault, + } + + addInput := tangled.RepoAddSecret_Input{ + Repo: "at://did:plc:repoowner/sh.tangled.repo/test-repo-rkey", + Key: "MY_SECRET", + Value: "supersecret", + } + + sendAdd := func(actor syntax.DID, input tangled.RepoAddSecret_Input) (*httptest.ResponseRecorder, int) { + body, _ := json.Marshal(input) + req := httptest.NewRequest(http.MethodPost, "/"+tangled.RepoAddSecretNSID, bytes.NewReader(body)) + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + w := httptest.NewRecorder() + x.AddSecret(w, req) + return w, w.Code + } + + w, code := sendAdd(pusherDid, addInput) + if code != http.StatusOK { + t.Fatalf("expected 200 for add secret, got %d (body: %s)", code, w.Body.String()) + } + + w, code = sendAdd(nonPusherDid, addInput) + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for unauthorized add secret, got %d", code) + } + + sendList := func(actor syntax.DID, repo string) (*httptest.ResponseRecorder, int) { + req := httptest.NewRequest(http.MethodGet, "/"+tangled.RepoListSecretsNSID+"?repo="+repo, nil) + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + w := httptest.NewRecorder() + x.ListSecrets(w, req) + return w, w.Code + } + + w, code = sendList(pusherDid, addInput.Repo) + if code != http.StatusOK { + t.Fatalf("expected 200 for list secrets, got %d (body: %s)", code, w.Body.String()) + } + + var listOut tangled.RepoListSecrets_Output + if err := json.Unmarshal(w.Body.Bytes(), &listOut); err != nil { + t.Fatalf("failed to decode list secrets output: %v", err) + } + if len(listOut.Secrets) != 1 || listOut.Secrets[0].Key != "MY_SECRET" { + t.Fatalf("unexpected secrets list: %+v", listOut.Secrets) + } + + w, code = sendList(nonPusherDid, addInput.Repo) + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for unauthorized list secrets, got %d", code) + } + + removeInput := tangled.RepoRemoveSecret_Input{ + Repo: addInput.Repo, + Key: "MY_SECRET", + } + + sendRemove := func(actor syntax.DID, input tangled.RepoRemoveSecret_Input) (*httptest.ResponseRecorder, int) { + body, _ := json.Marshal(input) + req := httptest.NewRequest(http.MethodPost, "/"+tangled.RepoRemoveSecretNSID, bytes.NewReader(body)) + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + w := httptest.NewRecorder() + x.RemoveSecret(w, req) + return w, w.Code + } + + w, code = sendRemove(pusherDid, removeInput) + if code != http.StatusOK { + t.Fatalf("expected 200 for remove secret, got %d (body: %s)", code, w.Body.String()) + } + + w, code = sendList(pusherDid, addInput.Repo) + if code != http.StatusOK { + t.Fatalf("list secrets failed: %d", code) + } + if err := json.Unmarshal(w.Body.Bytes(), &listOut); err != nil { + t.Fatalf("failed to decode: %v", err) + } + if len(listOut.Secrets) != 0 { + t.Fatal("secret was not removed") + } +}