This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

spindle: switch to rbac/v2

spindle members will be stored in db
repo collaborators will be managed by rbac/v2

Signed-off-by: Seongmin Lee <git@boltless.me>

author
Seongmin Lee
date (Jul 29, 2026, 5:25 PM +0900) commit 967247a8 parent 1b2cefa5 change-id wuqwonxy
+398 -532
+1
spindle/config/config.go
··· 24 24 QueueSize int `env:"QUEUE_SIZE, default=100"` 25 25 MaxJobCount int `env:"MAX_JOB_COUNT, default=2"` // max number of pipelines that run at a time 26 26 DockerSocket string `env:"DOCKER_SOCKET"` // path to a docker socket to expose to workflow containers 27 + InviteOnly bool `env:"INVITE_ONLY, default=true"` 27 28 } 28 29 29 30 type Tap struct {
+11
spindle/db/pipelines_test.go
··· 3 3 import ( 4 4 "context" 5 5 "encoding/json" 6 + "path/filepath" 6 7 "slices" 7 8 "testing" 8 9 "time" 9 10 10 11 "tangled.org/core/api/tangled" 11 12 ) 13 + 14 + func newTestDB(t *testing.T) *DB { 15 + t.Helper() 16 + d, err := Make(context.Background(), filepath.Join(t.TempDir(), "spindle.db")) 17 + if err != nil { 18 + t.Fatalf("Make: %v", err) 19 + } 20 + t.Cleanup(func() { d.Close() }) 21 + return d 22 + } 12 23 13 24 func seedPipelineEvent(t *testing.T, d *DB, rkey, repoDid, kind string, created int64) { 14 25 t.Helper()
+20 -64
spindle/db/repos.go
··· 15 15 CreatedAt string 16 16 } 17 17 18 - func (d *DB) AddRepo(repo Repo) error { 18 + func (d *DB) UpsertRepo(repo Repo) error { 19 19 var createdAt sql.NullString 20 20 if repo.CreatedAt != "" { 21 21 createdAt = sql.NullString{String: repo.CreatedAt, Valid: true} 22 22 } 23 23 _, err := d.Exec( 24 - `insert into repos (knot, owner, rkey, repo_did, created_at) 25 - values (?, ?, ?, ?, ?) 26 - on conflict(owner, rkey) do update set 27 - knot = excluded.knot, 28 - repo_did = excluded.repo_did, 29 - created_at = coalesce(excluded.created_at, repos.created_at)`, 30 - repo.Knot, repo.Owner.String(), repo.Rkey.String(), repo.RepoDid.String(), createdAt, 24 + `insert or replace into repos (repo_did, knot, owner, rkey, created_at) 25 + values (?, ?, ?, ?, ?)`, 26 + repo.RepoDid, repo.Knot, repo.Owner, repo.Rkey, createdAt, 31 27 ) 32 28 return err 33 29 } 34 30 35 - func (d *DB) CollapseRepoSiblings(owner, repoDid syntax.DID) (int64, error) { 36 - res, err := d.Exec( 37 - `delete from repos 38 - where owner = ? 39 - and repo_did = ? 40 - and ( 41 - (created_at is null and exists ( 42 - select 1 from repos r2 43 - where r2.owner = repos.owner 44 - and r2.repo_did = repos.repo_did 45 - and r2.created_at is not null 46 - and r2.rkey <> repos.rkey 47 - )) 48 - or (created_at is not null and created_at < ( 49 - select max(created_at) from repos 50 - where owner = ? and repo_did = ? and created_at is not null 51 - )) 52 - )`, 53 - owner.String(), repoDid.String(), owner.String(), repoDid.String(), 54 - ) 31 + func (d *DB) RepoOwners() ([]syntax.DID, error) { 32 + repos, err := d.AllRepos() 55 33 if err != nil { 56 - return 0, err 34 + return nil, err 57 35 } 58 - return res.RowsAffected() 36 + seen := make(map[syntax.DID]struct{}, len(repos)) 37 + dids := make([]syntax.DID, 0, len(repos)) 38 + for _, r := range repos { 39 + if r.Owner == "" { 40 + continue 41 + } 42 + if _, ok := seen[r.Owner]; ok { 43 + continue 44 + } 45 + seen[r.Owner] = struct{}{} 46 + dids = append(dids, r.Owner) 47 + } 48 + return dids, nil 59 49 } 60 50 61 51 func (d *DB) Knots() ([]string, error) { ··· 94 84 }, nil 95 85 } 96 86 97 - func (d *DB) SiblingRkeysForRepoDid(owner, repoDid syntax.DID, excludeRkey syntax.RecordKey) ([]string, error) { 98 - rows, err := d.Query( 99 - `select rkey from repos 100 - where owner = ? 101 - and coalesce(repo_did, '') = ? 102 - and rkey <> ?`, 103 - owner.String(), repoDid.String(), excludeRkey.String(), 104 - ) 105 - if err != nil { 106 - return nil, err 107 - } 108 - defer rows.Close() 109 - 110 - var collect func(acc []string) ([]string, error) 111 - collect = func(acc []string) ([]string, error) { 112 - if !rows.Next() { 113 - return acc, rows.Err() 114 - } 115 - var r string 116 - if err := rows.Scan(&r); err != nil { 117 - return acc, err 118 - } 119 - return collect(append(acc, r)) 120 - } 121 - return collect(nil) 122 - } 123 - 124 87 func (d *DB) GetRepoByDid(repoDid syntax.DID) (*Repo, error) { 125 88 return scanRepo(d.QueryRow( 126 89 `select knot, owner, rkey, repo_did from repos where repo_did = ?`, 127 90 repoDid.String(), 128 - )) 129 - } 130 - 131 - func (d *DB) GetRepoByOwnerRkey(owner syntax.DID, rkey syntax.RecordKey) (*Repo, error) { 132 - return scanRepo(d.QueryRow( 133 - `select knot, owner, rkey, repo_did from repos where owner = ? and rkey = ?`, 134 - owner.String(), rkey.String(), 135 91 )) 136 92 } 137 93
-143
spindle/db/repos_test.go
··· 1 - package db 2 - 3 - import ( 4 - "context" 5 - "path/filepath" 6 - "testing" 7 - 8 - "github.com/bluesky-social/indigo/atproto/syntax" 9 - ) 10 - 11 - func newTestDB(t *testing.T) *DB { 12 - t.Helper() 13 - d, err := Make(context.Background(), filepath.Join(t.TempDir(), "spindle.db")) 14 - if err != nil { 15 - t.Fatalf("Make: %v", err) 16 - } 17 - t.Cleanup(func() { d.Close() }) 18 - return d 19 - } 20 - 21 - func TestCollapseRepoSiblings_DeletesStaleNullCreatedAtWithDifferentRkey(t *testing.T) { 22 - d := newTestDB(t) 23 - owner := syntax.DID("did:plc:akshay") 24 - repoDid := syntax.DID("did:plc:boltless") 25 - 26 - if _, err := d.Exec(`insert into repos (knot, owner, rkey, repo_did, created_at) values 27 - ('k', ?, 'stale-bogus-rkey', ?, null), 28 - ('k', ?, 'fresh-pds-rkey', ?, '2024-06-01T00:00:00Z')`, 29 - owner.String(), repoDid.String(), 30 - owner.String(), repoDid.String()); err != nil { 31 - t.Fatalf("seed: %v", err) 32 - } 33 - 34 - n, err := d.CollapseRepoSiblings(owner, repoDid) 35 - if err != nil { 36 - t.Fatalf("CollapseRepoSiblings: %v", err) 37 - } 38 - if n != 1 { 39 - t.Errorf("expected 1 stale row deleted, got %d", n) 40 - } 41 - 42 - var rkey string 43 - if err := d.QueryRow(`select rkey from repos where owner = ? and repo_did = ?`, 44 - owner.String(), repoDid.String()).Scan(&rkey); err != nil { 45 - t.Fatalf("query: %v", err) 46 - } 47 - if rkey != "fresh-pds-rkey" { 48 - t.Errorf("expected fresh row preserved, got rkey=%q", rkey) 49 - } 50 - } 51 - 52 - func TestCollapseRepoSiblings_KeepsNullCreatedAtWhenAlone(t *testing.T) { 53 - d := newTestDB(t) 54 - owner := syntax.DID("did:plc:akshay") 55 - repoDid := syntax.DID("did:plc:boltless") 56 - 57 - if _, err := d.Exec(`insert into repos (knot, owner, rkey, repo_did, created_at) values 58 - ('k', ?, 'sole-row', ?, null)`, 59 - owner.String(), repoDid.String()); err != nil { 60 - t.Fatalf("seed: %v", err) 61 - } 62 - 63 - n, err := d.CollapseRepoSiblings(owner, repoDid) 64 - if err != nil { 65 - t.Fatalf("CollapseRepoSiblings: %v", err) 66 - } 67 - if n != 0 { 68 - t.Errorf("expected 0 deletions when only NULL row exists, got %d", n) 69 - } 70 - 71 - var count int 72 - if err := d.QueryRow(`select count(*) from repos where owner = ?`, owner.String()).Scan(&count); err != nil { 73 - t.Fatalf("count: %v", err) 74 - } 75 - if count != 1 { 76 - t.Errorf("sole NULL row should survive, got %d remaining", count) 77 - } 78 - } 79 - 80 - func TestCollapseRepoSiblings_OlderTimestampLoses(t *testing.T) { 81 - d := newTestDB(t) 82 - owner := syntax.DID("did:plc:akshay") 83 - repoDid := syntax.DID("did:plc:boltless") 84 - 85 - if _, err := d.Exec(`insert into repos (knot, owner, rkey, repo_did, created_at) values 86 - ('k', ?, 'older-rkey', ?, '2024-01-01T00:00:00Z'), 87 - ('k', ?, 'newer-rkey', ?, '2024-06-01T00:00:00Z')`, 88 - owner.String(), repoDid.String(), 89 - owner.String(), repoDid.String()); err != nil { 90 - t.Fatalf("seed: %v", err) 91 - } 92 - 93 - n, err := d.CollapseRepoSiblings(owner, repoDid) 94 - if err != nil { 95 - t.Fatalf("CollapseRepoSiblings: %v", err) 96 - } 97 - if n != 1 { 98 - t.Errorf("expected older row collapsed, got %d", n) 99 - } 100 - 101 - var rkey string 102 - if err := d.QueryRow(`select rkey from repos where owner = ? and repo_did = ?`, 103 - owner.String(), repoDid.String()).Scan(&rkey); err != nil { 104 - t.Fatalf("query: %v", err) 105 - } 106 - if rkey != "newer-rkey" { 107 - t.Errorf("expected newer row preserved, got rkey=%q", rkey) 108 - } 109 - } 110 - 111 - func TestCollapseRepoSiblings_KeepsNullRowWithMatchingRkey(t *testing.T) { 112 - d := newTestDB(t) 113 - owner := syntax.DID("did:plc:akshay") 114 - repoDid := syntax.DID("did:plc:boltless") 115 - 116 - if _, err := d.Exec(`insert into repos (knot, owner, rkey, repo_did, created_at) values 117 - ('k', ?, 'matched-rkey', ?, null)`, 118 - owner.String(), repoDid.String()); err != nil { 119 - t.Fatalf("seed: %v", err) 120 - } 121 - 122 - if err := d.AddRepo(Repo{ 123 - Knot: "k", 124 - Owner: owner, 125 - Rkey: "matched-rkey", 126 - RepoDid: repoDid, 127 - CreatedAt: "2024-06-01T00:00:00Z", 128 - }); err != nil { 129 - t.Fatalf("AddRepo upsert: %v", err) 130 - } 131 - 132 - if _, err := d.CollapseRepoSiblings(owner, repoDid); err != nil { 133 - t.Fatalf("CollapseRepoSiblings: %v", err) 134 - } 135 - 136 - var count int 137 - if err := d.QueryRow(`select count(*) from repos where owner = ?`, owner.String()).Scan(&count); err != nil { 138 - t.Fatalf("count: %v", err) 139 - } 140 - if count != 1 { 141 - t.Errorf("upserted row should be the single survivor, got %d", count) 142 - } 143 - }
+2 -1
spindle/embedtap.go
··· 63 63 RepoFetchTimeout: 5 * time.Minute, 64 64 IdentityCacheSize: 50_000, 65 65 EventCacheSize: 10_000, 66 - CollectionFilters: []string{tangled.RepoNSID, tangled.RepoCollaboratorNSID}, 66 + FullNetworkMode: !cfg.Server.InviteOnly, 67 + CollectionFilters: []string{tangled.RepoNSID}, 67 68 AdminPassword: cfg.Server.Tap.AdminPassword, 68 69 RetryTimeout: 60 * time.Second, 69 70 }
+50 -57
spindle/server.go
··· 11 11 "maps" 12 12 "net/http" 13 13 "path/filepath" 14 + "slices" 14 15 "sync" 15 16 "time" 16 17 ··· 119 120 // pull records are created by arbitrary users too, same hack as in tap 120 121 jc.ExemptCollection(tangled.RepoPullNSID) 121 122 122 - // Check if the spindle knows about any Dids; 123 - dids, err := d.ListAllowedMembers() 124 - if err != nil { 125 - return nil, fmt.Errorf("failed to get all dids: %w", err) 126 - } 127 - for _, d := range dids { 128 - jc.AddDid(d.String()) 129 - } 130 - 131 - knownRepos, err := d.AllRepos() 132 - if err != nil { 133 - return nil, fmt.Errorf("failed to get known repos: %w", err) 134 - } 135 - for _, r := range knownRepos { 136 - if r.Owner != "" { 137 - jc.AddDid(r.Owner.String()) 123 + if cfg.Server.InviteOnly { 124 + // listen to members and the collaborators on the repos we host 125 + dids, err := subscribedDids(d, e) 126 + if err != nil { 127 + return nil, fmt.Errorf("failed to build jetstream did filter: %w", err) 128 + } 129 + for _, did := range dids { 130 + jc.AddDid(did.String()) 138 131 } 132 + } else { 133 + // public spindle. listen to full network 134 + jc.ExemptCollection(tangled.RepoNSID) 139 135 } 140 136 141 137 resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl) ··· 236 232 return s.verify(ctx, repoident.RepoDid(repo)) 237 233 } 238 234 235 + // subscribedDids lists all allowed spindle members & all collaborators of registered repos 236 + func subscribedDids(d *db.DB, e *rbac.Enforcer) ([]syntax.DID, error) { 237 + members, err := d.ListAllowedMembers() 238 + if err != nil { 239 + return nil, fmt.Errorf("list members: %w", err) 240 + } 241 + repos, err := d.AllRepos() 242 + if err != nil { 243 + return nil, fmt.Errorf("list repos: %w", err) 244 + } 245 + 246 + dids := slices.Clone(members) 247 + for _, r := range repos { 248 + // includes the repo owner, via the repo:owner -> repo:collaborator grouping 249 + collaborators, err := e.GetRepoCollaborators(r.RepoDid) 250 + if err != nil { 251 + return nil, fmt.Errorf("list collaborators of %s: %w", r.RepoDid, err) 252 + } 253 + dids = append(dids, collaborators...) 254 + } 255 + 256 + slices.Sort(dids) 257 + return slices.Compact(dids), nil 258 + } 259 + 260 + func (s *Spindle) grantCollaborator(subject, repo syntax.DID) error { 261 + if err := s.e.AddRepoCollaborator(subject, repo); err != nil { 262 + return err 263 + } 264 + s.jc.AddDid(subject.String()) 265 + return nil 266 + } 267 + 239 268 // SetMotdContent sets custom MOTD content, replacing the embedded default. 240 269 func (s *Spindle) SetMotdContent(content []byte) { 241 270 s.motdMu.Lock() ··· 289 318 290 319 s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr) 291 320 return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router()) 292 - } 293 - 294 - func (s *Spindle) declareTapInterest(ctx context.Context) { 295 - repos, err := s.db.AllRepos() 296 - if err != nil { 297 - s.l.Warn("tap declare: failed to load known repos", "err", err) 298 - return 299 - } 300 - seen := make(map[syntax.DID]struct{}, len(repos)) 301 - dids := make([]syntax.DID, 0, len(repos)) 302 - for _, r := range repos { 303 - if r.Owner == "" { 304 - continue 305 - } 306 - if _, ok := seen[r.Owner]; ok { 307 - continue 308 - } 309 - seen[r.Owner] = struct{}{} 310 - dids = append(dids, r.Owner) 311 - } 312 - if err := s.tap.AddOwnerDIDs(ctx, dids); err != nil { 313 - s.l.Warn("tap declare: AddRepos rejected", "count", len(dids), "err", err) 314 - return 315 - } 316 - s.l.Info("tap declare: known owner DIDs registered", "count", len(dids)) 317 321 } 318 322 319 323 func Run(ctx context.Context) error { ··· 442 446 return nil 443 447 } 444 448 445 - func (s *Spindle) ingestKnotCollaborator(ctx context.Context, l *slog.Logger, src eventconsumer.Source, msg eventstream.Event) error { 449 + func (s *Spindle) ingestKnotCollaborator(_ context.Context, l *slog.Logger, src eventconsumer.Source, msg eventstream.Event) error { 446 450 var rec knotdb.RepoCollaboratorUpdate 447 451 if err := json.Unmarshal(msg.EventJson, &rec); err != nil { 448 452 l.Error("error unmarshalling collaboratorUpdate", "err", err) ··· 475 479 476 480 switch rec.Op { 477 481 case knotdb.AclOpAdd: 478 - if err := s.e.AddRepoCollaborator(subject, repoDid); err != nil { 482 + if err := s.grantCollaborator(subject, repoDid); err != nil { 479 483 return fmt.Errorf("add collaborator policy: %w", err) 480 484 } 481 485 l.Info("added knot-managed collaborator", "subject", subject, "repo", repoDid) 482 486 case knotdb.AclOpRemove: 487 + // ponytail: no jc.RemoveDid here - the subject may still be a member or a 488 + // collaborator elsewhere, and a stale filter entry is harmless (processRepo still 489 + // rejects non-members). Add refcounting across members + acl_2 if the filter grows. 483 490 if err := s.e.RemoveRepoCollaborator(subject, repoDid); err != nil { 484 491 return fmt.Errorf("remove collaborator policy: %w", err) 485 492 } ··· 853 860 } 854 861 return nil 855 862 } 856 - 857 - func (s *Spindle) resolvePipelineRepoDid(repo *tangled.Pipeline_TriggerRepo) (syntax.DID, error) { 858 - if repo.RepoDid == nil || *repo.RepoDid == "" { 859 - return "", fmt.Errorf("pipeline trigger missing repoDid") 860 - } 861 - repoDid, err := syntax.ParseDID(*repo.RepoDid) 862 - if err != nil { 863 - return "", fmt.Errorf("parse repoDid %s: %w", *repo.RepoDid, err) 864 - } 865 - if _, err := s.db.GetRepoByDid(repoDid); err != nil { 866 - return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) 867 - } 868 - return repoDid, nil 869 - }
+50
spindle/server_test.go
··· 1 1 package spindle 2 2 3 3 import ( 4 + "slices" 4 5 "testing" 5 6 7 + "github.com/bluesky-social/indigo/atproto/syntax" 6 8 kgit "tangled.org/core/knotserver/git" 9 + "tangled.org/core/spindle/db" 7 10 ) 8 11 9 12 func TestHasSkipCIPushOption(t *testing.T) { ··· 53 56 }) 54 57 } 55 58 } 59 + 60 + func TestSubscribedDids(t *testing.T) { 61 + d, e := newTestSpindleDB(t) 62 + 63 + var ( 64 + member = syntax.DID("did:plc:member") 65 + owner = syntax.DID("did:plc:owner") 66 + collab = syntax.DID("did:plc:collab") 67 + repoDid = syntax.DID("did:plc:repo") 68 + repo2Did = syntax.DID("did:plc:repo2") 69 + ) 70 + 71 + // the repo owner is also a member, so the union has to dedupe 72 + for _, did := range []syntax.DID{member, owner} { 73 + if err := d.AllowMember(t.Context(), did); err != nil { 74 + t.Fatalf("AllowMember(%s): %v", did, err) 75 + } 76 + } 77 + 78 + for _, r := range []syntax.DID{repoDid, repo2Did} { 79 + if err := d.UpsertRepo(db.Repo{ 80 + Knot: "knot.test", 81 + Owner: owner, 82 + Rkey: syntax.RecordKey("rkey-" + r.String()), 83 + RepoDid: r, 84 + }); err != nil { 85 + t.Fatalf("UpsertRepo(%s): %v", r, err) 86 + } 87 + if err := e.SetRepoOwner(owner, r); err != nil { 88 + t.Fatalf("SetRepoOwner(%s): %v", r, err) 89 + } 90 + } 91 + // a collaborator who is not a member - the case the members-only filter dropped 92 + if err := e.AddRepoCollaborator(collab, repoDid); err != nil { 93 + t.Fatalf("AddRepoCollaborator: %v", err) 94 + } 95 + 96 + got, err := subscribedDids(d, e) 97 + if err != nil { 98 + t.Fatalf("subscribedDids: %v", err) 99 + } 100 + 101 + want := []syntax.DID{collab, member, owner} // sorted and deduped 102 + if !slices.Equal(got, want) { 103 + t.Errorf("subscribedDids() = %v, want %v", got, want) 104 + } 105 + }
+19 -16
spindle/tapclient.go
··· 44 44 } 45 45 } 46 46 47 - func (t *Tap) AddOwnerDIDs(ctx context.Context, dids []syntax.DID) error { 48 - if len(dids) == 0 { 49 - return nil 50 - } 51 - return t.tap.AddRepos(ctx, dids) 52 - } 53 - 54 47 func (t *Tap) Start(connCtx context.Context) { 55 48 go t.tap.Connect(connCtx, &tapc.SimpleIndexer{ 56 49 EventHandler: t.processEvent, ··· 59 52 } 60 53 61 54 func (t *Tap) onConnect(ctx context.Context) { 62 - t.spindle.declareTapInterest(ctx) 55 + l := t.logger 56 + if t.spindle.cfg.Server.InviteOnly { 57 + // listen to owners of registered repositories 58 + owners, err := t.spindle.db.RepoOwners() 59 + if err != nil { 60 + l.Warn("tap declare: failed to load known repos", "err", err) 61 + return 62 + } 63 + if err := t.tap.AddRepos(ctx, owners); err != nil { 64 + l.Warn("tap declare: AddRepos rejected", "count", len(owners), "err", err) 65 + return 66 + } 67 + l.Info("tap declare: known owner DIDs registered", "count", len(owners)) 68 + } else { 69 + // public spindle. listen to full network 70 + l.Info("tap declare: listening to full network") 71 + } 63 72 } 64 73 65 74 func (t *Tap) processEvent(ctx context.Context, evt tapc.Event) error { ··· 98 107 return nil 99 108 } 100 109 101 - isMember, err := t.spindle.e.IsSpindleMember(ownerDid.String(), rbac.ThisServer) 110 + isMember, err := t.spindle.db.IsAllowedMember(ctx, ownerDid, !t.spindle.cfg.Server.InviteOnly) 102 111 if err != nil { 103 112 return fmt.Errorf("checking spindle membership: %w", err) 104 113 } ··· 150 159 CreatedAt: record.CreatedAt, 151 160 } 152 161 153 - if err := t.spindle.db.AddRepo(repo); err != nil { 162 + if err := t.spindle.db.UpsertRepo(repo); err != nil { 154 163 l.Error("failed to add repo row", "err", err) 155 164 return fmt.Errorf("add repo: %w", err) 156 165 } ··· 170 179 legacyName = *record.Name 171 180 } 172 181 migrateLegacyRepoSecrets(ctx, t.spindle.db, t.spindle.vault, l, ownerDid, legacyName, rkey, repoDid) 173 - 174 - if removed, err := t.spindle.db.CollapseRepoSiblings(ownerDid, repoDid); err != nil { 175 - l.Warn("collapse rename siblings failed", "err", err) 176 - } else if removed > 0 { 177 - l.Info("collapsed rename leftovers", "owner", ownerDid, "repo_did", repoDid, "removed", removed) 178 - } 179 182 180 183 if e := t.spindle.embedTap; e == nil || !e.closed.Load() { 181 184 if err := t.tap.AddRepos(ctx, []syntax.DID{ownerDid}); err != nil {
+232 -240
spindle/tapclient_test.go
··· 4 4 "context" 5 5 "encoding/json" 6 6 "log/slog" 7 + "net/http" 8 + "net/http/httptest" 9 + "net/url" 7 10 "strings" 8 - "tangled.org/core/jetstream" 9 11 "testing" 10 12 "time" 11 13 14 + "tangled.org/core/jetstream" 15 + "tangled.org/core/repoident" 16 + "tangled.org/core/repoverify" 17 + 12 18 "github.com/bluesky-social/indigo/atproto/identity" 13 19 "github.com/bluesky-social/indigo/atproto/syntax" 20 + "github.com/stretchr/testify/assert" 14 21 "tangled.org/core/api/tangled" 15 22 "tangled.org/core/eventconsumer" 16 23 "tangled.org/core/idresolver" 17 - "tangled.org/core/rbac" 18 24 "tangled.org/core/spindle/config" 19 25 "tangled.org/core/spindle/db" 20 26 ··· 41 47 return nil 42 48 } 43 49 50 + func mockRepoVerifier(res repoverify.Result) repoverify.Verifier { 51 + return func(ctx context.Context, repoDid repoident.RepoDid) (repoverify.Result, error) { 52 + return res, nil 53 + } 54 + } 55 + 56 + func TestProcessRepo_MembershipChange(t *testing.T) { 57 + member := syntax.DID("did:example:foo") 58 + 59 + var ok bool 60 + var err error 61 + d, _ := newTestSpindleDB(t) 62 + 63 + ok, err = d.IsAllowedMember(t.Context(), member, false) 64 + assert.NoError(t, err) 65 + assert.False(t, ok) 66 + 67 + ok, err = d.IsAllowedMember(t.Context(), member, true) 68 + assert.NoError(t, err) 69 + assert.True(t, ok) 70 + 71 + assert.NoError(t, d.AllowMember(t.Context(), member)) 72 + 73 + ok, err = d.IsAllowedMember(t.Context(), member, false) 74 + assert.NoError(t, err) 75 + assert.True(t, ok) 76 + 77 + ok, err = d.IsAllowedMember(t.Context(), member, true) 78 + assert.NoError(t, err) 79 + assert.True(t, ok) 80 + } 81 + 44 82 func TestProcessRepo_MembershipCheck(t *testing.T) { 45 83 d, e := newTestSpindleDB(t) 46 84 47 85 cfg := &config.Config{} 48 86 cfg.Server.Hostname = "spindle.test" 87 + cfg.Server.InviteOnly = true 49 88 50 89 ccfg := eventconsumer.NewConsumerConfig() 51 90 ccfg.Logger = slog.Default() ··· 63 102 ks: ks, 64 103 jc: jc, 65 104 rootCtx: context.Background(), 105 + verify: mockRepoVerifier(repoverify.Result{ 106 + RepoDid: "did:plc:testrepo123", 107 + OwnerDid: "did:plc:memberowner", 108 + Rkey: "test-repo-rkey", 109 + KnotURL: func() *url.URL { 110 + u, err := url.Parse("https://knot.test") 111 + assert.NoError(t, err) 112 + return u 113 + }(), 114 + }), 66 115 } 67 116 68 117 tap := &Tap{ ··· 72 121 73 122 ownerDid := syntax.DID("did:plc:memberowner") 74 123 nonMemberDid := syntax.DID("did:plc:nonmemberowner") 75 - repoDid := "did:plc:testrepo123" 124 + repoDid := syntax.DID("did:plc:testrepo123") 76 125 77 - err := e.AddSpindle(rbac.ThisServer) 78 - if err != nil { 79 - t.Fatalf("AddSpindle: %v", err) 80 - } 81 - err = e.AddSpindleMember(rbac.ThisServer, ownerDid.String()) 126 + err := d.AllowMember(t.Context(), ownerDid) 82 127 if err != nil { 83 128 t.Fatalf("AddSpindleMember: %v", err) 84 129 } 85 130 86 131 recNonMember := tangled.Repo{ 87 132 Knot: "knot.test", 88 - RepoDid: &repoDid, 133 + RepoDid: (*string)(&repoDid), 89 134 Spindle: &cfg.Server.Hostname, 90 135 CreatedAt: time.Now().Format(time.RFC3339), 91 136 } ··· 103 148 t.Fatalf("processRepo returned error for non-member: %v", err) 104 149 } 105 150 106 - _, err = d.GetRepoByOwnerRkey(nonMemberDid, "test-repo-rkey") 151 + _, err = d.GetRepoByDid(repoDid) 107 152 if err == nil { 108 153 t.Fatal("repo for non-member was registered in DB, expected rejection") 109 154 } 110 155 111 156 recMember := tangled.Repo{ 112 157 Knot: "knot.test", 113 - RepoDid: &repoDid, 158 + RepoDid: (*string)(&repoDid), 114 159 Spindle: &cfg.Server.Hostname, 115 160 CreatedAt: time.Now().Format(time.RFC3339), 116 161 } ··· 133 178 } 134 179 } 135 180 136 - func TestProcessPull_PushAllowedCheck(t *testing.T) { 181 + func TestProcessPull_IsCollaboratorCheck(t *testing.T) { 137 182 d, e := newTestSpindleDB(t) 138 183 139 184 cfg := &config.Config{} ··· 151 196 res: idresolver.DefaultResolver("https://plc.test"), 152 197 jc: jc, 153 198 rootCtx: context.Background(), 199 + verify: mockRepoVerifier(repoverify.Result{ 200 + RepoDid: "did:plc:testrepo123", 201 + OwnerDid: "did:plc:repoowner", 202 + Rkey: "test-repo-rkey", 203 + KnotURL: func() *url.URL { 204 + u, _ := url.Parse("knot.test") 205 + return u 206 + }(), 207 + }), 154 208 } 155 209 156 210 repoOwnerDid := syntax.DID("did:plc:repoowner") ··· 158 212 pusherDid := syntax.DID("did:plc:pusher") 159 213 repoDid := syntax.DID("did:plc:testrepo123") 160 214 161 - err := d.AddRepo(db.Repo{ 215 + err := d.UpsertRepo(db.Repo{ 162 216 Knot: "knot.test", 163 217 Owner: repoOwnerDid, 164 218 Rkey: "test-repo-rkey", ··· 169 223 t.Fatalf("AddRepo: %v", err) 170 224 } 171 225 172 - err = e.AddRepo(repoOwnerDid.String(), rbac.ThisServer, repoDid.String()) 226 + err = e.SetRepoOwner(repoOwnerDid, repoDid) 173 227 if err != nil { 174 228 t.Fatalf("AddRepo permissions: %v", err) 175 229 } 176 - err = e.AddCollaborator(pusherDid.String(), rbac.ThisServer, repoDid.String()) 230 + err = e.AddRepoCollaborator(pusherDid, repoDid) 177 231 if err != nil { 178 232 t.Fatalf("AddCollaborator: %v", err) 179 233 } ··· 242 296 ks: ks, 243 297 jc: jc, 244 298 rootCtx: context.Background(), 299 + verify: mockRepoVerifier(repoverify.Result{ 300 + RepoDid: "did:plc:sharedrepo", 301 + OwnerDid: "did:plc:alice", 302 + Rkey: "alice-repo", 303 + KnotURL: func() *url.URL { 304 + u, _ := url.Parse("knot.test") 305 + return u 306 + }(), 307 + }), 245 308 } 246 309 247 310 tap := &Tap{ ··· 251 314 252 315 aliceDid := syntax.DID("did:plc:alice") 253 316 bobDid := syntax.DID("did:plc:bob") 254 - repoDid := "did:plc:sharedrepo" 317 + repoDid := syntax.DID("did:plc:sharedrepo") 255 318 256 - err := e.AddSpindle(rbac.ThisServer) 257 - if err != nil { 258 - t.Fatalf("AddSpindle: %v", err) 259 - } 260 - err = e.AddSpindleMember(rbac.ThisServer, aliceDid.String()) 319 + err := d.AllowMember(t.Context(), aliceDid) 261 320 if err != nil { 262 321 t.Fatalf("AddSpindleMember alice: %v", err) 263 322 } 264 - err = e.AddSpindleMember(rbac.ThisServer, bobDid.String()) 323 + err = d.AllowMember(t.Context(), bobDid) 265 324 if err != nil { 266 325 t.Fatalf("AddSpindleMember bob: %v", err) 267 326 } 268 327 269 - err = d.AddRepo(db.Repo{ 328 + err = d.UpsertRepo(db.Repo{ 270 329 Knot: "knot.test", 271 330 Owner: aliceDid, 272 331 Rkey: "alice-repo", ··· 277 336 t.Fatalf("d.AddRepo: %v", err) 278 337 } 279 338 339 + if err := e.SetRepoOwner(aliceDid, repoDid); err != nil { 340 + t.Fatalf("SetRepoOwner: %v", err) 341 + } 342 + 280 343 // bob tries to register alice's repo did, must reject the hijack 281 344 recBob := tangled.Repo{ 282 345 Knot: "knot.test", 283 - RepoDid: &repoDid, 346 + RepoDid: (*string)(&repoDid), 284 347 Spindle: &cfg.Server.Hostname, 285 348 CreatedAt: time.Now().Format(time.RFC3339), 286 349 } ··· 298 361 t.Fatalf("processRepo returned error on duplicate repoDid hijack attempt: %v", err) 299 362 } 300 363 301 - _, err = d.GetRepoByOwnerRkey(bobDid, "bob-repo") 302 - if err == nil { 303 - t.Fatal("bob successfully hijacked alice's repoDid in DB, expected rejection") 304 - } 305 - } 306 - 307 - func TestProcessCollaborator_RBAC(t *testing.T) { 308 - d, e := newTestSpindleDB(t) 309 - 310 - cfg := &config.Config{} 311 - cfg.Server.Hostname = "spindle.test" 312 - 313 - ownerDid := syntax.DID("did:plc:repoowner") 314 - otherDid := syntax.DID("did:plc:otheractor") 315 - subjectDid := syntax.DID("did:plc:collabsubject") 316 - repoDid := syntax.DID("did:plc:testrepo123") 317 - 318 - h, err := syntax.ParseHandle("collabsubject.test") 364 + stored, err := d.GetRepoByDid(repoDid) 319 365 if err != nil { 320 - t.Fatalf("syntax.ParseHandle: %v", err) 366 + t.Fatalf("alice's repo row was destroyed by bob's hijack attempt: %v", err) 321 367 } 322 - mockIdent := &identity.Identity{ 323 - DID: subjectDid, 324 - Handle: h, 368 + if stored.Owner != aliceDid || stored.Rkey != "alice-repo" { 369 + t.Fatalf("bob hijacked alice's repoDid: owner=%s rkey=%s", stored.Owner, stored.Rkey) 325 370 } 326 - resolver := idresolver.NewMockResolver(&mockDirectory{ident: mockIdent}) 327 371 328 - jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) 329 - if jcerr != nil { 330 - t.Fatalf("NewJetstreamClient: %v", jcerr) 331 - } 332 - s := &Spindle{ 333 - db: d, 334 - e: e, 335 - l: slog.Default(), 336 - cfg: cfg, 337 - res: resolver, 338 - jc: jc, 339 - rootCtx: context.Background(), 340 - } 341 - 342 - tap := &Tap{ 343 - spindle: s, 344 - logger: slog.Default(), 345 - } 346 - 347 - err = d.AddRepo(db.Repo{ 372 + // bob points the same repoDid at another spindle, which must not tear alice's repo down 373 + otherSpindle := "other.test" 374 + recTeardown := tangled.Repo{ 348 375 Knot: "knot.test", 349 - Owner: ownerDid, 350 - Rkey: "test-repo-rkey", 351 - RepoDid: repoDid, 376 + RepoDid: (*string)(&repoDid), 377 + Spindle: &otherSpindle, 352 378 CreatedAt: time.Now().Format(time.RFC3339), 353 - }) 354 - if err != nil { 355 - t.Fatalf("AddRepo: %v", err) 356 379 } 357 - 358 - collabRecord := tangled.RepoCollaborator{ 359 - Subject: subjectDid.String(), 360 - Repo: repoDid.String(), 361 - } 362 - collabRecordJson, _ := json.Marshal(collabRecord) 363 - 364 - err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ 365 - Live: true, 366 - Did: otherDid, 367 - Rkey: "collab-rkey-1", 368 - Collection: syntax.NSID(tangled.RepoCollaboratorNSID), 369 - Action: tapc.RecordCreateAction, 370 - Record: collabRecordJson, 371 - }) 372 - if err != nil { 373 - t.Fatalf("processCollaborator returned error: %v", err) 374 - } 375 - 376 - _, err = d.GetRepoCollaborator(otherDid, "collab-rkey-1") 377 - if err == nil { 378 - t.Fatal("collaborator from non-owner was registered in DB") 379 - } 380 + recTeardownJson, _ := json.Marshal(recTeardown) 380 381 381 - err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ 382 + err = tap.processRepo(context.Background(), &tapc.RecordEventData{ 382 383 Live: true, 383 - Did: ownerDid, 384 - Rkey: "collab-rkey-2", 385 - Collection: syntax.NSID(tangled.RepoCollaboratorNSID), 386 - Action: tapc.RecordCreateAction, 387 - Record: collabRecordJson, 384 + Did: bobDid, 385 + Rkey: "bob-repo", 386 + Collection: syntax.NSID(tangled.RepoNSID), 387 + Action: tapc.RecordUpdateAction, 388 + Record: recTeardownJson, 388 389 }) 389 390 if err != nil { 390 - t.Fatalf("processCollaborator returned error: %v", err) 391 - } 392 - _, err = d.GetRepoCollaborator(ownerDid, "collab-rkey-2") 393 - if err == nil { 394 - t.Fatal("collaborator registered despite missing Casbin invite permission") 391 + t.Fatalf("processRepo returned error on forged teardown: %v", err) 395 392 } 396 393 397 - err = e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()) 398 - if err != nil { 399 - t.Fatalf("AddRepo permissions: %v", err) 400 - } 401 - 402 - err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ 403 - Live: true, 404 - Did: ownerDid, 405 - Rkey: "collab-rkey-3", 406 - Collection: syntax.NSID(tangled.RepoCollaboratorNSID), 407 - Action: tapc.RecordCreateAction, 408 - Record: collabRecordJson, 409 - }) 410 - if err != nil { 411 - t.Fatalf("processCollaborator failed for authorized owner: %v", err) 412 - } 413 - 414 - c, err := d.GetRepoCollaborator(ownerDid, "collab-rkey-3") 415 - if err != nil { 416 - t.Fatalf("GetRepoCollaborator error: %v", err) 417 - } 418 - if c.Subject != subjectDid || c.RepoDid != repoDid { 419 - t.Fatalf("unexpected collaborator: %+v", c) 394 + if _, err := d.GetRepoByDid(repoDid); err != nil { 395 + t.Fatalf("bob tore down alice's repo by naming her repoDid: %v", err) 420 396 } 421 - 422 - ok, err := e.IsRepoCollaborator(subjectDid.String(), rbac.ThisServer, repoDid.String()) 397 + ok, err := e.IsRepoOwner(aliceDid, repoDid) 423 398 if err != nil || !ok { 424 - t.Fatalf("Casbin policy for collaborator missing or err: %v", err) 425 - } 426 - 427 - err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ 428 - Live: true, 429 - Did: ownerDid, 430 - Rkey: "collab-rkey-3", 431 - Collection: syntax.NSID(tangled.RepoCollaboratorNSID), 432 - Action: tapc.RecordDeleteAction, 433 - }) 434 - if err != nil { 435 - t.Fatalf("delete collaborator process returned error: %v", err) 436 - } 437 - 438 - _, err = d.GetRepoCollaborator(ownerDid, "collab-rkey-3") 439 - if err == nil { 440 - t.Fatal("collaborator DB row remained after deletion") 441 - } 442 - 443 - ok, err = e.IsRepoCollaborator(subjectDid.String(), rbac.ThisServer, repoDid.String()) 444 - if err != nil || ok { 445 - t.Fatal("Casbin policy for collaborator remained after deletion") 399 + t.Fatal("bob's forged teardown removed alice's owner policy") 446 400 } 447 401 } 448 402 ··· 463 417 cfg: cfg, 464 418 jc: jc, 465 419 rootCtx: context.Background(), 420 + verify: mockRepoVerifier(repoverify.Result{ 421 + RepoDid: "did:plc:testrepo123", 422 + OwnerDid: "did:plc:repoowner", 423 + Rkey: "test-repo-rkey", 424 + KnotURL: func() *url.URL { 425 + u, _ := url.Parse("knot.test") 426 + return u 427 + }(), 428 + }), 466 429 } 467 430 468 431 tap := &Tap{ ··· 474 437 repoDid := syntax.DID("did:plc:testrepo123") 475 438 collabDid := syntax.DID("did:plc:collab") 476 439 477 - err := d.AddRepo(db.Repo{ 440 + err := d.UpsertRepo(db.Repo{ 478 441 Knot: "knot.test", 479 442 Owner: ownerDid, 480 443 Rkey: "test-repo-rkey", ··· 485 448 t.Fatalf("AddRepo DB: %v", err) 486 449 } 487 450 488 - err = e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()) 451 + err = e.SetRepoOwner(ownerDid, repoDid) 489 452 if err != nil { 490 453 t.Fatalf("AddRepo policy: %v", err) 491 454 } 492 455 493 - err = d.AddRepoCollaborator(db.RepoCollaborator{ 494 - OwnerDid: ownerDid, 495 - Rkey: "collab-rkey", 496 - Subject: collabDid, 497 - RepoDid: repoDid, 498 - }) 499 - if err != nil { 500 - t.Fatalf("AddCollaborator DB: %v", err) 501 - } 502 - 503 - err = e.AddCollaborator(collabDid.String(), rbac.ThisServer, repoDid.String()) 456 + err = e.AddRepoCollaborator(collabDid, repoDid) 504 457 if err != nil { 505 458 t.Fatalf("AddCollaborator policy: %v", err) 506 459 } ··· 516 469 t.Fatalf("processRepo delete returned error: %v", err) 517 470 } 518 471 519 - _, err = d.GetRepoByOwnerRkey(ownerDid, "test-repo-rkey") 472 + _, err = d.GetRepoByDid(repoDid) 520 473 if err == nil { 521 474 t.Fatal("repo remained in DB after delete") 522 475 } 523 476 524 - collabs, err := d.ListCollaboratorsByRepoDid(repoDid) 525 - if err != nil { 526 - t.Fatalf("ListCollaboratorsByRepoDid: %v", err) 527 - } 528 - if len(collabs) > 0 { 529 - t.Fatal("collaborators remained in DB after delete") 530 - } 531 - 532 - ok, err := e.IsRepoOwner(ownerDid.String(), rbac.ThisServer, repoDid.String()) 477 + ok, err := e.IsRepoOwner(ownerDid, repoDid) 533 478 if err != nil || ok { 534 479 t.Fatal("repo owner policy remained in Casbin after delete") 535 480 } 536 481 537 - ok, err = e.IsRepoCollaborator(collabDid.String(), rbac.ThisServer, repoDid.String()) 482 + ok, err = e.IsRepoCollaborator(collabDid, repoDid) 538 483 if err != nil || ok { 539 484 t.Fatal("collaborator policy remained in Casbin after delete") 540 485 } ··· 558 503 cfg: cfg, 559 504 jc: jc, 560 505 rootCtx: context.Background(), 506 + verify: mockRepoVerifier(repoverify.Result{ 507 + RepoDid: "did:plc:sharedrepo", 508 + OwnerDid: "did:plc:alice", 509 + Rkey: "test-repo-rkey", 510 + KnotURL: func() *url.URL { 511 + u, _ := url.Parse("knot.test") 512 + return u 513 + }(), 514 + }), 561 515 } 562 516 563 517 tap := &Tap{ ··· 569 523 bobDid := syntax.DID("did:plc:bob") 570 524 repoDid := syntax.DID("did:plc:sharedrepo") 571 525 572 - err := d.AddRepo(db.Repo{ 526 + err := d.UpsertRepo(db.Repo{ 573 527 Knot: "knot.test", 574 528 Owner: aliceDid, 575 529 Rkey: "test-repo-rkey", ··· 580 534 t.Fatalf("AddRepo DB: %v", err) 581 535 } 582 536 583 - err = e.AddRepo(aliceDid.String(), rbac.ThisServer, repoDid.String()) 537 + err = e.SetRepoOwner(aliceDid, repoDid) 584 538 if err != nil { 585 539 t.Fatalf("AddRepo policy: %v", err) 586 540 } ··· 597 551 t.Fatalf("processRepo returned error on delete: %v", err) 598 552 } 599 553 600 - _, err = d.GetRepoByOwnerRkey(aliceDid, "test-repo-rkey") 554 + _, err = d.GetRepoByDid(repoDid) 601 555 if err != nil { 602 556 t.Fatalf("Alice's repo was deleted or error: %v", err) 603 557 } 604 558 605 - ok, err := e.IsRepoOwner(aliceDid.String(), rbac.ThisServer, repoDid.String()) 559 + ok, err := e.IsRepoOwner(aliceDid, repoDid) 606 560 if err != nil || !ok { 607 561 t.Fatal("Alice's owner policy was removed from Casbin by forged delete") 608 562 } 609 563 } 610 564 611 - func TestProcessCollaborator_ForgeDeleteRejection(t *testing.T) { 565 + func TestReconcileCollaborators(t *testing.T) { 612 566 d, e := newTestSpindleDB(t) 613 567 614 - cfg := &config.Config{} 615 - cfg.Server.Hostname = "spindle.test" 568 + ownerDid := syntax.DID("did:plc:owner") 569 + repoDid := syntax.DID("did:plc:repo") 570 + staleDid := syntax.DID("did:plc:stale") 571 + keptDid := syntax.DID("did:plc:kept") 572 + newDid := syntax.DID("did:plc:new") 616 573 617 - jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) 618 - if jcerr != nil { 619 - t.Fatalf("NewJetstreamClient: %v", jcerr) 574 + if err := e.SetRepoOwner(ownerDid, repoDid); err != nil { 575 + t.Fatalf("SetRepoOwner: %v", err) 576 + } 577 + // spindle's view: one collaborator the knot dropped, one it still has 578 + for _, did := range []syntax.DID{staleDid, keptDid} { 579 + if err := e.AddRepoCollaborator(did, repoDid); err != nil { 580 + t.Fatalf("AddRepoCollaborator(%s): %v", did, err) 581 + } 620 582 } 621 583 622 - s := &Spindle{ 623 - db: d, 624 - e: e, 625 - l: slog.Default(), 626 - cfg: cfg, 627 - res: idresolver.DefaultResolver("https://plc.test"), 628 - jc: jc, 629 - rootCtx: context.Background(), 584 + // the knot's view: keptDid and a collaborator spindle never saw. Paginated, and it lists 585 + // the owner too - reconcile must not treat that as an explicit grant to remove. 586 + pages := [][]string{ 587 + {ownerDid.String(), keptDid.String()}, 588 + {newDid.String()}, 630 589 } 590 + var gotSubject string 591 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { 592 + if r.URL.Path != "/xrpc/"+tangled.RepoListCollaboratorsNSID { 593 + http.NotFound(w, r) 594 + return 595 + } 596 + gotSubject = r.URL.Query().Get("subject") 597 + page := 0 598 + if c := r.URL.Query().Get("cursor"); c != "" { 599 + page = 1 600 + } 601 + out := tangled.RepoListCollaborators_Output{} 602 + for _, s := range pages[page] { 603 + out.Items = append(out.Items, &tangled.RepoListCollaborators_ListItem{Subject: s}) 604 + } 605 + if page == 0 { 606 + next := "page2" 607 + out.Cursor = &next 608 + } 609 + json.NewEncoder(w).Encode(out) 610 + })) 611 + defer srv.Close() 631 612 613 + cfg := &config.Config{} 614 + cfg.Server.Dev = true // http, not https 615 + jc, jcerr := jetstream.NewJetstreamClient("", "", nil, nil, slog.Default(), nil, false, false) 616 + if jcerr != nil { 617 + t.Fatalf("NewJetstreamClient: %v", jcerr) 618 + } 632 619 tap := &Tap{ 633 - spindle: s, 620 + spindle: &Spindle{db: d, e: e, l: slog.Default(), cfg: cfg, jc: jc}, 634 621 logger: slog.Default(), 635 622 } 636 623 637 - ownerDid := syntax.DID("did:plc:repoowner") 638 - bobDid := syntax.DID("did:plc:bob") 639 - collabDid := syntax.DID("did:plc:collab") 640 - repoDid := syntax.DID("did:plc:testrepo123") 624 + knot := strings.TrimPrefix(srv.URL, "http://") 625 + tap.reconcileCollaborators(context.Background(), slog.Default(), knot, repoDid, ownerDid) 641 626 642 - err := d.AddRepo(db.Repo{ 643 - Knot: "knot.test", 644 - Owner: ownerDid, 645 - Rkey: "test-repo-rkey", 646 - RepoDid: repoDid, 647 - CreatedAt: time.Now().Format(time.RFC3339), 648 - }) 649 - if err != nil { 650 - t.Fatalf("AddRepo: %v", err) 627 + if gotSubject != repoDid.String() { 628 + t.Errorf("knot queried with subject %q, want %q", gotSubject, repoDid) 651 629 } 652 630 653 - err = e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()) 654 - if err != nil { 655 - t.Fatalf("AddRepo permissions: %v", err) 631 + for _, tc := range []struct { 632 + did syntax.DID 633 + want bool 634 + }{ 635 + {ownerDid, true}, // owner keeps access via role inheritance 636 + {keptDid, true}, // still on the knot 637 + {newDid, true}, // added from the knot's roster 638 + {staleDid, false}, // removed on the knot while spindle was down 639 + } { 640 + ok, err := e.IsRepoCollaborator(tc.did, repoDid) 641 + if err != nil { 642 + t.Fatalf("IsRepoCollaborator(%s): %v", tc.did, err) 643 + } 644 + if ok != tc.want { 645 + t.Errorf("IsRepoCollaborator(%s) = %v, want %v", tc.did, ok, tc.want) 646 + } 656 647 } 648 + } 649 + 650 + func TestReconcileCollaboratorsKeepsGrantsOnFetchFailure(t *testing.T) { 651 + d, e := newTestSpindleDB(t) 657 652 658 - err = d.AddRepoCollaborator(db.RepoCollaborator{ 659 - OwnerDid: ownerDid, 660 - Rkey: "collab-rkey", 661 - Subject: collabDid, 662 - RepoDid: repoDid, 663 - }) 664 - if err != nil { 653 + ownerDid := syntax.DID("did:plc:owner") 654 + repoDid := syntax.DID("did:plc:repo") 655 + collabDid := syntax.DID("did:plc:collab") 656 + 657 + if err := e.SetRepoOwner(ownerDid, repoDid); err != nil { 658 + t.Fatalf("SetRepoOwner: %v", err) 659 + } 660 + if err := e.AddRepoCollaborator(collabDid, repoDid); err != nil { 665 661 t.Fatalf("AddRepoCollaborator: %v", err) 666 662 } 667 663 668 - err = e.AddCollaborator(collabDid.String(), rbac.ThisServer, repoDid.String()) 669 - if err != nil { 670 - t.Fatalf("AddCollaborator policy: %v", err) 664 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { 665 + http.Error(w, "boom", http.StatusInternalServerError) 666 + })) 667 + defer srv.Close() 668 + 669 + cfg := &config.Config{} 670 + cfg.Server.Dev = true 671 + tap := &Tap{ 672 + spindle: &Spindle{db: d, e: e, l: slog.Default(), cfg: cfg}, 673 + logger: slog.Default(), 671 674 } 672 675 673 - // bob tries to delete alice's collaborator, must reject forged delete 674 - err = tap.processCollaborator(context.Background(), &tapc.RecordEventData{ 675 - Live: true, 676 - Did: bobDid, 677 - Rkey: "collab-rkey", 678 - Collection: syntax.NSID(tangled.RepoCollaboratorNSID), 679 - Action: tapc.RecordDeleteAction, 680 - }) 681 - if err != nil { 682 - t.Fatalf("processCollaborator delete returned error: %v", err) 683 - } 676 + tap.reconcileCollaborators(context.Background(), slog.Default(), 677 + strings.TrimPrefix(srv.URL, "http://"), repoDid, ownerDid) 684 678 685 - _, err = d.GetRepoCollaborator(ownerDid, "collab-rkey") 679 + ok, err := e.IsRepoCollaborator(collabDid, repoDid) 686 680 if err != nil { 687 - t.Fatalf("collaborator was deleted from DB: %v", err) 681 + t.Fatalf("IsRepoCollaborator: %v", err) 688 682 } 689 - 690 - ok, err := e.IsRepoCollaborator(collabDid.String(), rbac.ThisServer, repoDid.String()) 691 - if err != nil || !ok { 692 - t.Fatal("collaborator policy was removed from Casbin by forged delete") 683 + if !ok { 684 + t.Error("an unreachable knot wiped the collaborator roster") 693 685 } 694 686 }
+10 -11
spindle/xrpc/xrpc_test.go
··· 16 16 "github.com/bluesky-social/indigo/atproto/syntax" 17 17 "tangled.org/core/api/tangled" 18 18 "tangled.org/core/idresolver" 19 - "tangled.org/core/rbac" 19 + "tangled.org/core/rbac/v2" 20 20 "tangled.org/core/spindle/config" 21 21 "tangled.org/core/spindle/db" 22 22 "tangled.org/core/spindle/models" ··· 44 44 if err != nil { 45 45 t.Fatalf("rbac.NewEnforcer: %v", err) 46 46 } 47 - e.E.EnableAutoSave(true) 48 47 return d, e 49 48 } 50 49 ··· 56 55 pusherDid := syntax.DID("did:plc:pusher") 57 56 repoDid := syntax.DID("did:plc:testrepo123") 58 57 59 - err := d.AddRepo(db.Repo{ 58 + err := d.UpsertRepo(db.Repo{ 60 59 Knot: "knot.test", 61 60 Owner: repoOwnerDid, 62 61 Rkey: "test-repo-rkey", ··· 67 66 t.Fatalf("AddRepo: %v", err) 68 67 } 69 68 70 - err = e.AddRepo(repoOwnerDid.String(), rbac.ThisServer, repoDid.String()) 69 + err = e.SetRepoOwner(repoOwnerDid, repoDid) 71 70 if err != nil { 72 71 t.Fatalf("AddRepo permissions: %v", err) 73 72 } 74 - err = e.AddCollaborator(pusherDid.String(), rbac.ThisServer, repoDid.String()) 73 + err = e.AddRepoCollaborator(pusherDid, repoDid) 75 74 if err != nil { 76 75 t.Fatalf("AddCollaborator: %v", err) 77 76 } ··· 149 148 pusherDid := syntax.DID("did:plc:pusher") 150 149 repoDid := syntax.DID("did:plc:testrepo123") 151 150 152 - err := d.AddRepo(db.Repo{ 151 + err := d.UpsertRepo(db.Repo{ 153 152 Knot: "knot.test", 154 153 Owner: repoOwnerDid, 155 154 Rkey: "test-repo-rkey", ··· 160 159 t.Fatalf("AddRepo: %v", err) 161 160 } 162 161 163 - err = e.AddRepo(repoOwnerDid.String(), rbac.ThisServer, repoDid.String()) 162 + err = e.SetRepoOwner(repoOwnerDid, repoDid) 164 163 if err != nil { 165 164 t.Fatalf("AddRepo permissions: %v", err) 166 165 } 167 - err = e.AddCollaborator(pusherDid.String(), rbac.ThisServer, repoDid.String()) 166 + err = e.AddRepoCollaborator(pusherDid, repoDid) 168 167 if err != nil { 169 168 t.Fatalf("AddCollaborator: %v", err) 170 169 } ··· 260 259 pusherDid := syntax.DID("did:plc:pusher") 261 260 repoDid := syntax.DID("did:plc:testrepo123") 262 261 263 - err := d.AddRepo(db.Repo{ 262 + err := d.UpsertRepo(db.Repo{ 264 263 Knot: "knot.test", 265 264 Owner: repoOwnerDid, 266 265 Rkey: "test-repo-rkey", ··· 271 270 t.Fatalf("AddRepo: %v", err) 272 271 } 273 272 274 - err = e.AddRepo(repoOwnerDid.String(), rbac.ThisServer, repoDid.String()) 273 + err = e.SetRepoOwner(repoOwnerDid, repoDid) 275 274 if err != nil { 276 275 t.Fatalf("AddRepo permissions: %v", err) 277 276 } 278 - err = e.AddCollaborator(pusherDid.String(), rbac.ThisServer, repoDid.String()) 277 + err = e.AddRepoCollaborator(pusherDid, repoDid) 279 278 if err != nil { 280 279 t.Fatalf("AddCollaborator: %v", err) 281 280 }
+3
tapc/tap.go
··· 41 41 } 42 42 43 43 func (c *Client) AddRepos(ctx context.Context, dids []syntax.DID) error { 44 + if len(dids) == 0 { 45 + return nil 46 + } 44 47 body, err := json.Marshal(map[string][]syntax.DID{"dids": dids}) 45 48 if err != nil { 46 49 return err