forked from
tangled.org/core
Monorepo for Tangled
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}