From 99995bba3bf35adde9cfee3fadb10f7dd476b065 Mon Sep 17 00:00:00 2001 From: Bretton Date: Tue, 29 Sep 2026 03:33:43 -0700 Subject: [PATCH] feat(posts): rematerializer skips instance-removed legacy posts Implements Q-REMAT (user, 2026-09-28: skip). Without it, the one-shot rematerializer would turn an instance-removed legacy post into a fresh accepted postv2 URI with unblocked media. Changes: - Rematerializer requires a posts-local removal-lookup seam (satisfied by posts.Repository.ActiveRemovalsByURIs) and InstanceDID. A post is skipped only when it has an active removal with scope instance and authority InstanceDID; other authorities and scopes do not skip. - The check runs in the census pass before Ledger.Discover and at the top of RematerializeOne before its own Discover, covering the source pass, ledger resume, direct callers and rows reopened by -reopen-fallbacks. A post the census skipped is not retried later in the same run. - A skipped post gets no ledger row or advance, no author repo resolution, blob upload, postv2 write, acceptance, legacy delete or tombstone. It stays legacy and counts in RemainingLegacy, so ScopeComplete and Complete stay false. SkippedRemoved counts distinct URIs once per run. - Scoped runs never report Complete; main.go sends the operator to an unscoped run before PRD_AUTHOR_OWNED_POSTS section 11 step 6. - The removal is rechecked right before WriteAcceptance; a hit leaves the row at postv2_written and names the postv2 in the progress note. - Active removals that do not match (INSTANCE_DID, instance) are counted as UnmatchedRemovals and the census prints a warning, to catch INSTANCE_DID drift. - A removal-lookup error stops the run before mutating the record being checked; progress on earlier records is kept. - DryRunOf copies the new fields and walks the same path. - cmd/rematerialize-posts wires the post repository and cfg.Instance.DID, prints skipped-removed in the census and final verdict, and documents that a removal made after a postv2 was written does not carry to the new URI. - Tests: T0 coverage for P1 removed / P2 migrated in real and dry runs, other-authority and other-scope removals, resumed rows past discovered, reopened fallback rows, and lookup errors on the first and a later record; cmd census output tests. Reviewed with /second-opinion; findings fixed in review round 1. Co-Authored-By: Claude Opus 5.5 (1M context) --- cmd/rematerialize-posts/main.go | 81 +- cmd/rematerialize-posts/main_test.go | 214 +++ internal/core/posts/rematerialize.go | 216 ++- .../posts/rematerialize_blob_guard_test.go | 2 +- .../posts/rematerialize_credentials_test.go | 1 + internal/core/posts/rematerialize_dryrun.go | 2 + .../core/posts/rematerialize_dryrun_test.go | 4 +- .../core/posts/rematerialize_guard_test.go | 12 + .../core/posts/rematerialize_outer_test.go | 3 +- .../core/posts/rematerialize_removed_test.go | 1481 +++++++++++++++++ internal/core/posts/rematerialize_test.go | 29 +- 11 files changed, 2001 insertions(+), 44 deletions(-) create mode 100644 internal/core/posts/rematerialize_removed_test.go diff --git a/cmd/rematerialize-posts/main.go b/cmd/rematerialize-posts/main.go index 0f6504f..e9eb492 100644 --- a/cmd/rematerialize-posts/main.go +++ b/cmd/rematerialize-posts/main.go @@ -7,6 +7,17 @@ // index row before marking it done. The firehose no longer ingests legacy post // deletes, so the tool must converge both stores itself. // +// A legacy post with an active removal by the configured instance DID at +// instance scope is left as legacy and counted as skipped-removed, and as +// remaining legacy by any run whose scope includes it. Only an unscoped run can +// report the migration complete: a -community run's final re-scan covers only +// its own community, and a post skipped in another community leaves no ledger +// row. The removal is checked before the post's ledger row is discovered and +// again right before the community acceptance is written; when the second +// check fires, the postv2 already written in the author's repo is left there +// unaccepted and the legacy post stays in place. A removal made after the +// acceptance is written does not carry over to the new postv2 URI. +// // # THIS COMMAND DELETES PRODUCTION USER DATA IRREVERSIBLY // // It is run by hand, once, during a maintenance window, by someone who has been @@ -229,13 +240,16 @@ func main() { } progress := newProgressLogger() + postRepository := postgresRepo.NewPostRepository(db) tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authorFactory, Acceptances: writer, CommunityRepos: repoFactory, - Index: postgresRepo.NewPostRepository(db), + Index: postRepository, + Removals: postRepository, + InstanceDID: cfg.Instance.DID, CommunityScope: *communityFilter, // The blob copy's dev gate, decided HERE rather than inside the state // machine, exactly as blobs.PrivateHostOptions is decided above. @@ -294,25 +308,57 @@ func main() { os.Exit(1) } - // TWO SIGNALS, REPORTED SEPARATELY. A staged -community run finishing its own - // scope is a success even though the migration as a whole is not done, and - // collapsing the two taught the operator to ignore a red exit code on every - // staged run — which would leave §11 step 6 with no machine-checkable gate at - // all. + if exitCode := logVerdict(report); exitCode != 0 { + os.Exit(exitCode) + } +} + +// logVerdict prints the run's final verdict and returns the process exit code: +// 1 when the run's own scope is incomplete, 0 otherwise. +// +// TWO SIGNALS, REPORTED SEPARATELY. A staged -community run finishing its own +// scope is a success even though the migration as a whole is not done, and +// collapsing the two taught the operator to ignore a red exit code on every +// staged run — which would leave §11 step 6 with no machine-checkable gate at +// all. +// +// A scoped run never certifies the whole migration: its final re-scan covers +// only its own community, and a post skipped for a removal in another community +// leaves no ledger row. Only an unscoped run's verdict gates §11 step 6. +func logVerdict(report posts.RematerializeReport) (exitCode int) { if !report.ScopeComplete { - log.Printf("rematerialize-posts: SCOPE INCOMPLETE — %d of %d row(s) in scope reached done, %d fallback(s), %d legacy record(s) still standing", - report.Done, report.Discovered, report.Fallbacks, report.RemainingLegacy) - os.Exit(1) + log.Printf("rematerialize-posts: SCOPE INCOMPLETE — %d of %d row(s) in scope reached done, %d fallback(s), %d legacy record(s) still standing, %d skipped for an active instance removal", + report.Done, report.Discovered, report.Fallbacks, report.RemainingLegacy, report.SkippedRemoved) + logSkippedRemovedGuidance(report.SkippedRemoved) + return 1 } - log.Printf("rematerialize-posts: scope complete — every post in %s was re-materialized", scopeName(*communityFilter)) + log.Printf("rematerialize-posts: scope complete — every post in %s was re-materialized", scopeName(report.CommunityScope)) + if report.CommunityScope != "" { + log.Printf("rematerialize-posts: this run was scoped to %s, and a scoped run never certifies the whole migration.", report.CommunityScope) + log.Printf("rematerialize-posts: DO NOT run the legacy-removal follow-up (PRD_AUTHOR_OWNED_POSTS §11 step 6) until an UNSCOPED run (no -community) reports MIGRATION COMPLETE.") + return 0 + } if !report.Complete { - log.Printf("rematerialize-posts: THE MIGRATION AS A WHOLE IS NOT COMPLETE — %d of %d ledger row(s) done, %d fallback(s), %d legacy record(s) still standing.", - report.GlobalDone, report.GlobalDiscovered, report.GlobalFallbacks, report.RemainingLegacy) + log.Printf("rematerialize-posts: THE MIGRATION AS A WHOLE IS NOT COMPLETE — %d of %d ledger row(s) done, %d fallback(s), %d legacy record(s) still standing, %d skipped for an active instance removal.", + report.GlobalDone, report.GlobalDiscovered, report.GlobalFallbacks, report.RemainingLegacy, report.SkippedRemoved) + logSkippedRemovedGuidance(report.SkippedRemoved) log.Printf("rematerialize-posts: DO NOT run the legacy-removal follow-up (PRD §11 step 6) until this line says complete.") - return + return 0 } log.Printf("rematerialize-posts: MIGRATION COMPLETE — every discovered post was re-materialized and no legacy record remains") + return 0 +} + +// logSkippedRemovedGuidance tells the operator what the posts skipped for an +// active instance removal mean for the migration, when there are any. Without +// it every run exits red while a removal stands, with nothing saying why. +func logSkippedRemovedGuidance(skippedRemoved int) { + if skippedRemoved == 0 { + return + } + log.Printf("rematerialize-posts: the %d post(s) skipped for an active instance removal stay legacy while the removal stands and keep the migration incomplete; "+ + "restoring one and re-running migrates it, and PRD_AUTHOR_OWNED_POSTS §11 step 6 must not run while any stands.", skippedRemoved) } // repoClient is the narrow PDS surface the legacy source needs: enumerate, @@ -699,9 +745,14 @@ func logCensus(report posts.RematerializeReport, dryRun bool) { if dryRun { prefix = "census (dry run — no state was persisted)" } - log.Printf("rematerialize-posts: %s — scope=%s discovered=%d done=%d fallbacks=%d remaining-legacy=%d scope-complete=%v", + log.Printf("rematerialize-posts: %s — scope=%s discovered=%d done=%d fallbacks=%d remaining-legacy=%d skipped-removed=%d unmatched-removals=%d scope-complete=%v", prefix, scopeName(report.CommunityScope), report.Discovered, report.Done, report.Fallbacks, - report.RemainingLegacy, report.ScopeComplete) + report.RemainingLegacy, report.SkippedRemoved, report.UnmatchedRemovals, report.ScopeComplete) + if report.UnmatchedRemovals > 0 { + log.Printf("rematerialize-posts: WARNING — %d legacy post(s) carry an active removal that is not by this tool's INSTANCE_DID at instance scope. "+ + "They are hidden on the read path but were NOT skipped. A non-zero count usually means INSTANCE_DID differs from the server's.", + report.UnmatchedRemovals) + } for _, state := range censusOrder { if n, ok := report.ByState[state]; ok { log.Printf(" %-22s %d", state, n) diff --git a/cmd/rematerialize-posts/main_test.go b/cmd/rematerialize-posts/main_test.go index 9d51961..20998d8 100644 --- a/cmd/rematerialize-posts/main_test.go +++ b/cmd/rematerialize-posts/main_test.go @@ -1,8 +1,10 @@ package main import ( + "bytes" "context" "fmt" + "log" "strings" "testing" @@ -355,3 +357,215 @@ func TestOAuthScopes_GrantTheLegacyDeleteTheDrainDependsOn(t *testing.T) { assert.Containsf(t, legacy, "action=delete", "the legacy-post scope %q grants no delete; the drain's final step is exactly that delete", legacy) } + +// ---- the operator's census and verdict ------------------------------------ + +// captureLog sends the standard logger to a buffer, without timestamps, for the +// rest of the test. +func captureLog(t *testing.T) *bytes.Buffer { + t.Helper() + previousOutput, previousFlags := log.Writer(), log.Flags() + t.Cleanup(func() { + log.SetOutput(previousOutput) + log.SetFlags(previousFlags) + }) + log.SetFlags(0) + var output bytes.Buffer + log.SetOutput(&output) + return &output +} + +// lineContaining returns the first output line containing phrase, or "". +func lineContaining(output, phrase string) string { + for _, line := range strings.Split(output, "\n") { + if strings.Contains(line, phrase) { + return line + } + } + return "" +} + +// censusFields parses the label=value tokens of the census summary line. +func censusFields(t *testing.T, output string) map[string]string { + t.Helper() + line := lineContaining(output, "remaining-legacy=") + require.NotEmptyf(t, line, "no census line in output:\n%s", output) + fields := map[string]string{} + for _, token := range strings.Fields(line) { + if label, value, ok := strings.Cut(token, "="); ok { + fields[label] = value + } + } + return fields +} + +// Every count is distinct, so a census that prints one count under another +// count's label fails here. +func TestLogCensus_PrintsEachCountUnderItsOwnLabel(t *testing.T) { + output := captureLog(t) + report := posts.RematerializeReport{ + Discovered: 11, + Done: 2, + Fallbacks: 3, + RemainingLegacy: 5, + SkippedRemoved: 7, + UnmatchedRemovals: 13, + } + for _, mode := range []struct { + name string + dryRun bool + }{ + {name: "live", dryRun: false}, + {name: "dry_run", dryRun: true}, + } { + t.Run(mode.name, func(t *testing.T) { + output.Reset() + logCensus(report, mode.dryRun) + fields := censusFields(t, output.String()) + for label, want := range map[string]string{ + "discovered": "11", + "done": "2", + "fallbacks": "3", + "remaining-legacy": "5", + "skipped-removed": "7", + "unmatched-removals": "13", + "scope-complete": "false", + } { + assert.Equalf(t, want, fields[label], "census label %q in:\n%s", label, output.String()) + } + assert.Contains(t, output.String(), "skipped-removed=7 unmatched-removals=13", + "unmatched-removals belongs right after skipped-removed") + }) + } +} + +// A removal that is not (INSTANCE_DID, instance) hides the post on the read path +// but does not stop the tool migrating it. Moderation writes only that pair, so +// the census must say so when it sees one. +func TestLogCensus_WarnsAboutUnmatchedRemovalsOnlyWhenThereAreAny(t *testing.T) { + output := captureLog(t) + for _, tc := range []struct { + name string + unmatched int + wantWarning bool + }{ + {name: "none", unmatched: 0, wantWarning: false}, + {name: "some", unmatched: 13, wantWarning: true}, + } { + t.Run(tc.name, func(t *testing.T) { + output.Reset() + logCensus(posts.RematerializeReport{Discovered: 11, UnmatchedRemovals: tc.unmatched}, false) + warning := lineContaining(output.String(), "INSTANCE_DID") + if !tc.wantWarning { + assert.Emptyf(t, warning, "no unmatched removal was seen, so no warning belongs in:\n%s", output.String()) + return + } + require.NotEmptyf(t, warning, "no INSTANCE_DID warning in:\n%s", output.String()) + assert.Contains(t, warning, "13 legacy post(s)") + assert.Contains(t, warning, "hidden on the read path") + assert.Contains(t, warning, "NOT skipped") + assert.Contains(t, warning, "differs from the server's") + }) + } +} + +// A scoped run's final re-scan covers only its own community, so it can never +// certify the whole migration. Its verdict must send the operator to an +// unscoped run, not tell them to wait for a line a scoped run never prints. +func TestLogVerdict_AScopedRunSendsTheOperatorToAnUnscopedRun(t *testing.T) { + output := captureLog(t) + for _, complete := range []bool{false, true} { + t.Run(fmt.Sprintf("complete=%v", complete), func(t *testing.T) { + output.Reset() + exitCode := logVerdict(posts.RematerializeReport{ + CommunityScope: "did:plc:inscope22222222222222222", + Discovered: 4, + Done: 4, + ScopeComplete: true, + GlobalDiscovered: 9, + GlobalDone: 4, + Complete: complete, + }) + text := output.String() + assert.Equal(t, 0, exitCode, "a staged run that finished its own scope is a success") + assert.NotContains(t, text, "until this line says complete") + assert.NotContains(t, text, "rematerialize-posts: MIGRATION COMPLETE") + assert.Contains(t, text, "UNSCOPED") + assert.Contains(t, text, "§11 step 6") + }) + } +} + +// While an instance removal stands every run exits 1 or reports NOT COMPLETE, so +// the final line must say how many posts that is and what the operator does +// about it. +func TestLogVerdict_CountsThePostsSkippedForAnInstanceRemoval(t *testing.T) { + output := captureLog(t) + for _, tc := range []struct { + name string + report posts.RematerializeReport + verdictPhrase string + wantExitCode int + }{ + { + name: "scope incomplete", + report: posts.RematerializeReport{Discovered: 11, Done: 2, Fallbacks: 3, RemainingLegacy: 5}, + verdictPhrase: "SCOPE INCOMPLETE", + wantExitCode: 1, + }, + { + name: "whole migration not complete", + report: posts.RematerializeReport{ + Discovered: 11, Done: 11, ScopeComplete: true, + GlobalDiscovered: 11, GlobalDone: 10, GlobalFallbacks: 1, + }, + verdictPhrase: "NOT COMPLETE", + wantExitCode: 0, + }, + } { + for _, skipped := range []int{0, 7} { + t.Run(fmt.Sprintf("%s/skipped=%d", tc.name, skipped), func(t *testing.T) { + output.Reset() + report := tc.report + report.SkippedRemoved = skipped + exitCode := logVerdict(report) + text := output.String() + assert.Equal(t, tc.wantExitCode, exitCode) + verdict := lineContaining(text, tc.verdictPhrase) + require.NotEmptyf(t, verdict, "no %q line in:\n%s", tc.verdictPhrase, text) + assert.Contains(t, verdict, fmt.Sprintf("%d skipped for an active instance removal", skipped)) + guidance := lineContaining(text, "stay legacy while") + if skipped == 0 { + assert.Emptyf(t, guidance, "nothing was skipped, so no guidance belongs in:\n%s", text) + return + } + require.NotEmptyf(t, guidance, "no guidance for the skipped posts in:\n%s", text) + assert.Contains(t, guidance, "7 post(s)") + assert.Contains(t, guidance, "re-running migrates it") + assert.Contains(t, guidance, "§11 step 6 must not run") + }) + } + } +} + +// The unscoped run is the §11 step 6 gate, so its two final lines stay as they +// were. +func TestLogVerdict_AnUnscopedRunKeepsItsCompletionLines(t *testing.T) { + output := captureLog(t) + + exitCode := logVerdict(posts.RematerializeReport{ + Discovered: 11, Done: 11, ScopeComplete: true, + GlobalDiscovered: 11, GlobalDone: 10, GlobalFallbacks: 1, + }) + assert.Equal(t, 0, exitCode) + assert.Contains(t, output.String(), "THE MIGRATION AS A WHOLE IS NOT COMPLETE") + assert.Contains(t, output.String(), "until this line says complete") + + output.Reset() + exitCode = logVerdict(posts.RematerializeReport{ + Discovered: 11, Done: 11, ScopeComplete: true, + GlobalDiscovered: 11, GlobalDone: 11, Complete: true, + }) + assert.Equal(t, 0, exitCode) + assert.Contains(t, output.String(), "rematerialize-posts: MIGRATION COMPLETE") +} diff --git a/internal/core/posts/rematerialize.go b/internal/core/posts/rematerialize.go index 0310119..aa04bd5 100644 --- a/internal/core/posts/rematerialize.go +++ b/internal/core/posts/rematerialize.go @@ -70,7 +70,8 @@ import ( // RematerializeState is one legacy record's position in the ledger state machine // (migration 037). The happy path is discovered → postv2_written → verified → -// migrated → done; the fallback state is terminal. +// migrated → done; the fallback state is terminal. RematerializeSkippedRemoved +// is the one exception: an outcome RematerializeOne returns, never stored. type RematerializeState string const ( @@ -115,8 +116,23 @@ const ( // entry nothing produces is a trap for whoever writes recovery SQL at 2am // against a state that cannot exist. RematerializeFallbackLeftLegacy RematerializeState = "fallback_left_legacy" + + // RematerializeSkippedRemoved is returned for a legacy post with an active + // instance removal. It is an outcome, never a ledger state. + RematerializeSkippedRemoved RematerializeState = "skipped_removed" ) +// RematerializeRemovalLookup reports the active removal sources of each URI that +// has any. posts.Repository.ActiveRemovalsByURIs satisfies it. +type RematerializeRemovalLookup interface { + ActiveRemovalsByURIs(ctx context.Context, uris []string) (map[string][]RemovalSource, error) +} + +// instanceRemovalScopeKind is moderation.ScopeInstance. Posts cannot import +// moderation (moderation imports posts); migration 050's scope_kind CHECK +// ('instance', 'community') keeps the two from drifting. +const instanceRemovalScopeKind = "instance" + // IsFallback reports whether a state is a terminal fallback state, so the census // can gate "complete" on any of them surviving without enumerating each string // at every call site. @@ -312,7 +328,7 @@ type RematerializeLedger interface { // while the migration as a whole still has thousands of posts to go (Complete). // Reporting only the second makes every staged run look like a failure, and an // operator who has learned to ignore a red exit code is an operator with no gate -// on §11 step 6 at all. +// on §11 step 6 at all. Only an unscoped run can report Complete. type RematerializeReport struct { // CommunityScope is the community DID this run was restricted to, or "" for // every hosted community. @@ -325,7 +341,9 @@ type RematerializeReport struct { ByState map[RematerializeState]int // RemainingLegacy is how many legacy records a FINAL RE-SCAN of the source - // still saw that are not accounted for by a fallback row. It is what turns + // still saw that are not accounted for by a fallback row. A record this run + // skipped for an instance removal counts even when it has a fallback row, + // because it stays legacy. It is what turns // "the ledger says we are done" into "the source agrees" — a ledger-only // completion check cannot see a record written after the run began, or one // the discovery pass never listed. @@ -346,7 +364,24 @@ type RematerializeReport struct { // follow-up (§11 step 6). It requires the whole migration — not this run's // scope — to have reached done, with no fallback surviving and nothing left in // the source. + // + // A SCOPED RUN NEVER REPORTS IT. Its re-scan lists only its own community, so + // it cannot see a legacy record standing elsewhere that no ledger row + // accounts for: a post another scoped run skipped for an instance removal, or + // one no run has discovered yet. The §11 step 6 gate must come from an + // unscoped run. Complete bool + + // SkippedRemoved counts the distinct legacy posts this run skipped because + // they have an active instance removal. + SkippedRemoved int + + // UnmatchedRemovals counts the distinct legacy posts the census pass saw + // with an active removal, none of them by InstanceDID at instance scope. + // Such a post is hidden on the read path but NOT skipped here. Moderation + // writes only (instance DID, instance) removals, so a non-zero count means + // this tool's INSTANCE_DID differs from the server's. + UnmatchedRemovals int } // RematerializeProgress is one observable transition, handed to the caller's @@ -423,6 +458,13 @@ type Rematerializer struct { // set, the run stops and names the authors instead. AbortOnFallback bool + // Removals looks up active removals. Required. + Removals RematerializeRemovalLookup + + // InstanceDID is the authority whose instance-scoped removals skip a post. + // Required. + InstanceDID string + // credentials caches ONE resolution per distinct author DID for the lifetime // of this Rematerializer. Each resolution is a refresh-token rotation against // the PDS; an aggregator with 5,000 posts would otherwise rotate 5,000 times @@ -525,11 +567,25 @@ func RematerializeRkey(legacyPostURI string) string { // permission to run the irreversible legacy-removal step while a record sits // half-migrated. A no-creds fallback is NOT such an error — it is an expected // terminal outcome the census counts. +// +// A post with an active instance removal by InstanceDID is skipped before its +// ledger row is touched (RematerializeOne says where it checks): it stays +// legacy, and SkippedRemoved counts it. Every run whose scope includes it +// counts it in RemainingLegacy while its record stands, so ScopeComplete stays +// false; only an unscoped run can report Complete. A post the census skips is +// not retried by the source pass or the ledger reconcile later in the same run, +// even if its removal is lifted meanwhile: its author never went through the +// credential census. func (r *Rematerializer) Run(ctx context.Context) (RematerializeReport, error) { + if err := r.requireRemovalSeam(); err != nil { + return RematerializeReport{}, err + } legacies, err := r.Source.ListLegacyPosts(ctx) if err != nil { return RematerializeReport{}, fmt.Errorf("enumerating legacy posts: %w", err) } + skippedRemoved := make(map[string]bool) + unmatchedRemovals := make(map[string]bool) // Pass 1 — the census. Resolve EVERY not-yet-started author before ANY repo is // mutated. A row already past discovered had its credentials confirmed on an @@ -541,6 +597,29 @@ func (r *Rematerializer) Run(ctx context.Context) (RematerializeReport, error) { if err := ctx.Err(); err != nil { return RematerializeReport{}, err } + removed, unmatched, err := r.instanceRemoved(ctx, legacy.URI) + if err != nil { + return RematerializeReport{}, err + } + if removed { + note, err := r.removedSkipNote(ctx, legacy.URI) + if err != nil { + return RematerializeReport{}, err + } + skippedRemoved[legacy.URI] = true + r.report(RematerializeProgress{ + OldURI: legacy.URI, To: RematerializeSkippedRemoved, + Index: i + 1, Total: len(legacies), Note: note, + }) + continue + } + if unmatched { + unmatchedRemovals[legacy.URI] = true + r.report(RematerializeProgress{ + OldURI: legacy.URI, Index: i + 1, Total: len(legacies), + Note: "active removal not by the instance DID at instance scope; migrating", + }) + } row, err := r.Ledger.Discover(ctx, legacy.URI, legacy.CommunityDID, legacy.AuthorDID) if err != nil { return RematerializeReport{}, err @@ -576,15 +655,23 @@ func (r *Rematerializer) Run(ctx context.Context) (RematerializeReport, error) { } // Pass 2 — the source pass. A record whose census marked it a fallback is left - // untouched by RematerializeOne (it returns early on a terminal row). + // untouched by RematerializeOne (it returns early on a terminal row). A record + // the census skipped for a removal is not handed to it at all: the census + // already reported it, and its author was never preflighted. for i, legacy := range legacies { if err := ctx.Err(); err != nil { return RematerializeReport{}, err } + if skippedRemoved[legacy.URI] { + continue + } state, err := r.rematerializeOneBounded(ctx, legacy) if err != nil { return RematerializeReport{}, fmt.Errorf("re-materializing %s: %w", legacy.URI, err) } + if state == RematerializeSkippedRemoved { + skippedRemoved[legacy.URI] = true + } r.report(RematerializeProgress{OldURI: legacy.URI, To: state, Index: i + 1, Total: len(legacies)}) } @@ -606,15 +693,26 @@ func (r *Rematerializer) Run(ctx context.Context) (RematerializeReport, error) { if ledgerRow.State == RematerializeDiscovered { continue } + // The census verdict holds here too: a listed post it skipped is not + // finished from its ledger row even if the removal was lifted since. + if skippedRemoved[ledgerRow.OldURI] { + continue + } legacy := legacyFromLedgerRow(ledgerRow) state, err := r.rematerializeOneBounded(ctx, legacy) if err != nil { return RematerializeReport{}, fmt.Errorf("reconciling %s: %w", ledgerRow.OldURI, err) } + if state == RematerializeSkippedRemoved { + // RematerializeOne reported the skip, naming any postv2 that stands. The + // row did not move, so there is no reconcile transition to report. + skippedRemoved[ledgerRow.OldURI] = true + continue + } r.report(RematerializeProgress{OldURI: ledgerRow.OldURI, From: ledgerRow.State, To: state, Index: i + 1, Total: len(resumable), Note: "reconciled from the ledger"}) } - return r.census(ctx) + return r.census(ctx, skippedRemoved, unmatchedRemovals) } // rematerializeOneBounded runs one record under PerRecordTimeout, so a single @@ -633,7 +731,7 @@ func (r *Rematerializer) rematerializeOneBounded(ctx context.Context, legacy Leg // census builds the report: the scoped tally, the global tally, and a FINAL // RE-SCAN of the source that the completion signals are gated on. -func (r *Rematerializer) census(ctx context.Context) (RematerializeReport, error) { +func (r *Rematerializer) census(ctx context.Context, skippedRemoved, unmatchedRemovals map[string]bool) (RematerializeReport, error) { scoped, err := r.Ledger.CountByState(ctx, r.CommunityScope) if err != nil { return RematerializeReport{}, fmt.Errorf("taking the census: %w", err) @@ -647,9 +745,11 @@ func (r *Rematerializer) census(ctx context.Context) (RematerializeReport, error } report := RematerializeReport{ - CommunityScope: r.CommunityScope, - ByState: scoped, - GlobalByState: global, + CommunityScope: r.CommunityScope, + ByState: scoped, + GlobalByState: global, + SkippedRemoved: len(skippedRemoved), + UnmatchedRemovals: len(unmatchedRemovals), } for state, n := range scoped { report.Discovered += n @@ -683,9 +783,10 @@ func (r *Rematerializer) census(ctx context.Context) (RematerializeReport, error if err != nil { return report, fmt.Errorf("checking the ledger for the re-scanned %s: %w", legacy.URI, err) } - // A record deliberately left as legacy is accounted for, not remaining. A - // record with no row at all, or one not yet done, is remaining. - if found && IsFallback(row.State) { + // A credential fallback is accounted for, except when this run skipped + // the post for an instance removal: that legacy record still remains. + // A record with no row at all, or one not yet done, is remaining. + if found && IsFallback(row.State) && !skippedRemoved[legacy.URI] { continue } report.RemainingLegacy++ @@ -696,15 +797,27 @@ func (r *Rematerializer) census(ctx context.Context) (RematerializeReport, error // no fallback surviving, and the source re-scan agreeing. A row stranded in any // non-terminal state, a surviving fallback, or a legacy record still standing all // leave it false, and the operator's irreversible legacy-removal step (§11 step 6) - // must not run while any of them is true. + // must not run while any of them is true. A scoped run never reports it: its + // re-scan cannot see a legacy record standing in another community. report.Complete = report.GlobalDone == report.GlobalDiscovered && report.GlobalFallbacks == 0 && - report.RemainingLegacy == 0 + report.RemainingLegacy == 0 && + r.CommunityScope == "" return report, nil } // RematerializeOne drives a single legacy record from wherever its ledger row -// stands to a terminal state, and returns the state it reached. +// stands to a terminal state, and returns the ledger state the row reached, or +// RematerializeSkippedRemoved, which is not a ledger state. +// +// The removal is checked twice. At the top, before Discover: a removed post +// gets no ledger row, and an existing row stays exactly where it stood. Right +// before the acceptance write, the step that admits the postv2: a removal that +// landed in between leaves the row at postv2_written, with the postv2 it names +// standing in the author's repo without this call's acceptance, and the legacy +// record untouched. +// Either skip returns RematerializeSkippedRemoved. A removal that lands after +// the acceptance is not honoured by this call. // // The steps are guarded on the ledger state each moves FROM, so a resumed run // re-enters at exactly the step its predecessor stopped before and re-does none @@ -713,6 +826,9 @@ func (r *Rematerializer) census(ctx context.Context) (RematerializeReport, error // because the ledger records that verification passed at a moment now in the // past and the delete needs it to be true now. func (r *Rematerializer) RematerializeOne(ctx context.Context, legacy LegacyPost) (RematerializeState, error) { + if err := r.requireRemovalSeam(); err != nil { + return "", err + } if r.CommunityScope != "" && legacy.CommunityDID != r.CommunityScope { // A SCOPED RUN NEVER TOUCHES ANOTHER COMMUNITY'S POSTS. This is the last // gate before a delete, and it is checked here rather than only at discovery @@ -721,6 +837,20 @@ func (r *Rematerializer) RematerializeOne(ctx context.Context, legacy LegacyPost "refusing to re-materialize %s: it belongs to %s but this run is scoped to %s", legacy.URI, legacy.CommunityDID, r.CommunityScope) } + // An instance-removed post is left legacy, checked before Discover so it gets + // no ledger row and a resumed row is not advanced. + removed, _, err := r.instanceRemoved(ctx, legacy.URI) + if err != nil { + return "", err + } + if removed { + note, err := r.removedSkipNote(ctx, legacy.URI) + if err != nil { + return "", err + } + r.report(RematerializeProgress{OldURI: legacy.URI, To: RematerializeSkippedRemoved, Note: note}) + return RematerializeSkippedRemoved, nil + } row, err := r.Ledger.Discover(ctx, legacy.URI, legacy.CommunityDID, legacy.AuthorDID) if err != nil { @@ -764,6 +894,21 @@ func (r *Rematerializer) RematerializeOne(ctx context.Context, legacy LegacyPost // mean "verified" on the strength of nothing. The state advances only after // the read-back below. if row.State == RematerializePostV2Written { + // The acceptance is what admits the postv2, and it can follow the check at + // the top by up to PerRecordTimeout. A removal that landed meanwhile stops + // here, with the row left at postv2_written. + removed, _, err := r.instanceRemoved(ctx, legacy.URI) + if err != nil { + return row.State, err + } + if removed { + r.report(RematerializeProgress{ + OldURI: legacy.URI, To: RematerializeSkippedRemoved, + Note: "active instance removal landed before the acceptance; leaving legacy post untouched; the postv2 at " + + row.NewURI + " stands without this run's acceptance and is NOT removed", + }) + return RematerializeSkippedRemoved, nil + } if _, err := r.Acceptances.WriteAcceptance(ctx, CommunityWriteCommand{ CommunityDID: legacy.CommunityDID, PostURI: row.NewURI, @@ -833,6 +978,47 @@ func (r *Rematerializer) RematerializeOne(ctx context.Context, legacy LegacyPost return row.State, nil } +// instanceRemoved reports whether uri has an active removal by InstanceDID at +// instance scope (removed), and whether it has active removals of which none is +// that (unmatched). A lookup error is returned, never treated as "not removed". +func (r *Rematerializer) instanceRemoved(ctx context.Context, uri string) (removed, unmatched bool, err error) { + removals, err := r.Removals.ActiveRemovalsByURIs(ctx, []string{uri}) + if err != nil { + return false, false, fmt.Errorf("checking active removals for %s: %w", uri, err) + } + for _, removal := range removals[uri] { + if removal.AuthorityDID == r.InstanceDID && removal.ScopeKind == instanceRemovalScopeKind { + return true, false, nil + } + } + return false, len(removals[uri]) > 0, nil +} + +// removedSkipNote is the progress note for a post skipped for an instance +// removal. When an earlier run already wrote the post's postv2, the note names +// it: that postv2 stands and the removal does not cover it. +func (r *Rematerializer) removedSkipNote(ctx context.Context, uri string) (string, error) { + note := "active instance removal; leaving legacy post untouched" + existing, found, err := r.Ledger.Get(ctx, uri) + if err != nil { + return "", fmt.Errorf("reading the ledger row of the removed %s: %w", uri, err) + } + if found && existing.NewURI != "" { + note += "; a postv2 already stands at " + existing.NewURI + " and is NOT removed" + } + return note, nil +} + +func (r *Rematerializer) requireRemovalSeam() error { + if r.Removals == nil { + return errors.New("rematerializer requires Removals") + } + if r.InstanceDID == "" { + return errors.New("rematerializer requires InstanceDID") + } + return nil +} + // writePostV2 is step 1: resolve the author, re-read the legacy record, copy its // blobs, build the lossless conversion and write it at the deterministic rkey. // diff --git a/internal/core/posts/rematerialize_blob_guard_test.go b/internal/core/posts/rematerialize_blob_guard_test.go index ce90f3f..05aa25d 100644 --- a/internal/core/posts/rematerialize_blob_guard_test.go +++ b/internal/core/posts/rematerialize_blob_guard_test.go @@ -163,7 +163,7 @@ func TestDefaultRematerializeBlobClient_GuardedIsTheDefaultForTheStateMachine(t host := newCountingBlobHost(t) - fallback := (&Rematerializer{}).blobClient() + fallback := (&Rematerializer{Removals: noRematerializeRemovals{}, InstanceDID: "did:web:coves-instance.invalid"}).blobClient() require.NotNil(t, fallback, "a Rematerializer with no injected Blobs must still have a client") _, err := fallback.Fetch(context.Background(), host.server.URL, diff --git a/internal/core/posts/rematerialize_credentials_test.go b/internal/core/posts/rematerialize_credentials_test.go index 425fd32..ce68def 100644 --- a/internal/core/posts/rematerialize_credentials_test.go +++ b/internal/core/posts/rematerialize_credentials_test.go @@ -93,6 +93,7 @@ func TestClassifyResumeFailure_NilErrorIsNotAFailure(t *testing.T) { func TestRematerializer_ResolvesCredentialsOncePerAuthor(t *testing.T) { resolutions := map[string]int{} tool := &Rematerializer{ + Removals: noRematerializeRemovals{}, InstanceDID: "did:web:coves-instance.invalid", AuthorRepos: func(_ context.Context, did string, _ *oauth.ClientSessionData) (AuthorRepo, error) { resolutions[did]++ return nil, fmt.Errorf("no repo in this unit test: %w", ErrNoAuthorCredentials) diff --git a/internal/core/posts/rematerialize_dryrun.go b/internal/core/posts/rematerialize_dryrun.go index e1516ca..c821a46 100644 --- a/internal/core/posts/rematerialize_dryrun.go +++ b/internal/core/posts/rematerialize_dryrun.go @@ -72,6 +72,8 @@ func DryRunOf(tool *Rematerializer) *Rematerializer { Progress: tool.Progress, PerRecordTimeout: tool.PerRecordTimeout, AbortOnFallback: tool.AbortOnFallback, + Removals: tool.Removals, + InstanceDID: tool.InstanceDID, } return dry } diff --git a/internal/core/posts/rematerialize_dryrun_test.go b/internal/core/posts/rematerialize_dryrun_test.go index 7b9b1e3..5131367 100644 --- a/internal/core/posts/rematerialize_dryrun_test.go +++ b/internal/core/posts/rematerialize_dryrun_test.go @@ -265,7 +265,9 @@ func dryRunFixture() (*Rematerializer, *memSource, *memLedger, *memAuthorRepo, * CommunityRepos: func(context.Context, string) (CommunityRepo, error) { return communityRepo, nil }, - Blobs: blobClient, + Blobs: blobClient, + Removals: noRematerializeRemovals{}, + InstanceDID: "did:web:coves-instance.invalid", } return tool, source, ledger, authorRepo, writer, blobClient } diff --git a/internal/core/posts/rematerialize_guard_test.go b/internal/core/posts/rematerialize_guard_test.go index b36599d..a257ad6 100644 --- a/internal/core/posts/rematerialize_guard_test.go +++ b/internal/core/posts/rematerialize_guard_test.go @@ -64,6 +64,7 @@ func newGuardHarness(t *testing.T, authorDID string) *guardHarness { h.newCID = deterministicCID(h.rkey) h.tool = &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), } return h @@ -543,6 +544,7 @@ func TestRematerialize_ScopedRun_RefusesARecordFromAnotherCommunity(t *testing.T tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), CommunityScope: rematCommunityDID, } @@ -581,6 +583,7 @@ func TestRematerialize_ScopedRun_ReportsScopeAndWholeMigrationSeparately(t *test tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), CommunityScope: rematCommunityDID, } @@ -622,6 +625,7 @@ func TestRematerialize_Complete_IsGatedOnARescanOfTheSource(t *testing.T) { tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), } @@ -663,6 +667,7 @@ func TestRematerialize_NoCredentials_WritesNothingEvenWhenOtherIdentitiesAreWrit source := newFakeLegacySource(legacy) tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), } @@ -700,6 +705,7 @@ func TestRematerialize_RetryableCredentialFailure_FailsTheRunAndSentencesNothing source := newFakeLegacySource(one, two) tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), } @@ -745,6 +751,7 @@ func TestRematerialize_AbortOnFallback_StopsBeforeAnyRepoIsMutated(t *testing.T) tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), AbortOnFallback: true, } @@ -778,6 +785,7 @@ func TestRematerialize_ReopenFallback_LetsAReAuthorizedAuthorBeRetried(t *testin source := newFakeLegacySource(legacy) tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), } @@ -801,6 +809,7 @@ func TestRematerialize_ReopenFallback_LetsAReAuthorizedAuthorBeRetried(t *testin // of extra token rotations for an author who has none. retry := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), } state, err := retry.RematerializeOne(ctx, legacy) @@ -851,6 +860,7 @@ func TestRematerialize_BlobProbeFails_RefusesToDelete(t *testing.T) { tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), Blobs: blobClient, } @@ -878,6 +888,7 @@ func TestRematerialize_FetchesCommunityBlobsFromTheCommunitysHost(t *testing.T) tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), Blobs: blobClient, } @@ -921,6 +932,7 @@ func TestRematerialize_ResumeAtPostV2Written_StillVerifiesTheBlobs(t *testing.T) tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", Acceptances: writer, CommunityRepos: writer.repos(), Blobs: blobClient, } diff --git a/internal/core/posts/rematerialize_outer_test.go b/internal/core/posts/rematerialize_outer_test.go index d9c4873..4ea4012 100644 --- a/internal/core/posts/rematerialize_outer_test.go +++ b/internal/core/posts/rematerialize_outer_test.go @@ -175,7 +175,7 @@ func TestRematerialize_OuterContract_RealPDS_MovesPostAndIsIdempotent(t *testing source := &realLegacySource{community: communityGeneric, staged: []posts.LegacyPost{legacy}} ledger := postgres.NewRematerializeLedger(testkit.DB(t)) communityRepos := func(_ context.Context, _ string) (posts.CommunityRepo, error) { return communityRepo, nil } - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authorFactory, Acceptances: writer, CommunityRepos: communityRepos} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authorFactory, Acceptances: writer, CommunityRepos: communityRepos, Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} // ---- run ----------------------------------------------------------------- state, err := tool.RematerializeOne(ctx, legacy) @@ -325,6 +325,7 @@ func TestRematerialize_OuterContract_CopiesEmbedBlobToAuthorRepo(t *testing.T) { Source: source, Ledger: ledger, AuthorRepos: authorFactory, Acceptances: writer, CommunityRepos: communityRepos, Blobs: posts.DefaultRematerializeBlobClient(true), + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", } _, err = tool.RematerializeOne(ctx, legacy) diff --git a/internal/core/posts/rematerialize_removed_test.go b/internal/core/posts/rematerialize_removed_test.go new file mode 100644 index 0000000..9dbc037 --- /dev/null +++ b/internal/core/posts/rematerialize_removed_test.go @@ -0,0 +1,1481 @@ +package posts + +import ( + "context" + "errors" + "fmt" + "strings" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/auth/oauth" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "Coves/internal/atproto/pds" +) + +// The source must actually shrink on delete: the final re-scan is what keeps a +// skipped legacy record in RemainingLegacy while excluding the migrated one. +type removedContractSource struct { + *memSource + deletedURIs []string +} + +func (s *removedContractSource) ListLegacyPosts(ctx context.Context) ([]LegacyPost, error) { + posts, err := s.memSource.ListLegacyPosts(ctx) + return append([]LegacyPost(nil), posts...), err +} + +func (s *removedContractSource) DeleteLegacyPost(ctx context.Context, legacy LegacyPost, swapCID string) error { + if err := s.memSource.DeleteLegacyPost(ctx, legacy, swapCID); err != nil { + return err + } + s.deletedURIs = append(s.deletedURIs, legacy.URI) + for index, post := range s.posts { + if post.URI == legacy.URI { + s.posts = append(s.posts[:index], s.posts[index+1:]...) + break + } + } + return nil +} + +type removedContractLedger struct { + *memLedger + discoveredURIs []string +} + +func (l *removedContractLedger) Discover(ctx context.Context, oldURI, communityDID, authorDID string) (RematerializeLedgerRow, error) { + l.discoveredURIs = append(l.discoveredURIs, oldURI) + return l.memLedger.Discover(ctx, oldURI, communityDID, authorDID) +} + +type removedContractAcceptanceWriter struct { + *memAcceptanceWriter + postURIs []string +} + +func (w *removedContractAcceptanceWriter) WriteAcceptance(ctx context.Context, command CommunityWriteCommand) (CommunityWriteResult, error) { + w.postURIs = append(w.postURIs, command.PostURI) + return w.memAcceptanceWriter.WriteAcceptance(ctx, command) +} + +type removedContractIndex struct{ tombstonedURIs []string } + +func (i *removedContractIndex) SoftDelete(_ context.Context, uri string) error { + i.tombstonedURIs = append(i.tombstonedURIs, uri) + return nil +} + +type removedContractLookup struct{ byURI map[string][]RemovalSource } + +func (l *removedContractLookup) ActiveRemovalsByURIs(_ context.Context, uris []string) (map[string][]RemovalSource, error) { + result := make(map[string][]RemovalSource) + for _, uri := range uris { + if sources, found := l.byURI[uri]; found { + result[uri] = sources + } + } + return result, nil +} + +func TestRematerialize_OuterContract_SkipsInstanceRemovedLegacyPost(t *testing.T) { + ctx := context.Background() + instanceDID := "did:web:coves-instance.invalid" + communityDID := "did:plc:community2222222222222222" + removedAuthorDID := "did:plc:removedauthor111111111111" + remainingAuthorDID := "did:plc:remainingauthor2222222222" + legacyPost := func(authorDID, rkey string) LegacyPost { + return LegacyPost{ + URI: "at://" + communityDID + "/" + LegacyPostCollection + "/" + rkey, + CID: "bafylegacy" + rkey, + CommunityDID: communityDID, + AuthorDID: authorDID, + RawRecord: map[string]any{ + "$type": LegacyPostCollection, + "community": communityDID, + "author": authorDID, + "title": "a legacy post", + "createdAt": "2026-01-02T03:04:05Z", + }, + } + } + removed := legacyPost(removedAuthorDID, "3kremoved") + remaining := legacyPost(remainingAuthorDID, "3kremaining") + source := &removedContractSource{memSource: &memSource{posts: []LegacyPost{removed, remaining}}} + ledger := &removedContractLedger{memLedger: newMemLedger()} + communityRepo := &memCommunityRepo{did: communityDID, records: map[string]*pds.RecordResponse{}} + writer := &removedContractAcceptanceWriter{memAcceptanceWriter: &memAcceptanceWriter{repo: communityRepo}} + index := &removedContractIndex{} + authorRepos := map[string]*memAuthorRepo{ + removedAuthorDID: newMemAuthorRepo(removedAuthorDID), + remainingAuthorDID: newMemAuthorRepo(remainingAuthorDID), + } + var resolvedAuthorDIDs []string + var progress []RematerializeProgress + tool := &Rematerializer{ + Source: source, + Ledger: ledger, + AuthorRepos: func(_ context.Context, authorDID string, _ *oauth.ClientSessionData) (AuthorRepo, error) { + resolvedAuthorDIDs = append(resolvedAuthorDIDs, authorDID) + return authorRepos[authorDID], nil + }, + Acceptances: writer, + CommunityRepos: func(context.Context, string) (CommunityRepo, error) { + return communityRepo, nil + }, + Index: index, + Removals: &removedContractLookup{byURI: map[string][]RemovalSource{removed.URI: {{AuthorityDID: instanceDID, ScopeKind: "instance"}}}}, + InstanceDID: instanceDID, + Progress: func(event RematerializeProgress) { + progress = append(progress, event) + }, + } + + report, err := tool.Run(ctx) + require.NoError(t, err) + + _, foundRemoved, err := ledger.Get(ctx, removed.URI) + require.NoError(t, err) + require.False(t, foundRemoved, "instance-removed legacy post must never get a ledger row") + assert.NotContains(t, ledger.discoveredURIs, removed.URI, "Discover must never be called for an instance-removed post") + assert.NotContains(t, resolvedAuthorDIDs, removedAuthorDID, "removed author's repo must never be resolved") + assert.Zero(t, authorRepos[removedAuthorDID].puts, "removed author's repo must never receive a postv2") + for _, postURI := range writer.postURIs { + assert.False(t, strings.HasPrefix(postURI, "at://"+removedAuthorDID+"/"), "acceptance must not be written for the removed author's post: %s", postURI) + } + assert.NotContains(t, source.deletedURIs, removed.URI) + assert.NotContains(t, index.tombstonedURIs, removed.URI) + _, stillPresent, err := source.ReadLegacyPost(ctx, removed.URI) + require.NoError(t, err) + assert.True(t, stillPresent, "removed post must remain in the legacy source") + + row, foundRemaining, err := ledger.Get(ctx, remaining.URI) + require.NoError(t, err) + require.True(t, foundRemaining, "unremoved post must be discovered") + assert.Equal(t, RematerializeDone, row.State) + assert.Contains(t, source.deletedURIs, remaining.URI) + assert.Contains(t, index.tombstonedURIs, remaining.URI) + assert.Equal(t, []string{"at://" + remainingAuthorDID + "/" + PostV2Collection + "/" + RematerializeRkey(remaining.URI)}, writer.postURIs) + _, stillPresent, err = source.ReadLegacyPost(ctx, remaining.URI) + require.NoError(t, err) + assert.False(t, stillPresent, "migrated post must no longer be in the legacy source") + + assert.Equal(t, 1, report.SkippedRemoved, "census and source passes must count the removed URI only once") + assert.Equal(t, 1, report.RemainingLegacy) + assert.False(t, report.ScopeComplete) + assert.False(t, report.Complete) + skips := skipEvents(progress, removed.URI) + require.Len(t, skips, 1, "the census reports the skip; the source pass must not report it again") + assert.NotEmpty(t, skips[0].Note, "the skipped post must be reported with a note") +} + +// skipEvents returns the progress events that report uri as skipped for a +// removal. +func skipEvents(progress []RematerializeProgress, uri string) []RematerializeProgress { + var events []RematerializeProgress + for _, event := range progress { + if event.OldURI == uri && event.To == RematerializeSkippedRemoved { + events = append(events, event) + } + } + return events +} + +type noRematerializeRemovals struct{} + +func (noRematerializeRemovals) ActiveRemovalsByURIs(context.Context, []string) (map[string][]RemovalSource, error) { + return nil, nil +} + +type refusalSource struct { + *memSource + calls []string +} + +func (s *refusalSource) ListLegacyPosts(ctx context.Context) ([]LegacyPost, error) { + s.calls = append(s.calls, "ListLegacyPosts") + return s.memSource.ListLegacyPosts(ctx) +} + +func (s *refusalSource) ReadLegacyPost(ctx context.Context, uri string) (LegacyPost, bool, error) { + s.calls = append(s.calls, "ReadLegacyPost") + return s.memSource.ReadLegacyPost(ctx, uri) +} + +func (s *refusalSource) DeleteLegacyPost(ctx context.Context, legacy LegacyPost, swapCID string) error { + s.calls = append(s.calls, "DeleteLegacyPost") + return s.memSource.DeleteLegacyPost(ctx, legacy, swapCID) +} + +// callRecordingLedger records every ledger call, so a test can assert which +// ones a skipped or refused post reached. +type callRecordingLedger struct { + RematerializeLedger + calls []string +} + +func (l *callRecordingLedger) Discover(ctx context.Context, oldURI, communityDID, authorDID string) (RematerializeLedgerRow, error) { + l.calls = append(l.calls, "Discover") + return l.RematerializeLedger.Discover(ctx, oldURI, communityDID, authorDID) +} + +func (l *callRecordingLedger) Get(ctx context.Context, oldURI string) (RematerializeLedgerRow, bool, error) { + l.calls = append(l.calls, "Get") + return l.RematerializeLedger.Get(ctx, oldURI) +} + +func (l *callRecordingLedger) ListResumable(ctx context.Context, communityDID string) ([]RematerializeLedgerRow, error) { + l.calls = append(l.calls, "ListResumable") + return l.RematerializeLedger.ListResumable(ctx, communityDID) +} + +func (l *callRecordingLedger) RecordPostV2Written(ctx context.Context, oldURI, sourceCID, newURI, newCID, newRkey string) error { + l.calls = append(l.calls, "RecordPostV2Written") + return l.RematerializeLedger.RecordPostV2Written(ctx, oldURI, sourceCID, newURI, newCID, newRkey) +} + +func (l *callRecordingLedger) MarkVerified(ctx context.Context, oldURI string) error { + l.calls = append(l.calls, "MarkVerified") + return l.RematerializeLedger.MarkVerified(ctx, oldURI) +} + +func (l *callRecordingLedger) MarkMigrated(ctx context.Context, oldURI string) error { + l.calls = append(l.calls, "MarkMigrated") + return l.RematerializeLedger.MarkMigrated(ctx, oldURI) +} + +func (l *callRecordingLedger) MarkDone(ctx context.Context, oldURI string) error { + l.calls = append(l.calls, "MarkDone") + return l.RematerializeLedger.MarkDone(ctx, oldURI) +} + +func (l *callRecordingLedger) MarkFallback(ctx context.Context, oldURI string, state RematerializeState, reason string) error { + l.calls = append(l.calls, "MarkFallback") + return l.RematerializeLedger.MarkFallback(ctx, oldURI, state, reason) +} + +func (l *callRecordingLedger) ReopenFallback(ctx context.Context, communityDID string) (int, error) { + l.calls = append(l.calls, "ReopenFallback") + return l.RematerializeLedger.ReopenFallback(ctx, communityDID) +} + +func (l *callRecordingLedger) CountByState(ctx context.Context, communityDID string) (map[RematerializeState]int, error) { + l.calls = append(l.calls, "CountByState") + return l.RematerializeLedger.CountByState(ctx, communityDID) +} + +func TestRematerialize_RequiresRemovalSeamBeforeAnyWork(t *testing.T) { + for _, testCase := range []struct { + name string + run bool + missingLookup bool + }{ + {name: "Run/nil lookup", run: true, missingLookup: true}, + {name: "Run/empty instance DID", run: true}, + {name: "RematerializeOne/nil lookup", missingLookup: true}, + {name: "RematerializeOne/empty instance DID"}, + } { + t.Run(testCase.name, func(t *testing.T) { + ctx := context.Background() + communityDID := "did:plc:community2222222222222222" + authorDID := "did:plc:author11111111111111111" + legacy := LegacyPost{ + URI: "at://" + communityDID + "/" + LegacyPostCollection + "/3krefusal", + CID: "bafylegacyrefusal", + CommunityDID: communityDID, + AuthorDID: authorDID, + RawRecord: map[string]any{ + "$type": LegacyPostCollection, "community": communityDID, + "author": authorDID, "title": "a legacy post", "createdAt": "2026-01-02T03:04:05Z", + }, + } + source := &refusalSource{memSource: &memSource{posts: []LegacyPost{legacy}}} + ledger := &callRecordingLedger{RematerializeLedger: newMemLedger()} + authorRepo := newMemAuthorRepo(authorDID) + communityRepo := &memCommunityRepo{did: communityDID, records: map[string]*pds.RecordResponse{}} + writer := &removedContractAcceptanceWriter{memAcceptanceWriter: &memAcceptanceWriter{repo: communityRepo}} + index := &removedContractIndex{} + var authorResolutions int + tool := &Rematerializer{ + Source: source, Ledger: ledger, + AuthorRepos: func(context.Context, string, *oauth.ClientSessionData) (AuthorRepo, error) { + authorResolutions++ + return authorRepo, nil + }, + Acceptances: writer, + CommunityRepos: func(context.Context, string) (CommunityRepo, error) { + return communityRepo, nil + }, + Index: index, + Removals: noRematerializeRemovals{}, + InstanceDID: "did:web:coves-instance.invalid", + } + missingField := "InstanceDID" + if testCase.missingLookup { + tool.Removals = nil + missingField = "Removals" + } else { + tool.InstanceDID = "" + } + + var err error + if testCase.run { + _, err = tool.Run(ctx) + } else { + _, err = tool.RematerializeOne(ctx, legacy) + } + require.ErrorContains(t, err, missingField, "a missing removal lookup or instance DID must fail, naming it, before touching any seam") + assert.Empty(t, source.calls, "no legacy source operation may run") + assert.Empty(t, ledger.calls, "no ledger operation may run") + assert.Zero(t, authorResolutions, "no author repo may be resolved") + assert.Zero(t, authorRepo.puts, "no postv2 may be written") + assert.Empty(t, writer.postURIs, "no acceptance may be written") + assert.Empty(t, index.tombstonedURIs, "no index row may be tombstoned") + }) + } +} + +func TestRematerializeOne_OnlyInstanceRemovalSkipsBeforeLedgerMutation(t *testing.T) { + instanceDID := "did:web:coves-instance.invalid" + otherAuthorityDID := "did:web:other-instance.invalid" + for _, testCase := range []struct { + name string + removals []RemovalSource + seededRow bool + wantSkipped bool + }{ + { + name: "instance removal with no ledger row", + removals: []RemovalSource{{AuthorityDID: instanceDID, ScopeKind: "instance"}}, + wantSkipped: true, + }, + { + name: "instance removal among other authorities", + removals: []RemovalSource{ + {AuthorityDID: otherAuthorityDID, ScopeKind: "instance"}, + {AuthorityDID: instanceDID, ScopeKind: "instance"}, + }, + wantSkipped: true, + }, + { + name: "instance removal after postv2 was written", + removals: []RemovalSource{{AuthorityDID: instanceDID, ScopeKind: "instance"}}, + seededRow: true, + wantSkipped: true, + }, + { + name: "another authority instance removal migrates", + removals: []RemovalSource{{AuthorityDID: otherAuthorityDID, ScopeKind: "instance"}}, + }, + { + name: "own authority community removal migrates", + removals: []RemovalSource{{AuthorityDID: instanceDID, ScopeKind: "community"}}, + }, + {name: "no removal migrates"}, + } { + t.Run(testCase.name, func(t *testing.T) { + ctx := context.Background() + communityDID := "did:plc:community2222222222222222" + authorDID := "did:plc:author11111111111111111" + legacy := LegacyPost{ + URI: "at://" + communityDID + "/" + LegacyPostCollection + "/3kdirect", + CID: "bafylegacydirect", + CommunityDID: communityDID, + AuthorDID: authorDID, + RawRecord: map[string]any{ + "$type": LegacyPostCollection, + "community": communityDID, + "author": authorDID, + "title": "a legacy post", + "createdAt": "2026-01-02T03:04:05Z", + }, + } + source := &removedContractSource{memSource: &memSource{posts: []LegacyPost{legacy}}} + baseLedger := newMemLedger() + authorRepo := newMemAuthorRepo(authorDID) + var seededRow RematerializeLedgerRow + if testCase.seededRow { + _, err := baseLedger.Discover(ctx, legacy.URI, communityDID, authorDID) + require.NoError(t, err) + rkey := RematerializeRkey(legacy.URI) + newURI := "at://" + authorDID + "/" + PostV2Collection + "/" + rkey + newCID := "bafypostv2direct" + require.NoError(t, baseLedger.RecordPostV2Written(ctx, legacy.URI, legacy.CID, newURI, newCID, rkey)) + var seededFound bool + seededRow, seededFound, err = baseLedger.Get(ctx, legacy.URI) + require.NoError(t, err) + require.True(t, seededFound) + body, err := postV2Body(legacy) + require.NoError(t, err) + authorRepo.records[PostV2Collection+"/"+rkey] = &pds.RecordResponse{URI: newURI, CID: newCID, Value: body} + } + ledger := &callRecordingLedger{RematerializeLedger: baseLedger} + communityRepo := &memCommunityRepo{did: communityDID, records: map[string]*pds.RecordResponse{}} + writer := &removedContractAcceptanceWriter{memAcceptanceWriter: &memAcceptanceWriter{repo: communityRepo}} + index := &removedContractIndex{} + var resolvedAuthorDIDs []string + var progress []RematerializeProgress + tool := &Rematerializer{ + Source: source, Ledger: ledger, + AuthorRepos: func(_ context.Context, did string, _ *oauth.ClientSessionData) (AuthorRepo, error) { + resolvedAuthorDIDs = append(resolvedAuthorDIDs, did) + return authorRepo, nil + }, + Acceptances: writer, + CommunityRepos: func(context.Context, string) (CommunityRepo, error) { + return communityRepo, nil + }, + Index: index, + Removals: &removedContractLookup{byURI: map[string][]RemovalSource{legacy.URI: testCase.removals}}, + InstanceDID: instanceDID, + Progress: func(event RematerializeProgress) { + progress = append(progress, event) + }, + } + + state, err := tool.RematerializeOne(ctx, legacy) + require.NoError(t, err) + if testCase.wantSkipped { + require.Equal(t, RematerializeSkippedRemoved, state, "an active instance removal must leave the legacy post untouched") + assert.NotContains(t, ledger.calls, "Discover", "skip must precede Discover, even for a resumable row") + for _, mutation := range []string{"RecordPostV2Written", "MarkVerified", "MarkMigrated", "MarkDone", "MarkFallback", "ReopenFallback"} { + assert.NotContains(t, ledger.calls, mutation, "skip must not advance the ledger") + } + row, found, err := baseLedger.Get(ctx, legacy.URI) + require.NoError(t, err) + if testCase.seededRow { + assert.True(t, found) + assert.Equal(t, seededRow, row, "the pre-existing ledger row must be byte-for-byte unchanged") + } else { + assert.False(t, found, "a skipped post must not acquire a ledger row") + } + assert.Empty(t, resolvedAuthorDIDs) + assert.Zero(t, authorRepo.puts) + assert.Empty(t, writer.postURIs) + assert.Empty(t, source.deletedURIs) + assert.Empty(t, index.tombstonedURIs) + skips := skipEvents(progress, legacy.URI) + require.Len(t, skips, 1, "the skipped post must be reported once") + if testCase.seededRow { + assert.Contains(t, skips[0].Note, "a postv2 already stands at "+seededRow.NewURI+" and is NOT removed", + "the skip must point the operator at the postv2 that escaped the removal") + } else { + assert.Equal(t, "active instance removal; leaving legacy post untouched", skips[0].Note) + } + return + } + + assert.Equal(t, RematerializeDone, state, "a non-matching removal must not prevent migration") + row, found, err := baseLedger.Get(ctx, legacy.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, row.State) + assert.Contains(t, ledger.calls, "Discover") + assert.Equal(t, []string{authorDID}, resolvedAuthorDIDs) + assert.Equal(t, 1, authorRepo.puts) + assert.Equal(t, []string{row.NewURI}, writer.postURIs) + assert.Equal(t, []string{legacy.URI}, source.deletedURIs) + assert.Equal(t, []string{legacy.URI}, index.tombstonedURIs) + }) + } +} + +type removedRunFixture struct { + removed LegacyPost + clean LegacyPost + source *removedContractSource + ledger *removedContractLedger + lookup *removedContractLookup + writer *removedContractAcceptanceWriter + index *removedContractIndex + authorRepos map[string]*memAuthorRepo + resolvedAuthorDIDs []string + unavailableAuthorDID string + tool *Rematerializer +} + +func newRemovedRunFixture() *removedRunFixture { + instanceDID := "did:web:coves-instance.invalid" + communityDID := "did:plc:community2222222222222222" + removedAuthorDID := "did:plc:removedauthor111111111111" + cleanAuthorDID := "did:plc:cleanauthor22222222222222" + post := func(authorDID, rkey string) LegacyPost { + return LegacyPost{ + URI: "at://" + communityDID + "/" + LegacyPostCollection + "/" + rkey, + CID: "bafylegacy" + rkey, + CommunityDID: communityDID, + AuthorDID: authorDID, + RawRecord: map[string]any{ + "$type": LegacyPostCollection, + "community": communityDID, + "author": authorDID, + "title": "a legacy post", + "createdAt": "2026-01-02T03:04:05Z", + }, + } + } + fixture := &removedRunFixture{ + removed: post(removedAuthorDID, "3kremovedrun"), + clean: post(cleanAuthorDID, "3kcleanrun"), + source: &removedContractSource{memSource: &memSource{}}, + ledger: &removedContractLedger{memLedger: newMemLedger()}, + index: &removedContractIndex{}, + authorRepos: map[string]*memAuthorRepo{ + removedAuthorDID: newMemAuthorRepo(removedAuthorDID), + cleanAuthorDID: newMemAuthorRepo(cleanAuthorDID), + }, + } + fixture.source.posts = []LegacyPost{fixture.removed, fixture.clean} + fixture.lookup = &removedContractLookup{byURI: map[string][]RemovalSource{ + fixture.removed.URI: {{AuthorityDID: instanceDID, ScopeKind: "instance"}}, + }} + communityRepo := &memCommunityRepo{did: communityDID, records: map[string]*pds.RecordResponse{}} + fixture.writer = &removedContractAcceptanceWriter{memAcceptanceWriter: &memAcceptanceWriter{repo: communityRepo}} + fixture.tool = &Rematerializer{ + Source: fixture.source, + Ledger: fixture.ledger, + AuthorRepos: func(_ context.Context, authorDID string, _ *oauth.ClientSessionData) (AuthorRepo, error) { + fixture.resolvedAuthorDIDs = append(fixture.resolvedAuthorDIDs, authorDID) + if authorDID == fixture.unavailableAuthorDID { + return nil, fmt.Errorf("no restorable credentials: %w", ErrNoAuthorCredentials) + } + return fixture.authorRepos[authorDID], nil + }, + Acceptances: fixture.writer, + CommunityRepos: func(context.Context, string) (CommunityRepo, error) { + return communityRepo, nil + }, + Index: fixture.index, + Removals: fixture.lookup, + InstanceDID: instanceDID, + } + return fixture +} + +// addPost lists one more legacy post, by its own author, in the fixture's +// community. +func (f *removedRunFixture) addPost(authorDID, rkey string) LegacyPost { + communityDID := f.removed.CommunityDID + post := LegacyPost{ + URI: "at://" + communityDID + "/" + LegacyPostCollection + "/" + rkey, + CID: "bafylegacy" + rkey, + CommunityDID: communityDID, + AuthorDID: authorDID, + RawRecord: map[string]any{ + "$type": LegacyPostCollection, + "community": communityDID, + "author": authorDID, + "title": "a legacy post", + "createdAt": "2026-01-02T03:04:05Z", + }, + } + f.authorRepos[authorDID] = newMemAuthorRepo(authorDID) + f.source.posts = append(f.source.posts, post) + return post +} + +func TestRematerializeRun_SkippedRemovedResetsOnTheSameRematerializer(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + + firstReport, err := fixture.tool.Run(ctx) + require.NoError(t, err) + require.Equal(t, 1, firstReport.SkippedRemoved, "the first run must count P1 once across its census and source passes") + _, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + assert.False(t, found, "the skipped post must have no ledger row") + cleanRow, found, err := fixture.ledger.Get(ctx, fixture.clean.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, cleanRow.State) + assert.Equal(t, 1, firstReport.RemainingLegacy) + + delete(fixture.lookup.byURI, fixture.removed.URI) + secondReport, err := fixture.tool.Run(ctx) + require.NoError(t, err) + assert.Equal(t, 0, secondReport.SkippedRemoved, "the count must be reset for each Run on the same Rematerializer") + removedRow, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + require.True(t, found, "a restored legacy post must enter the ledger on the second run") + assert.Equal(t, RematerializeDone, removedRow.State) + assert.Contains(t, fixture.source.deletedURIs, fixture.removed.URI) +} + +func TestRematerializeRun_RemovedAuthorWithoutCredentialsDoesNotTriggerAbort(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + fixture.unavailableAuthorDID = fixture.removed.AuthorDID + fixture.tool.AbortOnFallback = true + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err, "a removed author's missing credentials must not abort the run") + assert.Equal(t, 1, report.SkippedRemoved) + _, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + assert.False(t, found, "a removed post must not acquire even a fallback row") + assert.NotContains(t, fixture.ledger.discoveredURIs, fixture.removed.URI) + assert.NotContains(t, fixture.resolvedAuthorDIDs, fixture.removed.AuthorDID) + cleanRow, found, err := fixture.ledger.Get(ctx, fixture.clean.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, cleanRow.State) + assert.Contains(t, fixture.source.deletedURIs, fixture.clean.URI) + assert.NotContains(t, fixture.source.deletedURIs, fixture.removed.URI) +} + +func TestRematerializeRun_RemovedPostWithFallbackRowStillCountsAsRemaining(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + fixture.unavailableAuthorDID = fixture.removed.AuthorDID + seededRow := RematerializeLedgerRow{ + OldURI: fixture.removed.URI, CommunityDID: fixture.removed.CommunityDID, + AuthorDID: fixture.removed.AuthorDID, State: RematerializeFallbackLeftLegacy, + Reason: "credentials were unavailable on an earlier run", + } + fixture.ledger.rows[fixture.removed.URI] = seededRow + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + require.Equal(t, 1, report.SkippedRemoved, "a removed post with an existing fallback row still counts as skipped") + assert.NotContains(t, fixture.ledger.discoveredURIs, fixture.removed.URI, "census must skip before Discover even if a fallback row exists") + assert.NotContains(t, fixture.resolvedAuthorDIDs, fixture.removed.AuthorDID) + row, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, seededRow, row, "the terminal fallback row and reason must remain unchanged") + cleanRow, found, err := fixture.ledger.Get(ctx, fixture.clean.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, cleanRow.State) + assert.Equal(t, 1, report.RemainingLegacy, "a skipped post stays legacy even when an older fallback row accounts for its URI") + assert.False(t, report.Complete) + assert.NotContains(t, fixture.source.deletedURIs, fixture.removed.URI) +} + +type removedAfterDiscoveryLookup struct { + *removedContractLookup + ledger *removedContractLedger + targetURI string + checkedBeforeRow bool + checkedAfterRow bool +} + +func (l *removedAfterDiscoveryLookup) ActiveRemovalsByURIs(ctx context.Context, uris []string) (map[string][]RemovalSource, error) { + result, err := l.removedContractLookup.ActiveRemovalsByURIs(ctx, uris) + if err != nil { + return nil, err + } + for _, uri := range uris { + if uri != l.targetURI { + continue + } + _, discovered, err := l.ledger.Get(ctx, uri) + if err != nil { + return nil, err + } + if discovered { + l.checkedAfterRow = true + result[uri] = []RemovalSource{{AuthorityDID: "did:web:coves-instance.invalid", ScopeKind: "instance"}} + } else { + l.checkedBeforeRow = true + } + } + return result, nil +} + +func TestRematerializeRun_RemovalAfterCensusDiscoverySkipsSourcePass(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + const blobCID = "bafkreiembeddedblobcid" + fixture.clean.RawRecord["embed"] = map[string]any{ + "$type": "social.coves.embed.images", + "images": []any{map[string]any{ + "alt": "an image", + "image": map[string]any{"$type": "blob", "ref": map[string]any{"$link": blobCID}, "mimeType": "image/png"}, + }}, + } + fixture.source.posts = []LegacyPost{fixture.clean} + blobs := &countingBlobClient{bytesFor: map[string][]byte{blobCID: []byte("PNGDATA")}} + fixture.tool.Blobs = blobs + lookup := &removedAfterDiscoveryLookup{ + removedContractLookup: &removedContractLookup{byURI: map[string][]RemovalSource{}}, + ledger: fixture.ledger, targetURI: fixture.clean.URI, + } + fixture.tool.Removals = lookup + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + assert.True(t, lookup.checkedBeforeRow, "the census must see the post before discovery") + assert.True(t, lookup.checkedAfterRow, "the source pass must recheck the removal after discovery") + row, found, err := fixture.ledger.Get(ctx, fixture.clean.URI) + require.NoError(t, err) + require.True(t, found, "the clean census post must have been discovered") + assert.Equal(t, RematerializeDiscovered, row.State) + assert.Empty(t, row.SourceCID) + assert.Empty(t, row.NewURI) + assert.Equal(t, []string{fixture.clean.URI}, fixture.ledger.discoveredURIs, "the source pass must not call Discover again") + assert.Equal(t, 1, fixture.ledger.writes, "only the census discovery may mutate the ledger") + assert.Zero(t, fixture.authorRepos[fixture.clean.AuthorDID].puts) + assert.Zero(t, fixture.authorRepos[fixture.clean.AuthorDID].uploads) + assert.Zero(t, blobs.fetches) + assert.Empty(t, fixture.writer.postURIs) + assert.Empty(t, fixture.source.deletedURIs) + assert.Empty(t, fixture.index.tombstonedURIs) + assert.Equal(t, 1, report.SkippedRemoved) + assert.Equal(t, 1, report.RemainingLegacy) + assert.False(t, report.ScopeComplete) + assert.False(t, report.Complete) +} + +// seedResumableRemovedRow leaves the fixture's removed post as an earlier run +// would have: a ledger row at state, the postv2 standing in the author's repo, +// and, past postv2_written, the community's acceptance. +func seedResumableRemovedRow(t *testing.T, fixture *removedRunFixture, state RematerializeState) RematerializeLedgerRow { + t.Helper() + rkey := RematerializeRkey(fixture.removed.URI) + newURI := "at://" + fixture.removed.AuthorDID + "/" + PostV2Collection + "/" + rkey + newCID := "bafypostv2resumed" + seededRow := RematerializeLedgerRow{ + OldURI: fixture.removed.URI, CommunityDID: fixture.removed.CommunityDID, + AuthorDID: fixture.removed.AuthorDID, State: state, + SourceCID: fixture.removed.CID, NewURI: newURI, NewCID: newCID, NewRkey: rkey, + CreatedAt: time.Date(2026, time.January, 2, 3, 4, 5, 0, time.UTC), + UpdatedAt: time.Date(2026, time.January, 3, 4, 5, 6, 0, time.UTC), + } + fixture.ledger.rows[fixture.removed.URI] = seededRow + body, err := postV2Body(fixture.removed) + require.NoError(t, err) + fixture.authorRepos[fixture.removed.AuthorDID].records[PostV2Collection+"/"+rkey] = &pds.RecordResponse{ + URI: newURI, CID: newCID, Value: body, + } + if state != RematerializePostV2Written { + fixture.writer.repo.records[AcceptanceCollection+"/"+SubjectRkey(newURI)] = &pds.RecordResponse{ + CID: "bafyacceptresumed", + Value: map[string]any{"subject": map[string]any{"uri": newURI, "cid": newCID}}, + } + } + return seededRow +} + +func TestRematerializeRun_RemovedResumableRowsStayUnchanged(t *testing.T) { + for _, state := range []RematerializeState{RematerializePostV2Written, RematerializeVerified, RematerializeMigrated} { + for _, listed := range []bool{true, false} { + name := string(state) + "/ledger only" + if listed { + name = string(state) + "/source and ledger" + } + t.Run(name, func(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + if listed { + fixture.source.posts = []LegacyPost{fixture.removed} + } else { + fixture.source.posts = nil + } + seededRow := seedResumableRemovedRow(t, fixture, state) + newURI := seededRow.NewURI + ledgerSpy := &callRecordingLedger{RematerializeLedger: fixture.ledger} + fixture.tool.Ledger = ledgerSpy + var progress []RematerializeProgress + fixture.tool.Progress = func(event RematerializeProgress) { progress = append(progress, event) } + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + row, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, seededRow, row, "a removed resumable row must retain all fields") + assert.Contains(t, ledgerSpy.calls, "ListResumable", "the seeded row must reach the reconcile pass") + for _, mutation := range []string{"Discover", "RecordPostV2Written", "MarkVerified", "MarkMigrated", "MarkDone", "MarkFallback"} { + assert.NotContains(t, ledgerSpy.calls, mutation) + } + assert.Zero(t, fixture.ledger.writes) + assert.Empty(t, fixture.resolvedAuthorDIDs) + assert.Zero(t, fixture.authorRepos[fixture.removed.AuthorDID].puts) + assert.Zero(t, fixture.authorRepos[fixture.removed.AuthorDID].uploads) + assert.Empty(t, fixture.writer.postURIs) + assert.Empty(t, fixture.source.deletedURIs) + assert.Empty(t, fixture.index.tombstonedURIs) + skips := skipEvents(progress, fixture.removed.URI) + namingPostV2 := 0 + for _, event := range skips { + assert.NotEqual(t, "reconciled from the ledger", event.Note, + "Run must not re-report RematerializeOne's skip as a %s → %s reconcile", event.From, event.To) + if strings.Contains(event.Note, "a postv2 already stands at "+newURI+" and is NOT removed") { + namingPostV2++ + } + } + assert.Equal(t, 1, namingPostV2, + "the ledger reconcile must reach the removed row once, and its skip must name the postv2 that escaped the removal") + assert.Equal(t, 1, report.SkippedRemoved, "census, source and reconcile sightings count once") + assert.Equal(t, map[RematerializeState]int{state: 1}, report.ByState) + remaining := 0 + if listed { + remaining = 1 + } + assert.Equal(t, remaining, report.RemainingLegacy) + assert.False(t, report.ScopeComplete) + assert.False(t, report.Complete) + }) + } + } +} + +func TestRematerializeRun_ReopenedFallbackRemovedBeforeRetry(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + fixture.source.posts = []LegacyPost{fixture.removed} + // ReopenFallback leaves the row at discovered and clears the fallback reason. + seededRow := RematerializeLedgerRow{ + OldURI: fixture.removed.URI, CommunityDID: fixture.removed.CommunityDID, + AuthorDID: fixture.removed.AuthorDID, State: RematerializeDiscovered, + CreatedAt: time.Date(2026, time.January, 2, 3, 4, 5, 0, time.UTC), + UpdatedAt: time.Date(2026, time.January, 3, 4, 5, 6, 0, time.UTC), + } + fixture.ledger.rows[fixture.removed.URI] = seededRow + ledgerSpy := &callRecordingLedger{RematerializeLedger: fixture.ledger} + fixture.tool.Ledger = ledgerSpy + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + row, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, seededRow, row, "the reopened row must stay discovered") + assert.NotContains(t, ledgerSpy.calls, "Discover") + for _, mutation := range []string{"RecordPostV2Written", "MarkVerified", "MarkMigrated", "MarkDone", "MarkFallback"} { + assert.NotContains(t, ledgerSpy.calls, mutation) + } + assert.Zero(t, fixture.ledger.writes) + assert.Empty(t, fixture.resolvedAuthorDIDs) + assert.Zero(t, fixture.authorRepos[fixture.removed.AuthorDID].puts) + assert.Zero(t, fixture.authorRepos[fixture.removed.AuthorDID].uploads) + assert.Empty(t, fixture.writer.postURIs) + assert.Empty(t, fixture.source.deletedURIs) + assert.Empty(t, fixture.index.tombstonedURIs) + assert.Equal(t, 1, report.SkippedRemoved) + assert.Equal(t, 1, report.RemainingLegacy) + assert.False(t, report.ScopeComplete) + assert.False(t, report.Complete) +} + +type conditionalRemovalFailureLookup struct { + failure error + failURI string + afterDoneURI string + ledger *removedContractLedger + askedURIs []string + checkedBefore bool + failedAfter bool +} + +func (l *conditionalRemovalFailureLookup) ActiveRemovalsByURIs(ctx context.Context, uris []string) (map[string][]RemovalSource, error) { + for _, uri := range uris { + l.askedURIs = append(l.askedURIs, uri) + if l.failURI != "" && uri != l.failURI { + continue + } + if l.afterDoneURI != "" { + row, found, err := l.ledger.Get(ctx, l.afterDoneURI) + if err != nil { + return nil, err + } + if !found || row.State != RematerializeDone { + l.checkedBefore = true + continue + } + l.failedAfter = true + } + return nil, l.failure + } + return nil, nil +} + +type censusSnapshotLedger struct { + *callRecordingLedger + targetURI string + row RematerializeLedgerRow + seen bool +} + +func (l *censusSnapshotLedger) Discover(ctx context.Context, oldURI, communityDID, authorDID string) (RematerializeLedgerRow, error) { + row, err := l.callRecordingLedger.Discover(ctx, oldURI, communityDID, authorDID) + if err == nil && oldURI == l.targetURI { + l.row = row + l.seen = true + } + return row, err +} + +func TestRematerializeRun_RemovalLookupErrorOnFirstRecordStopsBeforeDiscovery(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + errLookupUnavailable := errors.New("removal lookup unavailable") + lookup := &conditionalRemovalFailureLookup{failure: errLookupUnavailable} + fixture.tool.Removals = lookup + ledgerSpy := &callRecordingLedger{RematerializeLedger: fixture.ledger} + fixture.tool.Ledger = ledgerSpy + + _, err := fixture.tool.Run(ctx) + require.ErrorIs(t, err, errLookupUnavailable) + assert.Equal(t, []string{fixture.removed.URI}, lookup.askedURIs, "the first lookup error must halt enumeration") + assert.NotContains(t, ledgerSpy.calls, "Discover") + assert.Zero(t, fixture.ledger.writes) + for _, legacy := range []LegacyPost{fixture.removed, fixture.clean} { + _, found, getErr := fixture.ledger.Get(ctx, legacy.URI) + require.NoError(t, getErr) + assert.False(t, found, "the failed run must not create a ledger row for %s", legacy.URI) + assert.Zero(t, fixture.authorRepos[legacy.AuthorDID].puts) + assert.Zero(t, fixture.authorRepos[legacy.AuthorDID].uploads) + } + assert.Empty(t, fixture.resolvedAuthorDIDs) + assert.Empty(t, fixture.writer.postURIs) + assert.Empty(t, fixture.source.deletedURIs) + assert.Empty(t, fixture.index.tombstonedURIs) +} + +func TestRematerializeRun_RemovalLookupErrorAfterCompletedPostKeepsProgress(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + first := fixture.clean + second := fixture.removed + errLookupUnavailable := errors.New("removal lookup unavailable") + lookup := &conditionalRemovalFailureLookup{ + failure: errLookupUnavailable, failURI: second.URI, + afterDoneURI: first.URI, ledger: fixture.ledger, + } + fixture.tool.Removals = lookup + ledgerSpy := &censusSnapshotLedger{ + callRecordingLedger: &callRecordingLedger{RematerializeLedger: fixture.ledger}, + targetURI: second.URI, + } + fixture.tool.Ledger = ledgerSpy + const blobCID = "bafkreierrorpathblob" + second.RawRecord["embed"] = map[string]any{ + "$type": "social.coves.embed.images", + "images": []any{map[string]any{ + "alt": "an image", + "image": map[string]any{"$type": "blob", "ref": map[string]any{"$link": blobCID}, "mimeType": "image/png"}, + }}, + } + fixture.source.posts = []LegacyPost{first, second} + blobs := &countingBlobClient{bytesFor: map[string][]byte{blobCID: []byte("PNGDATA")}} + fixture.tool.Blobs = blobs + + _, err := fixture.tool.Run(ctx) + require.ErrorIs(t, err, errLookupUnavailable) + assert.True(t, lookup.checkedBefore, "P2 must have passed the census lookup while P1 was not done") + assert.True(t, lookup.failedAfter, "P2 must fail its source-pass lookup only after P1 reached done") + require.True(t, ledgerSpy.seen, "the census must discover P2 before its later lookup fails") + firstRow, found, err := fixture.ledger.Get(ctx, first.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, firstRow.State) + assert.Equal(t, []string{firstRow.NewURI}, fixture.writer.postURIs, "P1 acceptance must persist") + assert.Equal(t, []string{first.URI}, fixture.source.deletedURIs, "P1 delete must persist") + assert.Equal(t, []string{first.URI}, fixture.index.tombstonedURIs, "P1 tombstone must persist") + assert.Equal(t, 1, fixture.authorRepos[first.AuthorDID].puts) + secondRow, found, err := fixture.ledger.Get(ctx, second.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, ledgerSpy.row, secondRow, "P2's census row must be unchanged in every field") + assert.Equal(t, RematerializeDiscovered, secondRow.State) + secondDiscoveries := 0 + for _, uri := range fixture.ledger.discoveredURIs { + if uri == second.URI { + secondDiscoveries++ + } + } + assert.Equal(t, 1, secondDiscoveries, "only the census may Discover P2; the source pass must not Discover it again") + assert.Zero(t, fixture.authorRepos[second.AuthorDID].puts) + assert.Zero(t, fixture.authorRepos[second.AuthorDID].uploads) + assert.Zero(t, blobs.fetches) + assert.NotContains(t, fixture.source.deletedURIs, second.URI) + assert.NotContains(t, fixture.index.tombstonedURIs, second.URI) +} + +func TestRematerializeRun_RemovalLookupErrorOnLedgerOnlyRowLeavesItUntouched(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + fixture.source.posts = nil + errLookupUnavailable := errors.New("removal lookup unavailable") + lookup := &conditionalRemovalFailureLookup{failure: errLookupUnavailable, failURI: fixture.removed.URI} + fixture.tool.Removals = lookup + rkey := RematerializeRkey(fixture.removed.URI) + newURI := "at://" + fixture.removed.AuthorDID + "/" + PostV2Collection + "/" + rkey + newCID := "bafypostv2ledgeronly" + seededRow := RematerializeLedgerRow{ + OldURI: fixture.removed.URI, CommunityDID: fixture.removed.CommunityDID, + AuthorDID: fixture.removed.AuthorDID, State: RematerializeMigrated, + SourceCID: fixture.removed.CID, NewURI: newURI, NewCID: newCID, NewRkey: rkey, + CreatedAt: time.Date(2026, time.January, 2, 3, 4, 5, 0, time.UTC), + UpdatedAt: time.Date(2026, time.January, 3, 4, 5, 6, 0, time.UTC), + } + fixture.ledger.rows[fixture.removed.URI] = seededRow + body, err := postV2Body(fixture.removed) + require.NoError(t, err) + fixture.authorRepos[fixture.removed.AuthorDID].records[PostV2Collection+"/"+rkey] = &pds.RecordResponse{ + URI: newURI, CID: newCID, Value: body, + } + fixture.writer.repo.records[AcceptanceCollection+"/"+SubjectRkey(newURI)] = &pds.RecordResponse{ + CID: "bafyacceptledgeronly", + Value: map[string]any{"subject": map[string]any{"uri": newURI, "cid": newCID}}, + } + ledgerSpy := &callRecordingLedger{RematerializeLedger: fixture.ledger} + fixture.tool.Ledger = ledgerSpy + + _, err = fixture.tool.Run(ctx) + require.ErrorIs(t, err, errLookupUnavailable) + assert.Contains(t, ledgerSpy.calls, "ListResumable") + assert.Equal(t, []string{fixture.removed.URI}, lookup.askedURIs, "only the reconcile pass may reach the missing legacy record") + row, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, seededRow, row, "the failed lookup must not change any ledger field") + for _, mutation := range []string{"Discover", "RecordPostV2Written", "MarkVerified", "MarkMigrated", "MarkDone", "MarkFallback"} { + assert.NotContains(t, ledgerSpy.calls, mutation) + } + assert.Zero(t, fixture.ledger.writes) + assert.Empty(t, fixture.resolvedAuthorDIDs) + assert.Zero(t, fixture.authorRepos[fixture.removed.AuthorDID].puts) + assert.Empty(t, fixture.writer.postURIs) + assert.Empty(t, fixture.source.deletedURIs) + assert.Empty(t, fixture.index.tombstonedURIs) +} + +func TestDryRun_SkipsInstanceRemovedLegacyPostLikeTheRealRun(t *testing.T) { + ctx := context.Background() + realFixture := newRemovedRunFixture() + dryFixture := newRemovedRunFixture() + require.Equal(t, realFixture.removed, dryFixture.removed) + require.Equal(t, realFixture.clean, dryFixture.clean) + var dryProgress []RematerializeProgress + dryFixture.tool.Progress = func(event RematerializeProgress) { + dryProgress = append(dryProgress, event) + } + + realReport, err := realFixture.tool.Run(ctx) + require.NoError(t, err) + dry := DryRunOf(dryFixture.tool) + dryReport, err := dry.Run(ctx) + require.NoError(t, err) + + require.Equal(t, 1, realReport.SkippedRemoved) + assert.Equal(t, realReport.SkippedRemoved, dryReport.SkippedRemoved) + assert.Equal(t, realReport.RemainingLegacy, dryReport.RemainingLegacy) + assert.Equal(t, realReport.Done, dryReport.Done) + assert.Equal(t, realReport.Discovered, dryReport.Discovered) + assert.Equal(t, realReport.ScopeComplete, dryReport.ScopeComplete) + assert.Equal(t, realReport.Complete, dryReport.Complete) + assert.Equal(t, 1, dryReport.RemainingLegacy, "P1 must remain in the source after the simulated P2 delete") + assert.Equal(t, 1, dryReport.Done) + assert.Equal(t, 1, dryReport.Discovered) + assert.False(t, dryReport.ScopeComplete) + assert.False(t, dryReport.Complete) + deletes, isDry := DryRunDeletes(dry) + require.True(t, isDry) + assert.Equal(t, 1, deletes, "only clean P2 should be scheduled for deletion") + tombstones, isDry := DryRunTombstones(dry) + require.True(t, isDry) + assert.Equal(t, 1, tombstones, "only clean P2 should be scheduled for a tombstone") + + assert.Empty(t, dryFixture.ledger.rows, "the dry run must not persist even P2's ledger row") + assert.Zero(t, dryFixture.ledger.writes) + assert.Empty(t, dryFixture.ledger.discoveredURIs) + assert.Empty(t, dryFixture.writer.postURIs) + assert.Zero(t, dryFixture.writer.calls) + assert.Empty(t, dryFixture.writer.repo.records) + assert.Zero(t, dryFixture.writer.repo.puts) + assert.Empty(t, dryFixture.source.deletedURIs) + assert.Zero(t, dryFixture.source.deletes) + assert.Empty(t, dryFixture.index.tombstonedURIs) + for _, authorDID := range []string{dryFixture.removed.AuthorDID, dryFixture.clean.AuthorDID} { + repo := dryFixture.authorRepos[authorDID] + assert.Zero(t, repo.puts) + assert.Zero(t, repo.uploads) + assert.Zero(t, repo.deletes) + assert.Empty(t, repo.records) + } + assert.Equal(t, []string{dryFixture.clean.AuthorDID}, dryFixture.resolvedAuthorDIDs, + "P1 must be skipped before resolving A1, while clean P2 must resolve A2") + drySkips := skipEvents(dryProgress, dryFixture.removed.URI) + require.Len(t, drySkips, 1, "the dry run must report P1's skip once") + assert.NotEmpty(t, drySkips[0].Note, "the dry run must report P1's skip with a note") +} + +// scopedRemovedSource lists only its scope's posts, as the production source +// does for a -community run. +type scopedRemovedSource struct { + *removedContractSource + scope string +} + +func (s *scopedRemovedSource) ListLegacyPosts(ctx context.Context) ([]LegacyPost, error) { + all, err := s.removedContractSource.ListLegacyPosts(ctx) + if err != nil { + return nil, err + } + var inScope []LegacyPost + for _, post := range all { + if s.scope == "" || post.CommunityDID == s.scope { + inScope = append(inScope, post) + } + } + return inScope, nil +} + +// scopedRemovedLedger resumes and counts by community, as the migration-037 +// ledger does. +type scopedRemovedLedger struct { + *removedContractLedger +} + +func (l scopedRemovedLedger) ListResumable(_ context.Context, communityDID string) ([]RematerializeLedgerRow, error) { + var rows []RematerializeLedgerRow + for _, row := range l.rows { + if (communityDID == "" || row.CommunityDID == communityDID) && row.State != RematerializeDone && !IsFallback(row.State) { + rows = append(rows, row) + } + } + return rows, nil +} + +func (l scopedRemovedLedger) CountByState(_ context.Context, communityDID string) (map[RematerializeState]int, error) { + counts := map[RematerializeState]int{} + for _, row := range l.rows { + if communityDID == "" || row.CommunityDID == communityDID { + counts[row.State]++ + } + } + return counts, nil +} + +// A scoped run re-scans only its own community, so it cannot see a legacy post +// standing in another one: here, one a scoped run of that community skipped for +// an instance removal. Only an unscoped run may report the migration Complete. +func TestRematerializeRun_ScopedRunNeverReportsComplete(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + communityB := fixture.removed.CommunityDID + communityA := "did:plc:communityaaaaaaaaaaaaaaaaa" + inA := fixture.clean + inA.URI = "at://" + communityA + "/" + LegacyPostCollection + "/3kcleanina" + inA.CommunityDID = communityA + inA.RawRecord = map[string]any{ + "$type": LegacyPostCollection, + "community": communityA, + "author": inA.AuthorDID, + "title": "a legacy post", + "createdAt": "2026-01-02T03:04:05Z", + } + fixture.source.posts = []LegacyPost{fixture.removed, inA} + source := &scopedRemovedSource{removedContractSource: fixture.source} + fixture.tool.Source = source + fixture.tool.Ledger = scopedRemovedLedger{removedContractLedger: fixture.ledger} + + source.scope, fixture.tool.CommunityScope = communityB, communityB + reportB, err := fixture.tool.Run(ctx) + require.NoError(t, err) + require.Equal(t, 1, reportB.SkippedRemoved, "the B run must skip the removed post") + assert.Equal(t, 1, reportB.RemainingLegacy) + assert.False(t, reportB.Complete) + + source.scope, fixture.tool.CommunityScope = communityA, communityA + reportA, err := fixture.tool.Run(ctx) + require.NoError(t, err) + rowA, found, err := fixture.ledger.Get(ctx, inA.URI) + require.NoError(t, err) + require.True(t, found) + require.Equal(t, RematerializeDone, rowA.State) + assert.True(t, reportA.ScopeComplete, "the A run drained its own scope") + assert.Zero(t, reportA.RemainingLegacy, "the A run's re-scan cannot see community B") + assert.Equal(t, reportA.GlobalDiscovered, reportA.GlobalDone, "every row the ledger knows is done") + assert.False(t, reportA.Complete, + "a scoped run reported the whole migration complete while the removed legacy post still stands in another community") + + source.scope, fixture.tool.CommunityScope = "", "" + reportAll, err := fixture.tool.Run(ctx) + require.NoError(t, err) + assert.Equal(t, 1, reportAll.RemainingLegacy, "an unscoped run sees the removed post") + assert.False(t, reportAll.Complete) + + delete(fixture.lookup.byURI, fixture.removed.URI) + reportAll, err = fixture.tool.Run(ctx) + require.NoError(t, err) + assert.Zero(t, reportAll.RemainingLegacy) + assert.True(t, reportAll.Complete, "an unscoped run over a drained migration is the §11 step 6 gate") +} + +// removalByCheckLookup answers the n-th removal check of targetURI with +// removedOnCheck[n-1], and repeats the last answer after that. Every other URI +// has no removal. +type removalByCheckLookup struct { + targetURI string + instanceDID string + removedOnCheck []bool + checks int +} + +func (l *removalByCheckLookup) ActiveRemovalsByURIs(_ context.Context, uris []string) (map[string][]RemovalSource, error) { + result := make(map[string][]RemovalSource) + for _, uri := range uris { + if uri != l.targetURI { + continue + } + l.checks++ + if l.removedOnCheck[min(l.checks, len(l.removedOnCheck))-1] { + result[uri] = []RemovalSource{{AuthorityDID: l.instanceDID, ScopeKind: "instance"}} + } + } + return result, nil +} + +// The census verdict holds for the whole run. A post the census skipped is not +// retried by the source pass even if the removal is lifted in between: its +// author never went through the credential census. +func TestRematerializeRun_CensusSkipHoldsWhenTheRemovalIsLiftedBeforeTheSourcePass(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + lookup := &removalByCheckLookup{ + targetURI: fixture.removed.URI, instanceDID: fixture.tool.InstanceDID, + removedOnCheck: []bool{true, false}, + } + fixture.tool.Removals = lookup + var progress []RematerializeProgress + fixture.tool.Progress = func(event RematerializeProgress) { progress = append(progress, event) } + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + assert.Equal(t, 1, lookup.checks, "only the census may check a post the census skipped") + _, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + assert.False(t, found, "a post the census skipped must not acquire a ledger row later in the run") + assert.NotContains(t, fixture.ledger.discoveredURIs, fixture.removed.URI) + assert.NotContains(t, fixture.resolvedAuthorDIDs, fixture.removed.AuthorDID, + "the skipped post's author never went through the credential census") + assert.Zero(t, fixture.authorRepos[fixture.removed.AuthorDID].puts) + assert.Zero(t, fixture.authorRepos[fixture.removed.AuthorDID].uploads) + assert.NotContains(t, fixture.source.deletedURIs, fixture.removed.URI) + assert.NotContains(t, fixture.index.tombstonedURIs, fixture.removed.URI) + cleanRow, found, err := fixture.ledger.Get(ctx, fixture.clean.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, cleanRow.State) + assert.Equal(t, []string{cleanRow.NewURI}, fixture.writer.postURIs, "only the clean post may be accepted") + assert.Equal(t, 1, report.SkippedRemoved) + assert.Equal(t, 1, report.RemainingLegacy, "the skipped post still stands as legacy") + assert.False(t, report.ScopeComplete) + assert.False(t, report.Complete) + assert.Len(t, skipEvents(progress, fixture.removed.URI), 1, + "the census reports the skip; the source pass must not report it again") +} + +// The census verdict also holds for the ledger reconcile. A listed post whose +// row an earlier run left past discovered, skipped by the census, is not +// finished by the reconcile pass when the removal is lifted in between, and the +// census's own skip names the postv2 that already stands. +func TestRematerializeRun_CensusSkipHoldsForAResumableRowWhenTheRemovalIsLifted(t *testing.T) { + for _, state := range []RematerializeState{RematerializePostV2Written, RematerializeVerified, RematerializeMigrated} { + t.Run(string(state), func(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + fixture.source.posts = []LegacyPost{fixture.removed} + seededRow := seedResumableRemovedRow(t, fixture, state) + lookup := &removalByCheckLookup{ + targetURI: fixture.removed.URI, instanceDID: fixture.tool.InstanceDID, + removedOnCheck: []bool{true, false}, + } + fixture.tool.Removals = lookup + var progress []RematerializeProgress + fixture.tool.Progress = func(event RematerializeProgress) { progress = append(progress, event) } + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + assert.Equal(t, 1, lookup.checks, "only the census may check a post the census skipped") + row, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, seededRow, row, "the reconcile pass must not advance a row the census skipped") + assert.Empty(t, fixture.writer.postURIs) + assert.Empty(t, fixture.source.deletedURIs) + assert.Empty(t, fixture.index.tombstonedURIs) + assert.Equal(t, 1, report.SkippedRemoved) + assert.Equal(t, 1, report.RemainingLegacy) + skips := skipEvents(progress, fixture.removed.URI) + require.Len(t, skips, 1, "the census reports the skip once") + assert.Contains(t, skips[0].Note, "a postv2 already stands at "+seededRow.NewURI+" and is NOT removed") + }) + } +} + +// A removal that lands after the check at the top of RematerializeOne but +// before the acceptance is honoured: the postv2 already written stays +// unaccepted, the row stays at postv2_written, and nothing is deleted. +func TestRematerializeOne_RemovalBeforeTheAcceptanceLeavesThePostV2Unaccepted(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + legacy := fixture.clean + lookup := &removalByCheckLookup{ + targetURI: legacy.URI, instanceDID: fixture.tool.InstanceDID, + removedOnCheck: []bool{false, true}, + } + fixture.tool.Removals = lookup + var progress []RematerializeProgress + fixture.tool.Progress = func(event RematerializeProgress) { progress = append(progress, event) } + + state, err := fixture.tool.RematerializeOne(ctx, legacy) + require.NoError(t, err) + assert.Equal(t, RematerializeSkippedRemoved, state) + assert.Equal(t, 2, lookup.checks, "the removal must be checked again right before the acceptance") + newURI := "at://" + legacy.AuthorDID + "/" + PostV2Collection + "/" + RematerializeRkey(legacy.URI) + row, found, err := fixture.ledger.Get(ctx, legacy.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializePostV2Written, row.State, "the row must stay at postv2_written") + assert.Equal(t, newURI, row.NewURI) + assert.Equal(t, 1, fixture.authorRepos[legacy.AuthorDID].puts, "the postv2 written before the removal stays in the author's repo") + assert.Zero(t, fixture.writer.calls, "the acceptance must not be written") + assert.Empty(t, fixture.writer.postURIs) + assert.Empty(t, fixture.source.deletedURIs) + assert.Empty(t, fixture.index.tombstonedURIs) + skips := skipEvents(progress, legacy.URI) + require.Len(t, skips, 1) + assert.Contains(t, skips[0].Note, newURI, "the skip must name the unaccepted postv2") +} + +// The same late removal, seen by a whole run: the post is counted as skipped, +// stays legacy, and the reconcile pass leaves its postv2_written row alone. +func TestRematerializeRun_RemovalBeforeTheAcceptanceCountsAsSkippedAndRemaining(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + lookup := &removalByCheckLookup{ + targetURI: fixture.removed.URI, instanceDID: fixture.tool.InstanceDID, + // The census, then the top of RematerializeOne, then the recheck before + // the acceptance. + removedOnCheck: []bool{false, false, true}, + } + fixture.tool.Removals = lookup + var progress []RematerializeProgress + fixture.tool.Progress = func(event RematerializeProgress) { progress = append(progress, event) } + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + row, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializePostV2Written, row.State) + cleanRow, found, err := fixture.ledger.Get(ctx, fixture.clean.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, []string{cleanRow.NewURI}, fixture.writer.postURIs, "the removed post's postv2 must not be accepted") + assert.NotContains(t, fixture.source.deletedURIs, fixture.removed.URI) + assert.NotContains(t, fixture.index.tombstonedURIs, fixture.removed.URI) + assert.Equal(t, 1, report.SkippedRemoved) + assert.Equal(t, 1, report.RemainingLegacy) + assert.Equal(t, map[RematerializeState]int{RematerializePostV2Written: 1, RematerializeDone: 1}, report.ByState) + assert.False(t, report.ScopeComplete) + assert.False(t, report.Complete) + for _, event := range skipEvents(progress, fixture.removed.URI) { + assert.NotEqual(t, "reconciled from the ledger", event.Note) + } +} + +// Mutant M24: a census that skipped on ANY active removal, while +// RematerializeOne stayed strict, would bypass the credential preflight and +// count the post as skipped. A removal by another authority, or by the instance +// at community scope, must go through the census and migrate. +func TestRematerializeRun_NonMatchingRemovalGoesThroughTheCensusAndMigrates(t *testing.T) { + for _, testCase := range []struct { + name string + removal RemovalSource + }{ + {name: "instance authority at community scope", removal: RemovalSource{AuthorityDID: "did:web:coves-instance.invalid", ScopeKind: "community"}}, + {name: "another authority at instance scope", removal: RemovalSource{AuthorityDID: "did:web:other-instance.invalid", ScopeKind: "instance"}}, + } { + t.Run(testCase.name+"/migrates", func(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + fixture.lookup.byURI[fixture.removed.URI] = []RemovalSource{testCase.removal} + fixture.tool.AbortOnFallback = true + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + assert.Zero(t, report.SkippedRemoved) + row, found, err := fixture.ledger.Get(ctx, fixture.removed.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, row.State) + assert.Contains(t, fixture.source.deletedURIs, fixture.removed.URI) + assert.Zero(t, report.RemainingLegacy) + assert.True(t, report.Complete) + }) + t.Run(testCase.name+"/credential census still aborts", func(t *testing.T) { + ctx := context.Background() + fixture := newRemovedRunFixture() + fixture.lookup.byURI[fixture.removed.URI] = []RemovalSource{testCase.removal} + fixture.unavailableAuthorDID = fixture.removed.AuthorDID + fixture.tool.AbortOnFallback = true + + _, err := fixture.tool.Run(ctx) + require.ErrorContains(t, err, fixture.removed.AuthorDID, + "the credential census must reach a post whose removal does not match, and abort on its stranded author") + assert.Empty(t, fixture.writer.postURIs, "the abort must precede every repo mutation") + assert.Empty(t, fixture.source.deletedURIs) + }) + } +} + +// Moderation writes only removals by the instance DID at instance scope, so an +// active removal that matches neither means this tool's INSTANCE_DID differs +// from the server's. Such a post still migrates, but the run counts it. +func TestRematerializeRun_CountsActiveRemovalsThatDoNotMatchTheInstance(t *testing.T) { + ctx := context.Background() + otherAuthorityDID := "did:web:other-instance.invalid" + build := func() (*removedRunFixture, LegacyPost, LegacyPost) { + fixture := newRemovedRunFixture() + unmatched := fixture.addPost("did:plc:unmatchedauthor3333333333", "3kunmatched") + both := fixture.addPost("did:plc:bothauthor444444444444444", "3kboth") + fixture.lookup.byURI[unmatched.URI] = []RemovalSource{{AuthorityDID: otherAuthorityDID, ScopeKind: "instance"}} + fixture.lookup.byURI[both.URI] = []RemovalSource{ + {AuthorityDID: otherAuthorityDID, ScopeKind: "instance"}, + {AuthorityDID: fixture.tool.InstanceDID, ScopeKind: "instance"}, + } + return fixture, unmatched, both + } + fixture, unmatched, both := build() + var progress []RematerializeProgress + fixture.tool.Progress = func(event RematerializeProgress) { progress = append(progress, event) } + + report, err := fixture.tool.Run(ctx) + require.NoError(t, err) + assert.Equal(t, 1, report.UnmatchedRemovals, + "only the post whose every active removal misses the instance DID at instance scope is unmatched") + assert.Equal(t, 2, report.SkippedRemoved, "the matched post and the post with both kinds of removal are skipped") + unmatchedRow, found, err := fixture.ledger.Get(ctx, unmatched.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, unmatchedRow.State, "an unmatched removal is counted, not skipped") + _, found, err = fixture.ledger.Get(ctx, both.URI) + require.NoError(t, err) + assert.False(t, found, "a matching removal skips the post whatever else stands") + cleanRow, found, err := fixture.ledger.Get(ctx, fixture.clean.URI) + require.NoError(t, err) + require.True(t, found) + assert.Equal(t, RematerializeDone, cleanRow.State) + assert.Condition(t, func() bool { + for _, event := range progress { + if event.OldURI == unmatched.URI && strings.Contains(event.Note, "not by the instance DID at instance scope") { + return true + } + } + return false + }, "the census must report the unmatched removal") + + dryFixture, _, _ := build() + dryReport, err := DryRunOf(dryFixture.tool).Run(ctx) + require.NoError(t, err) + assert.Equal(t, 1, dryReport.UnmatchedRemovals, "the dry run must count the same unmatched removal") + assert.Equal(t, 2, dryReport.SkippedRemoved) +} diff --git a/internal/core/posts/rematerialize_test.go b/internal/core/posts/rematerialize_test.go index 4e49a5c..53e4e14 100644 --- a/internal/core/posts/rematerialize_test.go +++ b/internal/core/posts/rematerialize_test.go @@ -21,6 +21,12 @@ import ( "github.com/stretchr/testify/require" ) +type noRematerializeRemovalsIntegration struct{} + +func (noRematerializeRemovalsIntegration) ActiveRemovalsByURIs(context.Context, []string) (map[string][]posts.RemovalSource, error) { + return nil, nil +} + // callLog is an ordered record of the load-bearing calls a run makes, shared by // the fake factory, author repo, and legacy source, so a test can assert on // ORDERING that no outcome value reveals — specifically that the credential @@ -781,7 +787,7 @@ func TestRematerialize_HappyPath_WalksToDoneVerifyBeforeDelete(t *testing.T) { legacy := legacyPost(t, rematCommunityDID, rematAuthorDID) source := newFakeLegacySource(legacy) - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} state, err := tool.RematerializeOne(context.Background(), legacy) require.NoError(t, err) @@ -832,6 +838,7 @@ func TestRematerialize_TombstonesLegacyIndexAfterDeleteBeforeDone(t *testing.T) tool := &posts.Rematerializer{ Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Index: index, + Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid", } state, err := tool.RematerializeOne(context.Background(), legacy) @@ -866,7 +873,7 @@ func TestRematerialize_ReRun_IsAPureNoOp(t *testing.T) { writer := &spyAcceptanceWriter{} legacy := legacyPost(t, rematCommunityDID, rematAuthorDID) source := newFakeLegacySource(legacy) - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} first, err := tool.RematerializeOne(context.Background(), legacy) require.NoError(t, err) @@ -911,7 +918,7 @@ func TestRematerialize_ResumeAfterDeleteFailure_RetriesOnlyTheDelete(t *testing. legacy := legacyPost(t, rematCommunityDID, rematAuthorDID) source := newFakeLegacySource(legacy) source.deleteErr[legacy.URI] = fmt.Errorf("transient: the community PDS returned 502 on delete") - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} // First pass: everything succeeds up to the delete, which fails once. The row // must stop at migrated — the checkpoint BEFORE the delete — never done. @@ -953,7 +960,7 @@ func TestRematerialize_CIDMismatch_DoesNotCheckpointOrDelete(t *testing.T) { writer := &spyAcceptanceWriter{} legacy := legacyPost(t, rematCommunityDID, rematAuthorDID) source := newFakeLegacySource(legacy) - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} // The verify re-read of the postv2 comes back with a DIFFERENT CID than the // one the acceptance pinned — a concurrent edit landing in the write→verify @@ -1000,7 +1007,7 @@ func TestRematerialize_NoCredentials_LeavesLegacyNeverForges(t *testing.T) { authors.noCreds[humanDID] = true legacy := legacyPost(t, rematCommunityDID, humanDID) source := newFakeLegacySource(legacy) - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} state, err := tool.RematerializeOne(context.Background(), legacy) require.NoError(t, err, "a no-creds record is an expected terminal outcome, not a run-failing error") @@ -1036,7 +1043,7 @@ func TestRematerialize_Run_CensusGatesCompletionWhileFallbackSurvives(t *testing stranded := legacyPost(t, rematCommunityDID, humanDID) source := newFakeLegacySource(migratable, stranded) writer := &spyAcceptanceWriter{} - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} report, err := tool.Run(context.Background()) require.NoError(t, err) @@ -1072,7 +1079,7 @@ func TestRematerialize_UsesDirectAcceptanceWriter_NeverReDecides(t *testing.T) { // community it currently sits in. This mirrors service_writeforward_test.go's // scriptedDecider trick, made structural: the acceptance is written for the // post's content unconditionally. - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} state, err := tool.RematerializeOne(context.Background(), legacy) require.NoError(t, err) @@ -1124,7 +1131,7 @@ func TestRematerialize_PreservesEveryPublishedField(t *testing.T) { RawRecord: raw, } source := newFakeLegacySource(legacy) - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} _, err := tool.RematerializeOne(context.Background(), legacy) require.NoError(t, err) @@ -1158,7 +1165,7 @@ func TestRematerialize_RefusesWhenADifferentRecordStandsAtTheRkey(t *testing.T) writer := &spyAcceptanceWriter{} legacy := legacyPost(t, rematCommunityDID, rematAuthorDID) source := newFakeLegacySource(legacy) - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} // A DIFFERENT record already stands at the target rkey — its CID is the one a // fresh write would get (so a CID-only verify passes), but its body is NOT this @@ -1223,7 +1230,7 @@ func TestRematerialize_Run_ReconcilesStrandedMigratedRowFromLedger(t *testing.T) // The source does NOT list the stranded record — its community.post is gone. source := newFakeLegacySource() - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} report, err := tool.Run(ctx) require.NoError(t, err) @@ -1276,7 +1283,7 @@ func TestRematerialize_Run_ResolvesAllCredentialsBeforeAnyMutation(t *testing.T) second := legacyPost(t, rematCommunityDID, noCreds) source := newFakeLegacySource(first, second) source.log = log - tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos()} + tool := &posts.Rematerializer{Source: source, Ledger: ledger, AuthorRepos: authors.factory(), Acceptances: writer, CommunityRepos: writer.repos(), Removals: noRematerializeRemovalsIntegration{}, InstanceDID: "did:web:coves-instance.invalid"} _, err := tool.Run(context.Background()) require.NoError(t, err) -- 2.51.2