Monorepo for Tangled
0

Configure Feed

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

core / appview / migration / backfill_entity_state_test.go
13 kB 412 lines
1package migration 2 3import ( 4 "context" 5 "encoding/json" 6 "io" 7 "log/slog" 8 "net/http" 9 "net/http/httptest" 10 "strings" 11 "sync" 12 "testing" 13 "time" 14 15 "github.com/bluesky-social/indigo/atproto/atclient" 16 "github.com/bluesky-social/indigo/atproto/syntax" 17 "github.com/samber/lo" 18 cbg "github.com/whyrusleeping/cbor-gen" 19 "tangled.org/core/api/tangled" 20 "tangled.org/core/appview/db" 21 "tangled.org/core/appview/models" 22 "tangled.org/core/orm" 23) 24 25func TestBackfillRecordRoundTrip(t *testing.T) { 26 createdAt := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC) 27 for _, tc := range []struct { 28 name string 29 subject syntax.ATURI 30 value models.StateValue 31 wantCollection string 32 fromRecord func(cbg.CBORMarshaler) (models.StateRecord, error) 33 }{ 34 { 35 name: "issue closed", 36 subject: "at://did:plc:boltless/sh.tangled.repo.issue/issue1", 37 value: models.StateClosed, 38 wantCollection: tangled.RepoIssueStateNSID, 39 fromRecord: func(r cbg.CBORMarshaler) (models.StateRecord, error) { 40 return models.IssueStateFromRecord("did:plc:akshay", "s1", *r.(*tangled.RepoIssueState)) 41 }, 42 }, 43 { 44 name: "pull merged", 45 subject: "at://did:plc:boltless/sh.tangled.repo.pull/pull1", 46 value: models.StateMerged, 47 wantCollection: tangled.RepoPullStatusNSID, 48 fromRecord: func(r cbg.CBORMarshaler) (models.StateRecord, error) { 49 return models.PullStatusFromRecord("did:plc:akshay", "s1", *r.(*tangled.RepoPullStatus)) 50 }, 51 }, 52 } { 53 t.Run(tc.name, func(t *testing.T) { 54 collection, record, err := backfillRecord(db.BackfillSubject{ 55 Subject: tc.subject, 56 Value: tc.value, 57 CreatedAt: createdAt.Format(time.RFC3339), 58 }, createdAt) 59 if err != nil { 60 t.Fatalf("backfillRecord: %v", err) 61 } 62 if collection != tc.wantCollection { 63 t.Fatalf("collection = %s, want %s", collection, tc.wantCollection) 64 } 65 rec, err := tc.fromRecord(record) 66 if err != nil { 67 t.Fatalf("fromRecord: %v", err) 68 } 69 if rec.Value != tc.value { 70 t.Fatalf("round-trip value = %q, want %q", rec.Value, tc.value) 71 } 72 if rec.Subject != tc.subject { 73 t.Fatalf("round-trip subject = %s, want %s", rec.Subject, tc.subject) 74 } 75 if rec.SortMicros != createdAt.UnixMicro() { 76 t.Fatalf("round-trip sort micros = %d, want %d", rec.SortMicros, createdAt.UnixMicro()) 77 } 78 }) 79 } 80} 81 82func TestBackfillRecordRejects(t *testing.T) { 83 for _, tc := range []struct { 84 name string 85 subject syntax.ATURI 86 value models.StateValue 87 }{ 88 {"issue subject with merged value", "at://did:plc:boltless/sh.tangled.repo.issue/issue1", models.StateMerged}, 89 {"non issue or pull subject", "at://did:plc:boltless/sh.tangled.feed.star/x", models.StateClosed}, 90 } { 91 t.Run(tc.name, func(t *testing.T) { 92 if _, _, err := backfillRecord(db.BackfillSubject{ 93 Subject: tc.subject, 94 Value: tc.value, 95 }, time.Unix(0, 0)); err == nil { 96 t.Fatalf("backfillRecord must reject %s", tc.name) 97 } 98 }) 99 } 100} 101 102type capturedPut struct { 103 Collection string `json:"collection"` 104 Repo string `json:"repo"` 105 Rkey string `json:"rkey"` 106 Record json.RawMessage `json:"record"` 107} 108 109type putRecordStore struct { 110 mu sync.Mutex 111 puts []capturedPut 112} 113 114func (s *putRecordStore) add(p capturedPut) { 115 s.mu.Lock() 116 defer s.mu.Unlock() 117 s.puts = append(s.puts, p) 118} 119 120func (s *putRecordStore) count() int { 121 s.mu.Lock() 122 defer s.mu.Unlock() 123 return len(s.puts) 124} 125 126func (s *putRecordStore) byRkey() map[string]capturedPut { 127 s.mu.Lock() 128 defer s.mu.Unlock() 129 return lo.SliceToMap(s.puts, func(p capturedPut) (string, capturedPut) { 130 return p.Rkey, p 131 }) 132} 133 134func discardLogger() *slog.Logger { 135 return slog.New(slog.NewTextHandler(io.Discard, nil)) 136} 137 138func newPutRecordServer(t *testing.T, failRkeys map[string]bool) (*putRecordStore, *atclient.APIClient) { 139 t.Helper() 140 store := &putRecordStore{} 141 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { 142 if !strings.HasSuffix(r.URL.Path, "/com.atproto.repo.putRecord") { 143 t.Errorf("unexpected request path: %s", r.URL.Path) 144 http.Error(w, "not found", http.StatusNotFound) 145 return 146 } 147 var in capturedPut 148 if err := json.NewDecoder(r.Body).Decode(&in); err != nil { 149 http.Error(w, err.Error(), http.StatusBadRequest) 150 return 151 } 152 if failRkeys[in.Rkey] { 153 http.Error(w, `{"error":"InternalServerError","message":"boom"}`, http.StatusInternalServerError) 154 return 155 } 156 store.add(in) 157 _ = json.NewEncoder(w).Encode(map[string]string{ 158 "uri": "at://" + in.Repo + "/" + in.Collection + "/" + in.Rkey, 159 "cid": "bafyreigxt-test", 160 }) 161 })) 162 t.Cleanup(srv.Close) 163 return store, &atclient.APIClient{Host: srv.URL, Client: srv.Client()} 164} 165 166func seedRepoForOwner(t *testing.T, d *db.DB, owner, repoDid string) { 167 t.Helper() 168 tx, err := d.Begin() 169 if err != nil { 170 t.Fatalf("Begin: %v", err) 171 } 172 if err := db.AddRepo(tx, &models.Repo{ 173 Did: owner, 174 Name: "anemone", 175 Knot: "knot.example", 176 Rkey: "anemone", 177 RepoDid: repoDid, 178 }); err != nil { 179 t.Fatalf("AddRepo: %v", err) 180 } 181 if err := tx.Commit(); err != nil { 182 t.Fatalf("Commit: %v", err) 183 } 184} 185 186func seedIssueRow(t *testing.T, d *db.DB, repoDid, author, rkey string) *models.Issue { 187 t.Helper() 188 tx, err := d.Begin() 189 if err != nil { 190 t.Fatalf("Begin: %v", err) 191 } 192 issue := &models.Issue{ 193 Did: author, 194 Rkey: rkey, 195 RepoDid: syntax.DID(repoDid), 196 Title: "title", 197 Body: "body", 198 Open: true, 199 } 200 if err := db.PutIssue(tx, issue); err != nil { 201 t.Fatalf("PutIssue: %v", err) 202 } 203 if err := tx.Commit(); err != nil { 204 t.Fatalf("Commit: %v", err) 205 } 206 return issue 207} 208 209func seedPullRow(t *testing.T, d *db.DB, repoDid, author, rkey string) *models.Pull { 210 t.Helper() 211 tx, err := d.Begin() 212 if err != nil { 213 t.Fatalf("Begin: %v", err) 214 } 215 pull := &models.Pull{ 216 RepoDid: syntax.DID(repoDid), 217 OwnerDid: author, 218 Rkey: rkey, 219 Title: "title", 220 Body: "body", 221 TargetBranch: "main", 222 State: models.PullOpen, 223 } 224 if err := db.PutPull(tx, pull); err != nil { 225 t.Fatalf("PutPull: %v", err) 226 } 227 if err := tx.Commit(); err != nil { 228 t.Fatalf("Commit: %v", err) 229 } 230 return pull 231} 232 233func parseIssueStatePut(did, rkey string, raw json.RawMessage) (models.StateRecord, error) { 234 var rec tangled.RepoIssueState 235 if err := json.Unmarshal(raw, &rec); err != nil { 236 return models.StateRecord{}, err 237 } 238 return models.IssueStateFromRecord(did, rkey, rec) 239} 240 241func parsePullStatusPut(did, rkey string, raw json.RawMessage) (models.StateRecord, error) { 242 var rec tangled.RepoPullStatus 243 if err := json.Unmarshal(raw, &rec); err != nil { 244 return models.StateRecord{}, err 245 } 246 return models.PullStatusFromRecord(did, rkey, rec) 247} 248 249func assertStatePut( 250 t *testing.T, 251 put capturedPut, 252 owner syntax.DID, 253 subject syntax.ATURI, 254 wantCollection string, 255 want models.StateValue, 256 parse func(did, rkey string, raw json.RawMessage) (models.StateRecord, error), 257) { 258 t.Helper() 259 if put.Collection != wantCollection { 260 t.Fatalf("collection = %q, want %q", put.Collection, wantCollection) 261 } 262 if put.Repo != owner.String() { 263 t.Fatalf("repo = %q, want owner %q", put.Repo, owner) 264 } 265 if put.Rkey != subject.RecordKey().String() { 266 t.Fatalf("rkey = %q, want subject rkey %q", put.Rkey, subject.RecordKey()) 267 } 268 sr, err := parse(owner.String(), put.Rkey, put.Record) 269 if err != nil { 270 t.Fatalf("parse state record: %v", err) 271 } 272 if sr.Subject != subject { 273 t.Fatalf("subject = %s, want %s", sr.Subject, subject) 274 } 275 if sr.Value != want { 276 t.Fatalf("value = %q, want %q", sr.Value, want) 277 } 278} 279 280func TestBackfillEntityStateWritesRecords(t *testing.T) { 281 d := newTestDB(t) 282 owner := syntax.DID("did:plc:akshay") 283 const repoDid = "did:plc:anemone" 284 seedRepoForOwner(t, d, string(owner), repoDid) 285 286 closedIssue := seedIssueRow(t, d, repoDid, "did:plc:boltless", "issueClosed") 287 seedIssueRow(t, d, repoDid, "did:plc:boltless", "issueOpen") 288 mergedPull := seedPullRow(t, d, repoDid, "did:plc:boltless", "pullMerged") 289 closedPull := seedPullRow(t, d, repoDid, "did:plc:boltless", "pullClosed") 290 seedPullRow(t, d, repoDid, "did:plc:boltless", "pullOpen") 291 292 if err := db.CloseIssues(d, orm.FilterEq("at_uri", closedIssue.AtUri())); err != nil { 293 t.Fatalf("CloseIssues: %v", err) 294 } 295 if err := db.MergePulls(d, orm.FilterEq("at_uri", mergedPull.AtUri())); err != nil { 296 t.Fatalf("MergePulls: %v", err) 297 } 298 if err := db.ClosePulls(d, orm.FilterEq("at_uri", closedPull.AtUri())); err != nil { 299 t.Fatalf("ClosePulls: %v", err) 300 } 301 302 store, client := newPutRecordServer(t, nil) 303 client.AccountDID = &owner 304 m := &Migration{db: d, logger: discardLogger()} 305 306 if err := m.backfillEntityState(context.Background(), client, owner, ""); err != nil { 307 t.Fatalf("backfillEntityState: %v", err) 308 } 309 310 got := store.byRkey() 311 if len(got) != 3 { 312 t.Fatalf("wrote %d records, want 3: %+v", len(got), got) 313 } 314 assertStatePut(t, got["issueClosed"], owner, closedIssue.AtUri(), tangled.RepoIssueStateNSID, models.StateClosed, parseIssueStatePut) 315 assertStatePut(t, got["pullMerged"], owner, mergedPull.AtUri(), tangled.RepoPullStatusNSID, models.StateMerged, parsePullStatusPut) 316 assertStatePut(t, got["pullClosed"], owner, closedPull.AtUri(), tangled.RepoPullStatusNSID, models.StateClosed, parsePullStatusPut) 317} 318 319func TestBackfillEntityStateSkipsCleanOwner(t *testing.T) { 320 d := newTestDB(t) 321 owner := syntax.DID("did:plc:akshay") 322 const repoDid = "did:plc:anemone" 323 seedRepoForOwner(t, d, string(owner), repoDid) 324 seedIssueRow(t, d, repoDid, "did:plc:boltless", "issueOpen") 325 326 store, client := newPutRecordServer(t, nil) 327 m := &Migration{db: d, logger: discardLogger()} 328 329 if err := m.backfillEntityState(context.Background(), client, owner, ""); err != nil { 330 t.Fatalf("backfillEntityState: %v", err) 331 } 332 if store.count() != 0 { 333 t.Fatalf("wrote %d records for an owner with no column-only closed state, want 0", store.count()) 334 } 335} 336 337func TestBackfillEntityStateCreatedAt(t *testing.T) { 338 created := time.Date(2024, 3, 4, 5, 6, 7, 0, time.UTC) 339 for _, tc := range []struct { 340 name string 341 stored string 342 want time.Time 343 }{ 344 {"uses subject created", created.Format(time.RFC3339), created}, 345 {"falls back to epoch on unparseable", "not-a-timestamp", time.Unix(0, 0).UTC()}, 346 } { 347 t.Run(tc.name, func(t *testing.T) { 348 d := newTestDB(t) 349 owner := syntax.DID("did:plc:akshay") 350 const repoDid = "did:plc:anemone" 351 seedRepoForOwner(t, d, string(owner), repoDid) 352 issue := seedIssueRow(t, d, repoDid, "did:plc:boltless", "issueClosed") 353 if err := db.CloseIssues(d, orm.FilterEq("at_uri", issue.AtUri())); err != nil { 354 t.Fatalf("CloseIssues: %v", err) 355 } 356 if _, err := d.Exec(`update issues set created = ? where at_uri = ?`, tc.stored, issue.AtUri()); err != nil { 357 t.Fatalf("update created: %v", err) 358 } 359 360 store, client := newPutRecordServer(t, nil) 361 client.AccountDID = &owner 362 m := &Migration{db: d, logger: discardLogger()} 363 364 if err := m.backfillEntityState(context.Background(), client, owner, ""); err != nil { 365 t.Fatalf("backfillEntityState: %v", err) 366 } 367 368 var rec tangled.RepoIssueState 369 if err := json.Unmarshal(store.byRkey()["issueClosed"].Record, &rec); err != nil { 370 t.Fatalf("unmarshal: %v", err) 371 } 372 want, err := models.AsIssueStateRecord(issue.AtUri(), models.StateClosed, tc.want) 373 if err != nil { 374 t.Fatalf("AsIssueStateRecord: %v", err) 375 } 376 if rec.CreatedAt != want.CreatedAt { 377 t.Fatalf("createdAt = %q, want %q", rec.CreatedAt, want.CreatedAt) 378 } 379 }) 380 } 381} 382 383func TestBackfillEntityStatePartialFailureReturnsError(t *testing.T) { 384 d := newTestDB(t) 385 owner := syntax.DID("did:plc:akshay") 386 const repoDid = "did:plc:anemone" 387 seedRepoForOwner(t, d, string(owner), repoDid) 388 good := seedIssueRow(t, d, repoDid, "did:plc:boltless", "issueGood") 389 bad := seedIssueRow(t, d, repoDid, "did:plc:boltless", "issueBad") 390 if err := db.CloseIssues(d, orm.FilterEq("at_uri", good.AtUri())); err != nil { 391 t.Fatalf("CloseIssues good: %v", err) 392 } 393 if err := db.CloseIssues(d, orm.FilterEq("at_uri", bad.AtUri())); err != nil { 394 t.Fatalf("CloseIssues bad: %v", err) 395 } 396 397 store, client := newPutRecordServer(t, map[string]bool{"issueBad": true}) 398 client.AccountDID = &owner 399 m := &Migration{db: d, logger: discardLogger()} 400 401 if err := m.backfillEntityState(context.Background(), client, owner, ""); err == nil { 402 t.Fatal("backfillEntityState must return an error when a subject write fails") 403 } 404 405 got := store.byRkey() 406 if _, ok := got["issueGood"]; !ok { 407 t.Fatal("the succeeding subject must still be written when another fails") 408 } 409 if _, ok := got["issueBad"]; ok { 410 t.Fatal("the failing subject must not be recorded as written") 411 } 412}