This repository has no description
0

Configure Feed

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

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