Monorepo for Tangled tangled.org
4

Configure Feed

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

core / appview / ingester_state_test.go
7.7 kB 246 lines
1package appview 2 3import ( 4 "context" 5 "database/sql" 6 "encoding/json" 7 "errors" 8 "log/slog" 9 "path/filepath" 10 "testing" 11 12 "github.com/bluesky-social/indigo/atproto/syntax" 13 jmodels "github.com/bluesky-social/jetstream/pkg/models" 14 "tangled.org/core/api/tangled" 15 "tangled.org/core/appview/db" 16 "tangled.org/core/appview/knotacl" 17 "tangled.org/core/appview/models" 18 "tangled.org/core/orm" 19) 20 21func newStateIngester(t *testing.T) *Ingester { 22 t.Helper() 23 path := filepath.Join(t.TempDir(), "test.db") 24 d, err := db.Make(t.Context(), path) 25 if err != nil { 26 t.Fatalf("db.Make: %v", err) 27 } 28 t.Cleanup(func() { d.Close() }) 29 return &Ingester{ 30 Ctx: t.Context(), 31 Db: d, 32 Logger: slog.New(slog.DiscardHandler), 33 } 34} 35 36type stubAcl struct { 37 allow bool 38 err error 39} 40 41func (s stubAcl) HasRepoPermissionErr(ctx context.Context, repo *models.Repo, userDid, perm string) (bool, error) { 42 return s.allow, s.err 43} 44 45func seedTx(t *testing.T, d *db.DB, fn func(tx *sql.Tx) error) { 46 t.Helper() 47 tx, err := d.Begin() 48 if err != nil { 49 t.Fatalf("Begin: %v", err) 50 } 51 defer tx.Rollback() 52 if err := fn(tx); err != nil { 53 t.Fatalf("seed: %v", err) 54 } 55 if err := tx.Commit(); err != nil { 56 t.Fatalf("Commit: %v", err) 57 } 58} 59 60func seedRepo(t *testing.T, d *db.DB, ownerDid, repoDid string) { 61 t.Helper() 62 seedTx(t, d, func(tx *sql.Tx) error { 63 return db.AddRepo(tx, &models.Repo{ 64 Did: ownerDid, 65 Name: "anemone", 66 Knot: "knot.example", 67 Rkey: "anemone", 68 RepoDid: repoDid, 69 }) 70 }) 71} 72 73func seedIssue(t *testing.T, d *db.DB, ownerDid, repoDid, issueRkey string) syntax.ATURI { 74 t.Helper() 75 issue := &models.Issue{ 76 Did: ownerDid, 77 Rkey: issueRkey, 78 RepoDid: syntax.DID(repoDid), 79 Title: "title", 80 Body: "body", 81 Open: true, 82 } 83 seedTx(t, d, func(tx *sql.Tx) error { 84 return db.PutIssue(tx, issue) 85 }) 86 return issue.AtUri() 87} 88 89func seedRepoAndIssue(t *testing.T, d *db.DB, ownerDid, repoDid, issueRkey string) syntax.ATURI { 90 t.Helper() 91 seedRepo(t, d, ownerDid, repoDid) 92 return seedIssue(t, d, ownerDid, repoDid, issueRkey) 93} 94 95func issueStateEvent(t *testing.T, op, did, rkey, subject, state, createdAt string) *jmodels.Event { 96 t.Helper() 97 raw, err := json.Marshal(tangled.RepoIssueState{ 98 Issue: subject, 99 State: state, 100 CreatedAt: createdAt, 101 }) 102 if err != nil { 103 t.Fatalf("marshal: %v", err) 104 } 105 return &jmodels.Event{ 106 Did: did, 107 Kind: jmodels.EventKindCommit, 108 Commit: &jmodels.Commit{ 109 Operation: op, 110 Collection: tangled.RepoIssueStateNSID, 111 RKey: rkey, 112 Record: raw, 113 }, 114 } 115} 116 117func ingestedIssueOpen(t *testing.T, d *db.DB, subject syntax.ATURI) bool { 118 t.Helper() 119 issues, err := db.GetIssues(d, orm.FilterEq("at_uri", subject)) 120 if err != nil || len(issues) != 1 { 121 t.Fatalf("GetIssues: %v len %d", err, len(issues)) 122 } 123 return issues[0].Open 124} 125 126func TestIngestState_Authorization(t *testing.T) { 127 owner := "did:plc:boltless" 128 cases := []struct { 129 name string 130 author string 131 acl RepoPermissionChecker 132 wantOpen bool 133 }{ 134 {"author closes own issue", owner, stubAcl{err: errors.New("acl must not be consulted for the author")}, false}, 135 {"collaborator with push closes", "did:plc:squid", stubAcl{allow: true}, false}, 136 {"stranger without push rejected", "did:plc:squid", stubAcl{allow: false}, true}, 137 {"unreachable knot fails open", "did:plc:squid", stubAcl{allow: false, err: knotacl.ErrKnotUnreachable}, false}, 138 } 139 for _, tc := range cases { 140 t.Run(tc.name, func(t *testing.T) { 141 ing := newStateIngester(t) 142 ing.Acl = tc.acl 143 at := seedRepoAndIssue(t, ing.Db, owner, "did:plc:anemone", "issue1") 144 145 ev := issueStateEvent(t, jmodels.CommitOperationCreate, tc.author, "s1", string(at), tangled.RepoIssueStateClosed, "2026-06-01T00:00:00Z") 146 if err := ing.ingestState(context.Background(), ev, ing.Logger, issueStateSpec); err != nil { 147 t.Fatalf("ingestState: %v", err) 148 } 149 150 if got := ingestedIssueOpen(t, ing.Db, at); got != tc.wantOpen { 151 t.Fatalf("issue open = %v, want %v", got, tc.wantOpen) 152 } 153 if pending, _ := db.PendingStateRecordsForSubject(ing.Db, at); len(pending) != 0 { 154 t.Fatalf("a record whose subject exists must resolve, not park: pending=%d", len(pending)) 155 } 156 }) 157 } 158} 159 160func TestIngestState_ParkThenReconcileDrains(t *testing.T) { 161 ing := newStateIngester(t) 162 owner := "did:plc:boltless" 163 subject := syntax.ATURI("at://" + owner + "/" + tangled.RepoIssueNSID + "/issue1") 164 165 ev := issueStateEvent(t, jmodels.CommitOperationCreate, owner, "s1", string(subject), tangled.RepoIssueStateClosed, "2026-06-01T00:00:00Z") 166 if err := ing.ingestState(ing.Ctx, ev, ing.Logger, issueStateSpec); err != nil { 167 t.Fatalf("park: %v", err) 168 } 169 if pending, _ := db.PendingStateRecordsForSubject(ing.Db, subject); len(pending) != 1 { 170 t.Fatalf("a state record whose subject is missing must be parked, pending=%d", len(pending)) 171 } 172 173 at := seedRepoAndIssue(t, ing.Db, owner, "did:plc:anemone", "issue1") 174 if at != subject { 175 t.Fatalf("subject mismatch: %s vs %s", at, subject) 176 } 177 ing.ReconcilePendingState() 178 179 if ingestedIssueOpen(t, ing.Db, at) { 180 t.Fatal("the reconciler must drain a parked record once its subject exists") 181 } 182 if pending, _ := db.PendingStateRecordsForSubject(ing.Db, at); len(pending) != 0 { 183 t.Fatalf("a drained record must be unparked, pending=%d", len(pending)) 184 } 185} 186 187func TestIngestState_DeleteUnparksBeforeSubjectArrives(t *testing.T) { 188 ing := newStateIngester(t) 189 ctx := context.Background() 190 owner := "did:plc:boltless" 191 subject := syntax.ATURI("at://" + owner + "/" + tangled.RepoIssueNSID + "/issue1") 192 193 create := issueStateEvent(t, jmodels.CommitOperationCreate, owner, "s1", string(subject), tangled.RepoIssueStateClosed, "2026-06-01T00:00:00Z") 194 if err := ing.ingestState(ctx, create, ing.Logger, issueStateSpec); err != nil { 195 t.Fatalf("park: %v", err) 196 } 197 198 del := &jmodels.Event{ 199 Did: owner, 200 Kind: jmodels.EventKindCommit, 201 Commit: &jmodels.Commit{ 202 Operation: jmodels.CommitOperationDelete, 203 Collection: tangled.RepoIssueStateNSID, 204 RKey: "s1", 205 }, 206 } 207 if err := ing.ingestState(ctx, del, ing.Logger, issueStateSpec); err != nil { 208 t.Fatalf("delete: %v", err) 209 } 210 211 if pending, _ := db.PendingStateRecordsForSubject(ing.Db, subject); len(pending) != 0 { 212 t.Fatalf("deleting a parked record must unpark it, pending=%d", len(pending)) 213 } 214 215 at := seedRepoAndIssue(t, ing.Db, owner, "did:plc:anemone", "issue1") 216 ing.drainPendingState(ctx, at, ing.Logger) 217 if !ingestedIssueOpen(t, ing.Db, at) { 218 t.Fatal("a deleted parked record must not apply after the subject arrives") 219 } 220} 221 222func TestIngestState_ParkedUnauthorizedDroppedOnDrain(t *testing.T) { 223 ing := newStateIngester(t) 224 ing.Acl = stubAcl{allow: false} 225 ctx := context.Background() 226 owner := "did:plc:boltless" 227 subject := syntax.ATURI("at://" + owner + "/" + tangled.RepoIssueNSID + "/issue1") 228 229 ev := issueStateEvent(t, jmodels.CommitOperationCreate, "did:plc:squid", "s1", string(subject), tangled.RepoIssueStateClosed, "2026-06-01T00:00:00Z") 230 if err := ing.ingestState(ctx, ev, ing.Logger, issueStateSpec); err != nil { 231 t.Fatalf("park: %v", err) 232 } 233 if pending, _ := db.PendingStateRecordsForSubject(ing.Db, subject); len(pending) != 1 { 234 t.Fatalf("a record whose subject is missing parks before any auth check: pending=%d", len(pending)) 235 } 236 237 at := seedRepoAndIssue(t, ing.Db, owner, "did:plc:anemone", "issue1") 238 ing.drainPendingState(ctx, at, ing.Logger) 239 240 if !ingestedIssueOpen(t, ing.Db, at) { 241 t.Fatal("a parked record that fails authorization on drain must not apply") 242 } 243 if pending, _ := db.PendingStateRecordsForSubject(ing.Db, at); len(pending) != 0 { 244 t.Fatalf("a parked record rejected on drain must be unparked, not left to re-accumulate: pending=%d", len(pending)) 245 } 246}