From e1efbcc1a0247fa829d16531336495cf19996ee9 Mon Sep 17 00:00:00 2001 From: Lewis Date: Fri, 03 Jul 2026 07:24:30 +0000 Subject: [PATCH] appview/db: entity-state backfill queries Lewis: May this revision serve well! --- appview/db/entity_state_backfill.go | 85 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ appview/db/entity_state_backfill_test.go | 181 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 2 file(s) changed, 266 insertion(s)(+), 0 deletion(s)(-) diff --git a/appview/db/entity_state_backfill.go b/appview/db/entity_state_backfill.go new file mode 100644 --- /dev/null +++ b/appview/db/entity_state_backfill.go @@ -0,0 +1,85 @@ +package db + +import ( + "context" + + "github.com/bluesky-social/indigo/atproto/syntax" + + "tangled.org/core/appview/models" +) + +const EntityStateBackfillName = "backfill-entity-state" + +type BackfillSubject struct { + Subject syntax.ATURI + Value models.StateValue + CreatedAt string +} + +func EnqueueEntityStateBackfill(ctx context.Context, e Execer) (int64, error) { + res, err := e.ExecContext(ctx, ` + insert into pds_migration (name, did, collection, rkey) + select distinct ?, r.did, '', '' + from repos r + where r.did like 'did:%' + and ( + exists ( + select 1 from issues i + where i.repo_did = r.repo_did and i.open = 0 and i.deleted is null and i.rkey != '' + and not exists (select 1 from issue_states s where s.subject = i.at_uri) + ) + or exists ( + select 1 from pulls p + where p.repo_did = r.repo_did and p.state in (0, 2) and p.rkey != '' + and not exists (select 1 from pull_states s where s.subject = p.at_uri) + ) + ) + on conflict(name, did, collection, rkey) do update set + status = 'pending', + retry_count = 0, + retry_after = 0, + error_msg = null + where pds_migration.status in ('done', 'failed') + `, EntityStateBackfillName) + if err != nil { + return 0, err + } + return res.RowsAffected() +} + +func ColumnOnlyClosedSubjectsForOwner(ctx context.Context, e Execer, owner syntax.DID) ([]BackfillSubject, error) { + rows, err := e.QueryContext(ctx, ` + select i.at_uri, ?, i.created + from issues i join repos r on i.repo_did = r.repo_did + where r.did = ? and i.open = 0 and i.deleted is null and i.rkey != '' + and not exists (select 1 from issue_states s where s.subject = i.at_uri) + union all + select p.at_uri, case p.state when 2 then ? else ? end, p.created + from pulls p join repos r on p.repo_did = r.repo_did + where r.did = ? and p.state in (0, 2) and p.rkey != '' + and not exists (select 1 from pull_states s where s.subject = p.at_uri) + `, + string(models.StateClosed), + owner, + string(models.StateMerged), string(models.StateClosed), + owner, + ) + if err != nil { + return nil, err + } + defer rows.Close() + + var subjects []BackfillSubject + for rows.Next() { + var subject, value, created string + if err := rows.Scan(&subject, &value, &created); err != nil { + return nil, err + } + subjects = append(subjects, BackfillSubject{ + Subject: syntax.ATURI(subject), + Value: models.StateValue(value), + CreatedAt: created, + }) + } + return subjects, rows.Err() +} diff --git a/appview/db/entity_state_backfill_test.go b/appview/db/entity_state_backfill_test.go new file mode 100644 --- /dev/null +++ b/appview/db/entity_state_backfill_test.go @@ -0,0 +1,181 @@ +package db + +import ( + "context" + "testing" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/samber/lo" + "tangled.org/core/appview/models" + "tangled.org/core/orm" +) + +func TestColumnOnlyClosedSubjectsForOwner(t *testing.T) { + d := newTestDB(t) + owner := syntax.DID("did:plc:akshay") + repo := seedRepo(t, d, string(owner), "knot.example", "anemone", "anemone", "did:plc:anemone") + + closedIssue := seedIssue(t, d, repo, "did:plc:boltless", "issueClosed") + seedIssue(t, d, repo, "did:plc:boltless", "issueOpen") + recordedIssue := seedIssue(t, d, repo, "did:plc:boltless", "issueRecorded") + + if err := CloseIssues(d, orm.FilterEq("at_uri", closedIssue.AtUri())); err != nil { + t.Fatalf("CloseIssues closed: %v", err) + } + if err := CloseIssues(d, orm.FilterEq("at_uri", recordedIssue.AtUri())); err != nil { + t.Fatalf("CloseIssues recorded: %v", err) + } + putIssueStateRec(t, d, issueRec("did:plc:boltless", "srec", recordedIssue.AtUri(), models.StateClosed, 100)) + + emptyRkey := seedIssue(t, d, repo, "did:plc:boltless", "") + if err := CloseIssues(d, orm.FilterEq("at_uri", emptyRkey.AtUri())); err != nil { + t.Fatalf("CloseIssues empty rkey: %v", err) + } + + mergedPull := seedPull(t, d, repo, "did:plc:boltless", "pullMerged") + closedPull := seedPull(t, d, repo, "did:plc:boltless", "pullClosed") + seedPull(t, d, repo, "did:plc:boltless", "pullOpen") + + if err := MergePulls(d, orm.FilterEq("at_uri", mergedPull.AtUri())); err != nil { + t.Fatalf("MergePulls: %v", err) + } + if err := ClosePulls(d, orm.FilterEq("at_uri", closedPull.AtUri())); err != nil { + t.Fatalf("ClosePulls: %v", err) + } + + subjects, err := ColumnOnlyClosedSubjectsForOwner(context.Background(), d, owner) + if err != nil { + t.Fatalf("ColumnOnlyClosedSubjectsForOwner: %v", err) + } + + got := lo.SliceToMap(subjects, func(s BackfillSubject) (syntax.ATURI, models.StateValue) { + return s.Subject, s.Value + }) + + if len(got) != 3 { + t.Fatalf("got %d subjects, want 3: %+v", len(got), got) + } + if got[closedIssue.AtUri()] != models.StateClosed { + t.Fatalf("closed issue: got %q want closed", got[closedIssue.AtUri()]) + } + if got[mergedPull.AtUri()] != models.StateMerged { + t.Fatalf("merged pull: got %q want merged", got[mergedPull.AtUri()]) + } + if got[closedPull.AtUri()] != models.StateClosed { + t.Fatalf("closed pull: got %q want closed", got[closedPull.AtUri()]) + } + if _, ok := got[emptyRkey.AtUri()]; ok { + t.Fatal("closed issue with empty rkey must be excluded") + } +} + +func mustEnqueueBackfill(t *testing.T, d *DB) { + t.Helper() + if _, err := EnqueueEntityStateBackfill(context.Background(), d); err != nil { + t.Fatalf("EnqueueEntityStateBackfill: %v", err) + } +} + +func markBackfillDone(t *testing.T, d *DB) { + t.Helper() + if _, err := d.Exec( + `update pds_migration set status = 'done' where name = ?`, EntityStateBackfillName, + ); err != nil { + t.Fatalf("mark done: %v", err) + } +} + +func TestEnqueueEntityStateBackfill(t *testing.T) { + owner := syntax.DID("did:plc:akshay") + for _, tc := range []struct { + name string + arrange func(t *testing.T, d *DB, issue *models.Issue) + wantRows int64 + wantNoRow bool + wantStatus models.PDSMigrationStatus + }{ + { + name: "enqueues owner with column-only closed work", + arrange: func(t *testing.T, d *DB, issue *models.Issue) {}, + wantRows: 1, + wantStatus: models.PDSMigrationStatusPending, + }, + { + name: "idempotent while still pending", + arrange: func(t *testing.T, d *DB, issue *models.Issue) { mustEnqueueBackfill(t, d) }, + wantRows: 0, + wantStatus: models.PDSMigrationStatusPending, + }, + { + name: "re-arms completed owner with fresh work", + arrange: func(t *testing.T, d *DB, issue *models.Issue) { + mustEnqueueBackfill(t, d) + markBackfillDone(t, d) + }, + wantRows: 1, + wantStatus: models.PDSMigrationStatusPending, + }, + { + name: "leaves settled owner done", + arrange: func(t *testing.T, d *DB, issue *models.Issue) { + mustEnqueueBackfill(t, d) + putIssueStateRec(t, d, issueRec(string(owner), "srec", issue.AtUri(), models.StateClosed, 100)) + markBackfillDone(t, d) + }, + wantRows: 0, + wantStatus: models.PDSMigrationStatusDone, + }, + { + name: "skips owner whose closed work is already recorded", + arrange: func(t *testing.T, d *DB, issue *models.Issue) { + putIssueStateRec(t, d, issueRec(string(owner), "srec", issue.AtUri(), models.StateClosed, 100)) + }, + wantRows: 0, + wantNoRow: true, + }, + } { + t.Run(tc.name, func(t *testing.T) { + d := newTestDB(t) + repo := seedRepo(t, d, string(owner), "knot.example", "anemone", "anemone", "did:plc:anemone") + issue := seedIssue(t, d, repo, "did:plc:boltless", "issue1") + if err := CloseIssues(d, orm.FilterEq("at_uri", issue.AtUri())); err != nil { + t.Fatalf("CloseIssues: %v", err) + } + tc.arrange(t, d, issue) + + n, err := EnqueueEntityStateBackfill(context.Background(), d) + if err != nil { + t.Fatalf("EnqueueEntityStateBackfill: %v", err) + } + if n != tc.wantRows { + t.Fatalf("enqueue affected %d rows, want %d", n, tc.wantRows) + } + + if tc.wantNoRow { + var count int + if err := d.QueryRow( + `select count(*) from pds_migration where name = ?`, EntityStateBackfillName, + ).Scan(&count); err != nil { + t.Fatalf("query: %v", err) + } + if count != 0 { + t.Fatalf("owner with no column-only closed work produced %d rows, want 0", count) + } + return + } + + var did, status string + if err := d.QueryRow( + `select did, status from pds_migration where name = ?`, EntityStateBackfillName, + ).Scan(&did, &status); err != nil { + t.Fatalf("query: %v", err) + } + if syntax.DID(did) != owner { + t.Fatalf("row did = %s, want owner %s", did, owner) + } + if status != string(tc.wantStatus) { + t.Fatalf("row status = %s, want %s", status, tc.wantStatus) + } + }) + } +} -- tangled.sh