diff --git a/.golangci.yml b/.golangci.yml index 86e6bf5..cd65e53 100644 --- a/.golangci.yml +++ b/.golangci.yml @@ -41,6 +41,7 @@ linters: - "**/pkg/decode/**" - "**/pkg/checkpoint/**" - "**/pkg/schemachange/**" + - "**/pkg/internal/**" - "!$test" allow: - $gostd @@ -55,6 +56,7 @@ linters: - github.com/block/pg-sprite/pkg/decode - github.com/block/pg-sprite/pkg/checkpoint - github.com/block/pg-sprite/pkg/schemachange + - github.com/block/pg-sprite/pkg/internal/chunksql - github.com/block/pg-sprite/pkg/statement - github.com/block/pg-sprite/pkg/schemadiff sloglint: diff --git a/SAFETY.md b/SAFETY.md index f3a0ff6..7d2ed06 100644 --- a/SAFETY.md +++ b/SAFETY.md @@ -20,8 +20,9 @@ The invariant registry (invariant IDs referenced below) lives in | `pkg/dbconn` — pool defaults, terminate-blockers, retries, RDS TLS, advisory table lock | ✅ core | table-lock primitive exists; wiring into executing modes planned | LK-1 primitive; LK-2 primitives; CO-9 (session hook and `LocalSearchPath`) | | `pkg/preflight` — precondition verifier, refusals | ✅ core | exists; copy-and-swap target proof (shape and dependents) and cluster/volume environment check exist, not yet wired into a route | ST-6, RF-1..RF-5 | | `pkg/executor` — bounded optimistic attempt; native concurrent index build with invalid-index recovery; native sequence executor for the safer idioms | ✅ core | exists (Phase 1: attempt-under-budget; Phase 3.1: concurrent index build; Phase 3.2: sequence executor) | LK-2 (attempt bound + the CONCURRENTLY wait-policy exception), CO-9 (qualified proof reads), ST-9 (create owner verified, never repaired) | -| `pkg/checksum` — chunk verifier, continuous checker, repair | ✅ core | the chunk `Verifier` exists (one read-only `REPEATABLE READ` transaction per chunk so both digests share a snapshot, every column cast to the shadow's type on both sides, the same guard as the copier — owner role, catalog-only `search_path`, `ACCESS SHARE` on both relations, lock confirmation, relation-OID check — and a `Report` of mismatched chunks up to the landed watermark); divergence policy, repair, the proof constructors, and the continuous checker planned | CO-1 (compare; the gate lands with the proof constructors), CO-2, CO-3, CO-9, LK-1 | -| `pkg/copier` — shadow-table chunked copy | ✅ core | contract types and the keyset `Chunker` exist (row-count chunks over the proven key, first chunk open below and last open above, a cut frontier for the applier's discard rule, time-targeted sizing) and the parallel `Copier` (one bounded never-overwriting insert per chunk under the table lock session, frontier-ordered in-flight registry, a resume that first clears the shadow above the watermark, `Position.Classify` for the applier, and the copy step's `progress.WorkSource` — rows from committed chunks, the source's catalog row count, both tables' measured sizes) exist | CO-4 (chunk coverage, copy SQL shape, in-flight registry), LK-1, LK-3 | +| `pkg/checksum` — chunk verifier, divergence policy, repair, continuous checker | ✅ core | the chunk `Verifier` exists (one read-only `REPEATABLE READ` transaction per chunk so both digests share a snapshot, every column cast to the shadow's type on both sides, the same guard as the copier — owner role, catalog-only `search_path`, `ACCESS SHARE` on both relations, lock confirmation, relation-OID check — and a `Report` of mismatched chunks up to the landed watermark); `Check` runs the pass under an explicit `DivergencePolicy` (zero value refused; `ParseDivergencePolicy` for configuration), aborts with the shadow untouched or recopies every differing chunk with the copier's chunk statement in one guarded transaction — every chunk's rows deleted before any are put back, so a unique value the source moved between two chunks lands — then rereads each chunk, returning the committed repairs with any refusal, and mints `CleanWatermark` and `VerifiedShadow` (private constructors) only from a pass that found nothing and repaired nothing; the continuous checker is planned | CO-1, CO-2, CO-3, CO-9, LK-1 | +| `pkg/copier` — shadow-table chunked copy | ✅ core | contract types and the keyset `Chunker` exist (row-count chunks over the proven key, first chunk open below and last open above, a cut frontier for the applier's discard rule, time-targeted sizing) and the parallel `Copier` (one bounded never-overwriting insert per chunk under the table lock session, each in a guarded transaction — owner role, catalog-only `search_path`, `ACCESS SHARE` on both relations, lock confirmation, relation-OID check — frontier-ordered in-flight registry, a resume that first clears the shadow above the watermark, `Position.Classify` for the applier, and the copy step's `progress.WorkSource` — rows from committed chunks, the source's catalog row count, both tables' measured sizes) exist | CO-4 (chunk coverage, copy SQL shape, in-flight registry), LK-1, LK-3 | +| `pkg/internal/chunksql` — the chunk insert statement the copier and the verifier's repair both run | ✅ core | exists; unexported from the module so no caller can run the statement outside the guard both packages wrap around it | CO-4 (copy SQL shape) | | `pkg/applier` — change apply, buffer, flush scheduling | ✅ core | package contract exists; applier planned | CO-4, CO-5, CO-6, CO-8, LK-3 | | `pkg/decode` — logical decoding, LSN/position accounting, per-column presence | ✅ core | contract types exist; decoder planned | ST-4, CO-4, CO-8 | | `pkg/checkpoint` — durable resume state | ✅ core | checkpoint contract exists; persistence planned | ST-1, ST-2 | diff --git a/docs/architecture.md b/docs/architecture.md index 9c22b40..fa547fb 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -201,7 +201,7 @@ different levels of commitment: | `pkg/executor` | Native backend with stable outcome codes: the bounded optimistic attempt, the concurrent index build and its invalid-index recovery, the autocommit safer-sequence runner, the greenfield `CREATE TABLE` path, and the accepted-blocking passthrough primitive; the full `Executor` contract (`Plan`/`Execute`/`Status`/`Abort`) arrives with the copy-and-swap backend | native execution exists | | `pkg/progress` | Strategy-wide, pollable progress snapshots: native phase/elapsed time, sequence position, retry attempt, and server-reported concurrent-index work, and engine-measured copy work polled from the copier (`WorkSource`) | native progress and the copy work source exist | | `pkg/copier` | PK-range chunker over one integer-family primary key with dynamic time-based sizing (produces `Chunk` and `Watermark`; composite keys refused in v1), and the parallel chunked copy into the shadow table (never overwrites; reports the cut frontier, in-flight chunks, and landed watermark for the applier, and fills the progress tracker's copy counters while it runs) — there is no separate chunker package | chunker, copy loop, and progress fillers exist | -| `pkg/checksum` | The mandatory correctness gate: the chunk verifier (single-snapshot per-chunk digests through the shadow's types, reporting the chunks that differ up to the landed watermark); divergence policy, repair primitive, proof constructors, and continuous checker to follow | verifier exists; gate, repair, and checker planned | +| `pkg/checksum` | The mandatory correctness gate: the chunk verifier (single-snapshot per-chunk digests through the shadow's types, reporting the chunks that differ up to the landed watermark), `Check` under an explicit divergence policy (abort, or repair each differing chunk with the copier's statement and reread it), and the proof constructors (`VerifiedShadow`, `CleanWatermark`) minted only by a pass that found nothing and repaired nothing; continuous checker to follow | verifier, policy, repair, and proofs exist; checker planned | | `pkg/decode` | Logical-decoding change capture, LSN accounting, slot lifecycle | Phase 6, 8 | | `pkg/applier` | Change apply onto the shadow (always wins), buffer/dedup, flush scheduling | Phase 6 | | `pkg/schemachange` | Orchestrator: lifecycle, cutover swap + fidelity gate, checkpoint/resume | Phase 7–8 | diff --git a/docs/copy-and-swap-design.md b/docs/copy-and-swap-design.md index b02e2bd..d8123b6 100644 --- a/docs/copy-and-swap-design.md +++ b/docs/copy-and-swap-design.md @@ -360,7 +360,11 @@ mode where known missing changes require repair. The mode must not be inferred f **Alternative considered → deferred.** A universal repair default can conceal executor defects. -**Where enforced.** `pkg/checksum` and `pkg/schemachange`; CO-2, CO-3, ST-4. +**Where enforced.** `pkg/checksum` `Check` takes a `DivergencePolicy` on every pass and refuses +the zero value (`ErrNoDivergencePolicy`); `abort` returns a `DivergenceError` with the shadow +untouched, `repair` recopies every differing chunk with the copier's own statement in one +transaction and reads each again, and a pass with repairs mints no proof. `pkg/schemachange` will select the policy per +lifecycle mode; CO-2, CO-3, ST-4. ### D15 — Capture changes with pgoutput @@ -382,8 +386,8 @@ decoding but adds write-path availability and amplification costs. | --- | --- | --- | | `pkg/dbconn` | Produces `TableLock`, carried by `TableLockSession`; `Confirm` is the in-transaction check every writer runs from its own connection before its first write. | LK-1 | | `pkg/preflight` | Produces `CopySwapTarget`, the copy-and-swap route's proof (the table facts `PreflightedTable` carries plus the v1 shape, replica identity, and dependent-object checks above, minted by `CheckCopySwapShape`); `CheckCopySwapEnvironment` then proves the decoding and headroom facts against that target; owns Tier-3 refusals. | ST-6, RF-1..RF-3 | -| `pkg/copier` | Produces `Chunk` and `Watermark`; `Chunker` (built only from a `CopySwapTarget`) cuts consecutive chunks that tile the whole int64 key space — first open below, last open above — so every key a row can carry belongs to exactly one chunk and a watermark at the largest value means the copy is complete. `Copier` (built from a `CopySwapTarget`, a `Shadow` — the shape `schemachange.BuiltShadow` satisfies — and the table's `TableLockSession`) copies chunks with several workers, each in its own bounded transaction under the owner's role that confirms the lock, takes `ACCESS SHARE` on both relations, and confirms both relation OIDs before one frozen never-overwriting insert; a resumed copy first deletes every shadow row above the watermark in batches of the same guarded shape — the whole shadow after a zero watermark — so the resumed cut frontier and the shadow agree on what is uncut, and its first batch takes `SHARE MODE` on the shadow so a chunk transaction of the earlier run still committing into it ends before the clear reads; the proof check refuses a shadow that is the source by name or by OID, since the clear deletes from one and the copy reads the other; the copy statement carries no conversion expression, so a column whose type differs between the two tables takes the server's assignment cast, and an `ALTER COLUMN TYPE … USING` change needs a conversion-aware copy before the planner's copy-and-swap route is wired to the copier; `Position` snapshots the cut frontier and landed watermark as two `Watermark` values plus the in-flight chunks, and `Position.Classify` is the applier's uncut / in-flight / landed rule. While `Run` is in progress the `Copier` is the `progress.Tracker`'s `WorkSource` (`Options.Tracker`): `rows_copied` is the rows this run's committed chunks inserted, `rows_total` the source's catalog row count read once at the start of the run, and `bytes_copied` / `bytes_total` the shadow's and the source's `pg_table_size` measured at each poll — every counter measured, none projected; on a resumed run the row ratio is that run's share, not the copy's completion. Each measurement runs in a read-only transaction of its own under the chunk budgets, with every catalog name `pg_catalog`-qualified, so a poll on a caller-built pool is still bounded and still reads the real catalog; the copier stops being the source before `Run` returns. The copier's refusals are fail-closed `ErrInvariantViolation` values tagged with their invariant; they take a typed cause in [refusal-classes.md](refusal-classes.md#shadow-operation-refusals-keyed-on-refusalcause) when the orchestrator that runs the copy is wired, so the cause lands with its first importer. | CO-4, LK-1, LK-3 | -| `pkg/checksum` | `Verifier` (built from a `CopySwapTarget`, a `copier.Shadow`, and the table's `TableLockSession`) compares the source with its shadow up to the copier's landed watermark and returns a `Report` of the chunks that differ. It cuts its own chunks with a `copier.Chunker` and digests each in one read-only `REPEATABLE READ` transaction, so a chunk's two digests describe one snapshot and no snapshot outlives one chunk; each transaction runs the copier's guard — owner role, catalog-only `search_path`, `ACCESS SHARE` on both relations taken before the snapshot, lock confirmation, relation-OID check — and pins `extra_float_digits` to its maximum, since the digest hashes each row's text rendering and a database or role configured at zero or below would render two floats that differ only in their last digits the same. Both sides run the identical frozen statement — row count plus `md5` of the key-ordered concatenation of each row's `md5(ROW(col::shadow_type, …)::text)` over `pk BETWEEN $1 AND $2` — so only the data can differ, and the cast on every column is D7: a converted column hashes as the value the shadow holds. Only the copy columns are compared; generated columns present on both sides are not yet hashed. A `Report` proves nothing — the divergence policy, repair, and the constructors of `VerifiedShadow` and `CleanWatermark` (private to this package) land on top of it. | CO-1, CO-2, CO-3, CO-9, LK-1 | +| `pkg/copier` | Produces `Chunk` and `Watermark`; `Chunker` (built only from a `CopySwapTarget`) cuts consecutive chunks that tile the whole int64 key space — first open below, last open above — so every key a row can carry belongs to exactly one chunk and a watermark at the largest value means the copy is complete. `Copier` (built from a `CopySwapTarget`, a `Shadow` — the shape `schemachange.BuiltShadow` satisfies — and the table's `TableLockSession`) copies chunks with several workers, each in its own bounded transaction under the owner's role with the catalog alone on its `search_path` that confirms the lock, takes `ACCESS SHARE` on both relations, and confirms both relation OIDs before one frozen never-overwriting insert — the chunk statement `pkg/internal/chunksql` holds for the copier and the verifier's repair alike; a resumed copy first deletes every shadow row above the watermark in batches of the same guarded shape — the whole shadow after a zero watermark — so the resumed cut frontier and the shadow agree on what is uncut, and its first batch takes `SHARE MODE` on the shadow so a chunk transaction of the earlier run still committing into it ends before the clear reads; the proof check refuses a shadow that is the source by name or by OID, since the clear deletes from one and the copy reads the other; the copy statement carries no conversion expression, so a column whose type differs between the two tables takes the server's assignment cast, and an `ALTER COLUMN TYPE … USING` change needs a conversion-aware copy before the planner's copy-and-swap route is wired to the copier; `Position` snapshots the cut frontier and landed watermark as two `Watermark` values plus the in-flight chunks, and `Position.Classify` is the applier's uncut / in-flight / landed rule. While `Run` is in progress the `Copier` is the `progress.Tracker`'s `WorkSource` (`Options.Tracker`): `rows_copied` is the rows this run's committed chunks inserted, `rows_total` the source's catalog row count read once at the start of the run, and `bytes_copied` / `bytes_total` the shadow's and the source's `pg_table_size` measured at each poll — every counter measured, none projected; on a resumed run the row ratio is that run's share, not the copy's completion. Each measurement runs in a read-only transaction of its own under the chunk budgets, with every catalog name `pg_catalog`-qualified, so a poll on a caller-built pool is still bounded and still reads the real catalog; the copier stops being the source before `Run` returns. The copier's refusals are fail-closed `ErrInvariantViolation` values tagged with their invariant; they take a typed cause in [refusal-classes.md](refusal-classes.md#shadow-operation-refusals-keyed-on-refusalcause) when the orchestrator that runs the copy is wired, so the cause lands with its first importer. | CO-4, LK-1, LK-3 | +| `pkg/checksum` | `Verifier` (built from a `CopySwapTarget`, a `copier.Shadow`, and the table's `TableLockSession`) compares the source with its shadow up to the copier's landed watermark and returns a `Report` of the chunks that differ. It cuts its own chunks with a `copier.Chunker` and digests each in one read-only `REPEATABLE READ` transaction, so a chunk's two digests describe one snapshot and no snapshot outlives one chunk; each transaction runs the copier's guard — owner role, catalog-only `search_path`, `ACCESS SHARE` on both relations taken before the snapshot, lock confirmation, relation-OID check — and pins `extra_float_digits` to its maximum, since the digest hashes each row's text rendering and a database or role configured at zero or below would render two floats that differ only in their last digits the same. Both sides run the identical frozen statement — row count plus `md5` of the key-ordered concatenation of each row's `md5(ROW(col::shadow_type, …)::text)` over `pk BETWEEN $1 AND $2` — so only the data can differ, and the cast on every column is D7: a converted column hashes as the value the shadow holds. Only the copy columns are compared; generated columns present on both sides are not yet hashed. A `Report` proves nothing. `Check` runs the same pass under a `DivergencePolicy` the caller states every time (the zero value is refused): `abort` returns a `DivergenceError` carrying the report with the shadow untouched; `repair` replaces every differing chunk inside one guarded read-write transaction — delete the shadow's rows over every chunk's key range, then the copy statement itself for every chunk, so a unique value the source moved from one differing chunk to another lands instead of colliding with the stale row — and digests each chunk again in a fresh snapshot, returning a `RepairError` for the first that still differs alongside the `Outcome` listing every repair that committed. A repair pass assumes nothing else writes the shadow and the source rows it recopies hold still until the rereads; a write inside the pass reads as a `RepairError`, never as a second repair. `ParseDivergencePolicy` lets a caller refuse a configured policy before any pass. Only a pass that found nothing and repaired nothing mints the proofs, whose constructors are private to the package: a `CleanWatermark` at the watermark it read through, and a `VerifiedShadow` only when that watermark is complete; `Outcome.Clean` is true only when the clean watermark was minted. A pass with repairs returns its `Repair`s and no proof; the next pass mints. | CO-1, CO-2, CO-3, CO-9, LK-1 | | `pkg/decode` | Produces `ChangeEvent`, including per-column presence and `OldKey` for an UPDATE that moved the primary key. | ST-3, ST-4, CO-4, CO-8 | | `pkg/applier` | Applies presence-aware events from the per-key buffer. | CO-4, CO-5, CO-6, CO-8, LK-3 | | `pkg/checkpoint` | Produces `Checkpoint`. | ST-1, ST-2 | diff --git a/docs/invariants.md b/docs/invariants.md index 687bf70..b911070 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -37,12 +37,15 @@ A migration that cannot prove shadow == source **must refuse to cut over**. No f capture mechanism removes the gate; it is also the repair primitive for [slot-loss reconciliation](low-level-design.md#failover-during-migration-what-survives-and-what-doesnt). *Enforced:* cutover entry condition. *Enforced today:* `pkg/checksum` `Verifier` — the -comparison the gate will demand: every chunk up to the landed watermark digested on both sides +comparison the gate demands: every chunk up to the landed watermark digested on both sides inside one read-only `REPEATABLE READ` transaction, with every column cast to the shadow's type, and every differing chunk reported with its two row counts (clean-copy, changed-row, missing-and-extra-row, converted-type, watermark-clamp, shadowing-`search_path`, -replaced-relation, and lost-lock tests). *Planned enforcement:* the `VerifiedShadow` constructor -accepts only a clean `Report`, and cutover accepts only a `VerifiedShadow`. +replaced-relation, and lost-lock tests); `Check` mints a `VerifiedShadow` only from a pass that +compared through the complete watermark, found no difference, and repaired nothing, and a +`CleanWatermark` from any clean pass — both constructors are private to the package, and a +partial clean pass mints only the watermark (clean-complete, partial-watermark, and +repairs-mint-nothing tests). *Planned enforcement:* cutover accepts only a `VerifiedShadow`. *Source:* [design-principles](design-principles.md#correctness-and-safety), risks-and-mitigations; Spirit's "never skip it". @@ -59,6 +62,11 @@ loop; a deferred-cutover mode, if ever built, would run the same checker longer) from a stale watermark after a continuous-checker repair would let a re-run "pass" by verifying only trailing chunks — silently neutralizing a deliberate divergence abort. *Enforced:* checkpoint writer (watermark dropped unless **all** active checkers are clean). +*Enforced today:* `pkg/checksum` `Check` — a pass that repaired any chunk returns its `Repair`s +and no `CleanWatermark`, even though every repaired chunk was read again and found equal; the +proof comes only from the next pass, which reads every chunk fresh (repairs-mint-nothing test). +A repaired chunk that still differs on the fresh read is a `RepairError`, not a second repair +(repair-did-not-take test). The checkpoint writer that persists the watermark is planned. *Source:* Spirit `pkg/migration/runner.go` + `pkg/move/runner.go` ("Safety invariant"). ### CO-3 — Divergence policy is an explicit setting, never inferred @@ -74,7 +82,15 @@ inferred from whether a recopier happens to be wired up: as fatal). - The two knobs stay decoupled: fatal-divergence aborts even if a recopier is supplied. -*Enforced:* checker configuration per lifecycle mode. *Source:* Spirit AGENTS.md +*Enforced:* checker configuration per lifecycle mode. *Enforced today:* `pkg/checksum` +`Check` takes a `DivergencePolicy` on every call and refuses the zero value and any string +that is not `abort` or `repair` with `ErrNoDivergencePolicy` before reading anything; under +`abort` a difference is a `DivergenceError` carrying the report with the shadow untouched, +under `repair` every differing chunk is recopied with the copier's own statement inside one +guarded transaction (delete every chunk's shadow rows, then the copy statement for every chunk) +and read again, with the committed repairs reported alongside any refusal (no-policy, +abort-leaves-shadow-alone, repairs-every-chunk, moved-unique-value, and +later-repair-fails tests). *Source:* Spirit AGENTS.md (block/spirit#994 policy) — maps directly onto our failover-reconcile design. ### CO-4 — The copy/apply ordering invariants diff --git a/docs/low-level-design.md b/docs/low-level-design.md index f06a7c4..35d9cbc 100644 --- a/docs/low-level-design.md +++ b/docs/low-level-design.md @@ -440,9 +440,12 @@ implicit in the implementation. The races to design against: the shadow regresses to the older image. - **Ghost-row resurrection.** The copier reads a row, the applier applies that row's `DELETE`, then the chunk insert lands — re-inserting a row that no longer exists on the source. -- **The same races during reconciliation.** The checksum-repair pass after +- **Reconciliation is kept out of the race.** The checksum-repair pass after [slot loss](#failover-during-migration-what-survives-and-what-doesnt) re-copies divergent - chunks while the new slot's stream is being applied — the same two races, a second exposure. + chunks by deleting and re-inserting their rows, which does overwrite; it runs before the new + slot's stream is applied, with no applier writing the shadow, and a source write that lands + inside the pass makes the repaired chunk read different again and fails the pass closed + rather than being repaired twice. The invariants that resolve them (Spirit's model, translated): diff --git a/pkg/checksum/check.go b/pkg/checksum/check.go new file mode 100644 index 0000000..d5c2b27 --- /dev/null +++ b/pkg/checksum/check.go @@ -0,0 +1,111 @@ +package checksum + +import ( + "context" + "fmt" + + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/block/pg-sprite/pkg/copier" +) + +// Outcome is one pass run under a divergence policy: what it compared and +// found, what it repaired, and — only when it found nothing and repaired +// nothing — the proofs the cutover and the checkpoint demand. A pass that +// repaired a chunk has read that chunk clean afterwards, but it mints no +// proof: the proof comes from the next pass, which reads every chunk +// fresh (CO-2). Check returns an Outcome with its error too, so what a +// failed pass compared and wrote to the shadow is never lost. +type Outcome struct { + // Report is the comparison the pass made before any repair. It is the + // zero Report when no comparison ran. + Report Report + // Repairs lists the chunks the pass recopied and committed under + // DivergenceRepair, in ascending key order, whether or not every one + // then read equal. It is empty under DivergenceAbort and for a clean + // pass. + Repairs []Repair + + clean CleanWatermark + verified VerifiedShadow +} + +// Clean reports whether the pass read every compared chunk equal and so +// minted the clean watermark. A pass that found a difference, repaired +// anything, or was refused is not clean, whatever its Report says. +func (o Outcome) Clean() bool { return o.clean.Watermark().Valid() } + +// CleanWatermark returns the proof that every chunk through the pass's +// watermark was read clean, and whether the pass minted one. +func (o Outcome) CleanWatermark() (CleanWatermark, bool) { + return o.clean, o.clean.Watermark().Valid() +} + +// VerifiedShadow returns the proof that the whole shadow equals its source, +// and whether the pass minted one: a clean pass through the complete +// watermark does; a clean pass through a partial watermark proves only +// its prefix and does not. +func (o Outcome) VerifiedShadow() (VerifiedShadow, bool) { + return o.verified, o.verified.Table() != "" +} + +// Check runs one verification pass through the copier's landed watermark +// and acts on what it finds as policy says. A clean pass mints a +// CleanWatermark, and a VerifiedShadow when the watermark is complete. +// Under DivergenceAbort a difference is returned as a *DivergenceError +// carrying the report, with the shadow untouched. Under DivergenceRepair +// every differing chunk is recopied from the source in one transaction and +// then read again in a fresh snapshot; a chunk still different after its +// repair is returned as a *RepairError, with the Outcome listing every +// repair that committed. A pass whose repairs all took returns a nil error +// and no proof. The policy is stated for every pass; there is no default +// (CO-3). +// +// A repair pass is a reconciliation: it assumes nothing else writes the +// shadow and no source row in a differing chunk changes between the pass +// and its rereads — the pass after slot loss runs before change capture +// resumes, and the application's writes to the source are what the resumed +// capture then carries. A write that lands inside the pass makes the +// repaired chunk read different again, and the pass refuses it as a +// RepairError rather than repair it twice. +func (v *Verifier) Check(ctx context.Context, pool *pgxpool.Pool, through copier.Watermark, policy DivergencePolicy) (Outcome, error) { + if err := policy.validate(); err != nil { + return Outcome{}, v.verifyError(err) + } + report, err := v.Verify(ctx, pool, through) + if err != nil { + return Outcome{}, err + } + if report.Clean() { + return v.prove(report), nil + } + // INV: CO-3 + switch policy { + case DivergenceAbort: + return Outcome{Report: report}, v.verifyError(&DivergenceError{Report: report}) + case DivergenceRepair: + repairs, err := v.repairAll(ctx, pool, report.Mismatches) + if err != nil { + return Outcome{Report: report, Repairs: repairs}, v.verifyError(err) + } + return Outcome{Report: report, Repairs: repairs}, nil + default: + return Outcome{Report: report}, fmt.Errorf("%w (CO-3): divergence policy %q passed validation", ErrInvariantViolation, string(policy)) + } +} + +// verifyError names the table a pass's refusal is about. +func (v *Verifier) verifyError(err error) error { + return fmt.Errorf("verify %s.%s: %w", v.target.Schema(), v.target.Table(), err) +} + +// prove mints the proofs a clean pass earns: the clean watermark always, +// and the verified shadow only when nothing lies beyond the watermark. +func (v *Verifier) prove(report Report) Outcome { + // INV: CO-1, CO-2 + outcome := Outcome{Report: report, clean: newCleanWatermark(report.Through)} + if report.Through.Complete() { + outcome.verified = newVerifiedShadow(v.shadow, report.Through, v.opts.Clock.Now()) + } + return outcome +} diff --git a/pkg/checksum/check_integration_test.go b/pkg/checksum/check_integration_test.go new file mode 100644 index 0000000..9556a32 --- /dev/null +++ b/pkg/checksum/check_integration_test.go @@ -0,0 +1,237 @@ +package checksum_test + +import ( + "errors" + "math" + "testing" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/checksum" + "github.com/block/pg-sprite/pkg/copier" + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/preflight" + "github.com/block/pg-sprite/pkg/schemachange" +) + +// check runs one pass over fixed chunks under policy. +func (f verifierFixture) check(t *testing.T, pool *pgxpool.Pool, target preflight.CopySwapTarget, shadow schemachange.BuiltShadow, lock *dbconn.TableLockSession, through copier.Watermark, policy checksum.DivergencePolicy) (checksum.Outcome, error) { + t.Helper() + v, err := checksum.NewVerifier(target, shadow, lock, checksum.Options{Chunker: chunkRows}) + require.NoError(t, err) + return v.Check(t.Context(), pool, through, policy) +} + +// corrupt puts three differences into the shadow, one per chunk: a changed +// value in the middle chunk, a missing row in the first, and a row the +// source does not hold in the last. +func (f verifierFixture) corrupt(t *testing.T, shadow schemachange.BuiltShadow) { + t.Helper() + f.exec(t, "DELETE FROM "+f.shadowName(shadow)+" WHERE id = 5") + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1500") + f.exec(t, "INSERT INTO "+f.shadowName(shadow)+" (id, qty) VALUES (2600, 2600)") +} + +// shadowRow reads one shadow row's qty, and whether the row exists. +func (f verifierFixture) shadowRow(t *testing.T, shadow schemachange.BuiltShadow, id int64) (qty int64, exists bool) { + t.Helper() + err := f.pool.QueryRow(t.Context(), "SELECT qty FROM "+f.shadowName(shadow)+" WHERE id = $1", id).Scan(&qty) + if errors.Is(err, pgx.ErrNoRows) { + return 0, false + } + require.NoError(t, err) + return qty, true +} + +// assertCorrupted asserts the three differences corrupt made are still in +// the shadow, so a pass that must not write did not. +func (f verifierFixture) assertCorrupted(t *testing.T, shadow schemachange.BuiltShadow) { + t.Helper() + _, exists := f.shadowRow(t, shadow, 5) + assert.False(t, exists, "the missing row is still missing") + qty, _ := f.shadowRow(t, shadow, 1500) + assert.Equal(t, int64(0), qty, "the changed row still differs") + _, exists = f.shadowRow(t, shadow, 2600) + assert.True(t, exists, "the extra row is still there") +} + +// assertNoProofs asserts the outcome minted neither proof (CO-2). +func assertNoProofs(t *testing.T, outcome checksum.Outcome) { + t.Helper() + _, minted := outcome.CleanWatermark() + assert.False(t, minted, "no clean watermark") + _, minted = outcome.VerifiedShadow() + assert.False(t, minted, "no verified shadow") +} + +// A clean pass through the complete watermark mints both proofs: the clean +// watermark, and the verified shadow naming the relations the pass compared +// (CO-1). Which policy was stated does not matter when nothing differed. +func TestCheckMintsBothProofsForACleanCompletePass(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + + outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceAbort) + require.NoError(t, err) + assert.True(t, outcome.Clean()) + assert.Empty(t, outcome.Repairs) + assert.Equal(t, 3, outcome.Report.Chunks) + clean, minted := outcome.CleanWatermark() + require.True(t, minted) + assert.Equal(t, copier.NewWatermark(math.MaxInt64), clean.Watermark()) + verified, minted := outcome.VerifiedShadow() + require.True(t, minted) + assert.Equal(t, f.schema, verified.Schema()) + assert.Equal(t, "orders", verified.Table()) + assert.Equal(t, shadow.ShadowTable(), verified.Shadow()) + assert.Equal(t, copier.NewWatermark(math.MaxInt64), verified.Watermark()) + assert.False(t, verified.VerifiedAt().IsZero()) +} + +// A clean pass through a partial watermark proves only the prefix it read: +// it mints the clean watermark at that frontier and no verified shadow, +// because the keys above it were never compared (CO-1). +func TestCheckMintsOnlyTheCleanWatermarkForAPartialPass(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1600") + + outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(1500), checksum.DivergenceAbort) + require.NoError(t, err) + assert.True(t, outcome.Clean(), "the difference at 1600 is above the watermark") + clean, minted := outcome.CleanWatermark() + require.True(t, minted) + assert.Equal(t, copier.NewWatermark(1500), clean.Watermark()) + _, minted = outcome.VerifiedShadow() + assert.False(t, minted, "a partial pass cannot vouch for the whole shadow") +} + +// Under DivergenceAbort a difference ends the pass at its report: the +// error carries every differing chunk, no proof is minted, and the shadow +// is left exactly as the pass found it for the operator to inspect. +func TestCheckAbortsOnDivergenceAndLeavesTheShadowAlone(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.corrupt(t, shadow) + + outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceAbort) + var divergence *checksum.DivergenceError + require.ErrorAs(t, err, &divergence) + require.Len(t, divergence.Report.Mismatches, 3) + assert.Equal(t, chunk(t, math.MinInt64, 1000), divergence.Report.Mismatches[0].Chunk) + assert.Equal(t, chunk(t, 1001, 2000), divergence.Report.Mismatches[1].Chunk) + assert.Equal(t, chunk(t, 2001, math.MaxInt64), divergence.Report.Mismatches[2].Chunk) + assert.EqualError(t, err, "verify "+f.schema+".orders: source and shadow differ in 3 of 3 chunks through watermark 9223372036854775807") + assert.False(t, outcome.Clean(), "a pass that found a difference is not clean") + assert.Equal(t, divergence.Report, outcome.Report, "the outcome carries the comparison too") + assert.Empty(t, outcome.Repairs) + assertNoProofs(t, outcome) + f.assertCorrupted(t, shadow) +} + +// Under DivergenceRepair every differing chunk is recopied from the source +// and read again: the outcome lists each repair with what it removed and +// put back, the shadow equals the source afterwards, and the pass still +// mints no proof (CO-2). The next pass, reading every chunk fresh, finds +// the shadow clean and mints both. +func TestCheckRepairsEveryDifferingChunkAndMintsNothing(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.corrupt(t, shadow) + + outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + require.NoError(t, err) + assert.False(t, outcome.Clean(), "the report is the comparison before the repairs") + require.Len(t, outcome.Repairs, 3) + missing, changed, extra := outcome.Repairs[0], outcome.Repairs[1], outcome.Repairs[2] + assert.Equal(t, chunk(t, math.MinInt64, 1000), missing.Mismatch.Chunk) + assert.Equal(t, int64(999), missing.Removed, "the shadow held 999 of the chunk's 1000 keys") + assert.Equal(t, int64(1000), missing.Inserted) + assert.Equal(t, chunk(t, 1001, 2000), changed.Mismatch.Chunk) + assert.Equal(t, int64(1000), changed.Removed) + assert.Equal(t, int64(1000), changed.Inserted) + assert.Equal(t, chunk(t, 2001, math.MaxInt64), extra.Mismatch.Chunk) + assert.Equal(t, int64(501), extra.Removed, "the extra row goes with the chunk") + assert.Equal(t, int64(500), extra.Inserted, "the source holds 500 keys above 2000") + assertNoProofs(t, outcome) + + _, exists := f.shadowRow(t, shadow, 5) + assert.True(t, exists, "the missing row is back") + qty, _ := f.shadowRow(t, shadow, 1500) + assert.Equal(t, int64(1500), qty, "the changed row holds the source's value") + _, exists = f.shadowRow(t, shadow, 2600) + assert.False(t, exists, "the extra row is gone") + + again, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + require.NoError(t, err) + assert.True(t, again.Clean()) + assert.Empty(t, again.Repairs) + _, minted := again.CleanWatermark() + assert.True(t, minted, "the fresh pass mints the clean watermark") + _, minted = again.VerifiedShadow() + assert.True(t, minted, "and the verified shadow") +} + +// A repair that does not take is refused, not retried: something other +// than the copier writes the shadow — here a trigger that zeroes every +// inserted qty — and recopying again would not say why. The recopy has +// committed before the fresh read finds it wrong, so the chunk is left as +// the trigger made it; the error names the chunk and the second reading. +func TestCheckRefusesARepairThatDidNotTake(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1500") + f.exec(t, ` + CREATE FUNCTION %s.zero_qty() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + NEW.qty := 0; + RETURN NEW; + END + $$`) + f.exec(t, "CREATE TRIGGER zero_qty BEFORE INSERT ON "+f.shadowName(shadow)+" FOR EACH ROW EXECUTE FUNCTION %s.zero_qty()") + + outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + var failed *checksum.RepairError + require.ErrorAs(t, err, &failed) + assert.Equal(t, chunk(t, 1001, 2000), failed.Repair.Mismatch.Chunk) + assert.Equal(t, int64(1000), failed.Repair.Removed) + assert.Equal(t, int64(1000), failed.Repair.Inserted) + assert.Equal(t, chunk(t, 1001, 2000), failed.After.Chunk) + assert.Equal(t, int64(1000), failed.After.Source.Rows) + assert.Equal(t, int64(1000), failed.After.Shadow.Rows, "every row is back, with the wrong value") + assert.NotEqual(t, failed.After.Source.Hash, failed.After.Shadow.Hash) + assert.EqualError(t, err, "verify "+f.schema+".orders: chunk [1001, 2000] still differs after its repair: source 1000 rows, shadow 1000 rows") + assert.False(t, outcome.Clean()) + assert.Equal(t, []checksum.Repair{failed.Repair}, outcome.Repairs, "the committed recopy is reported with the refusal") + assertNoProofs(t, outcome) + + var zeroed int64 + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT count(*) FROM "+f.shadowName(shadow)+" WHERE id BETWEEN 1001 AND 2000 AND qty = 0").Scan(&zeroed)) + assert.Equal(t, int64(1000), zeroed, "the committed recopy stands as the trigger wrote it") +} + +// A pass must say what a difference means before it reads anything: the +// zero policy and a string that is not one of the two are refused with no +// pass run and no write made (CO-3). +func TestCheckRefusesAPassWithoutAPolicy(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.corrupt(t, shadow) + + for name, policy := range map[string]checksum.DivergencePolicy{ + "zero policy": "", + "unknown policy": "fix", + } { + t.Run(name, func(t *testing.T) { + outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), policy) + require.ErrorIs(t, err, checksum.ErrNoDivergencePolicy) + assert.Zero(t, outcome.Report.Chunks, "no pass ran") + assert.False(t, outcome.Clean(), "a pass that compared nothing is not clean") + assertNoProofs(t, outcome) + }) + } + f.assertCorrupted(t, shadow) +} diff --git a/pkg/checksum/doc.go b/pkg/checksum/doc.go index f4ed965..9418e2d 100644 --- a/pkg/checksum/doc.go +++ b/pkg/checksum/doc.go @@ -1,13 +1,18 @@ -// Package checksum compares the shadow table with its source and defines -// the proofs the cutover will demand of that comparison (CO-1, CO-2, CO-3). +// Package checksum compares the shadow table with its source and mints the +// proofs the cutover demands of that comparison (CO-1, CO-2, CO-3). // A Verifier reads every key at or below the copier's landed watermark in // chunks, digesting both tables inside one read-only REPEATABLE READ // transaction per chunk so the two digests describe one snapshot, with // every column cast to the type the shadow declares so a schema change that // converts a column compares as the shadow holds it (D7). Each transaction // runs under the table's lock session and refuses relations that are no -// longer the ones the proofs describe. A pass reports what differed; acting -// on a difference is the caller's divergence policy. VerifiedShadow and -// CleanWatermark are the proofs a clean pass will mint; their constructors -// are private to this package. +// longer the ones the proofs describe. Verify reports what differed; Check +// runs the same pass under a DivergencePolicy the caller states for every +// pass — abort on a difference, or recopy every differing chunk from the +// source with the copier's own statement in one transaction and read each +// again, assuming nothing else writes the shadow and the source rows it +// recopies hold still until the rereads. Only a pass +// that found no difference and repaired nothing mints a CleanWatermark, +// and a VerifiedShadow when its watermark is complete; the proofs' +// constructors are private to this package. package checksum diff --git a/pkg/checksum/guard.go b/pkg/checksum/guard.go index 5bb2394..92840cf 100644 --- a/pkg/checksum/guard.go +++ b/pkg/checksum/guard.go @@ -17,14 +17,18 @@ import ( // names a relation that does not exist. const sqlstateUndefinedTable = "42P01" -// begin opens the read-only transaction one chunk's two digests run in and -// guards it. The transaction is REPEATABLE READ so its two reads see one -// snapshot: the source and the shadow are compared as they stood at the -// same instant, and the snapshot is taken by the first query after both -// relations are locked, so no rename or drop can slip between the lock and -// the reads. -func (v *Verifier) begin(ctx context.Context, pool *pgxpool.Pool) (pgx.Tx, error) { - tx, err := pool.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly}) +// snapshotRead is the transaction one chunk's two digests run in: read-only +// and REPEATABLE READ, so its two reads see one snapshot — the source and +// the shadow are compared as they stood at the same instant — and the +// snapshot is taken by the first query after both relations are locked, so +// no rename or drop can slip between the lock and the reads. +func snapshotRead() pgx.TxOptions { + return pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly} +} + +// begin opens a transaction with the given options and guards it. +func (v *Verifier) begin(ctx context.Context, pool *pgxpool.Pool, options pgx.TxOptions) (pgx.Tx, error) { + tx, err := pool.BeginTx(ctx, options) if err != nil { return nil, fmt.Errorf("begin verification of %s.%s: %w", v.target.Schema(), v.target.Table(), err) } @@ -37,13 +41,14 @@ func (v *Verifier) begin(ctx context.Context, pool *pgxpool.Pool) (pgx.Tx, error return tx, nil } -// guard prepares the transaction every read runs in: it bounds it, puts the -// catalog alone on its search_path, puts it under the source owner's role, -// locks both relations by name, and then confirms from this connection that -// the lock session's backend still holds the table and that the source and -// shadow are still the relations the proofs describe. Holding even ACCESS -// SHARE keeps a DROP or rename from completing until the transaction ends, -// so the identity check stays true for every statement after it. +// guard prepares the transaction every read and every repair runs in: it +// bounds it, puts the catalog alone on its search_path, puts it under the +// source owner's role, locks both relations by name, and then confirms from +// this connection that the lock session's backend still holds the table and +// that the source and shadow are still the relations the proofs describe. +// Holding even ACCESS SHARE keeps a DROP or rename from completing until the +// transaction ends, so the identity check stays true for every statement +// after it. func (v *Verifier) guard(ctx context.Context, tx pgx.Tx) error { if err := setVerifySession(ctx, tx, v.target.OwnerRole(), v.opts); err != nil { return err diff --git a/pkg/checksum/policy.go b/pkg/checksum/policy.go new file mode 100644 index 0000000..fa598d8 --- /dev/null +++ b/pkg/checksum/policy.go @@ -0,0 +1,81 @@ +package checksum + +import ( + "errors" + "fmt" +) + +// DivergencePolicy says what a pass does with a chunk whose source and +// shadow differ. It is a setting the caller states for every pass, never a +// default and never inferred from what happens to be wired up (CO-3): in +// the steady state a divergence is a defect and aborts; only a pass that +// expects differences — reconciliation after a lost replication slot — +// repairs them. The zero value names no policy and is refused. +type DivergencePolicy string + +const ( + // DivergenceAbort ends the pass at its report: nothing is written, and + // the caller refuses to cut over. + DivergenceAbort DivergencePolicy = "abort" + // DivergenceRepair recopies every differing chunk from the source and + // reads it again; the pass still mints no proof. + DivergenceRepair DivergencePolicy = "repair" +) + +// ErrNoDivergencePolicy reports a pass asked to run without saying what a +// difference means. +var ErrNoDivergencePolicy = errors.New("no divergence policy") + +// ParseDivergencePolicy returns the policy a configuration value or flag +// names, or ErrNoDivergencePolicy when it names neither, so a caller can +// refuse a setting when it loads it rather than at its first pass. +func ParseDivergencePolicy(value string) (DivergencePolicy, error) { + policy := DivergencePolicy(value) + if err := policy.validate(); err != nil { + return "", err + } + return policy, nil +} + +// validate refuses every value but the two named policies. +func (p DivergencePolicy) validate() error { + // INV: CO-3 + switch p { + case DivergenceAbort, DivergenceRepair: + return nil + case "": + return fmt.Errorf("%w: a pass must state %q or %q", ErrNoDivergencePolicy, DivergenceAbort, DivergenceRepair) + default: + return fmt.Errorf("%w: %q is not %q or %q", ErrNoDivergencePolicy, string(p), DivergenceAbort, DivergenceRepair) + } +} + +// DivergenceError is a pass that found differences under DivergenceAbort. +// It carries the report so the caller can say which chunks differed; the +// shadow is as the pass found it. +type DivergenceError struct { + // Report is the pass that found the differences. + Report Report +} + +func (e *DivergenceError) Error() string { + return fmt.Sprintf("source and shadow differ in %d of %d chunks through watermark %d", len(e.Report.Mismatches), e.Report.Chunks, e.Report.Through.Value()) +} + +// RepairError is a chunk that still differed when read again after its +// repair. A recopy that does not converge means something other than the +// copier writes the shadow, or the source changed between the recopy and +// the read, which a repair pass assumes does not happen; either way the +// pass cannot vouch for the chunk, and repairing it again would not say +// why. The recopy has committed, with every other chunk's: the Outcome +// returned alongside lists them. +type RepairError struct { + // Repair is the recopy that did not take. + Repair Repair + // After is the chunk as the fresh read found it. + After Mismatch +} + +func (e *RepairError) Error() string { + return fmt.Sprintf("chunk [%d, %d] still differs after its repair: source %d rows, shadow %d rows", e.After.Chunk.Lower(), e.After.Chunk.Upper(), e.After.Source.Rows, e.After.Shadow.Rows) +} diff --git a/pkg/checksum/policy_test.go b/pkg/checksum/policy_test.go new file mode 100644 index 0000000..d6c32ff --- /dev/null +++ b/pkg/checksum/policy_test.go @@ -0,0 +1,102 @@ +package checksum + +import ( + "math" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/copier" +) + +// A pass must state one of the two named policies; the zero value is no +// policy, and a string that happens to type-check is not one either (CO-3). +func TestDivergencePolicyValidate(t *testing.T) { + require.NoError(t, DivergenceAbort.validate()) + require.NoError(t, DivergenceRepair.validate()) + + err := DivergencePolicy("").validate() + require.ErrorIs(t, err, ErrNoDivergencePolicy) + assert.EqualError(t, err, `no divergence policy: a pass must state "abort" or "repair"`) + + err = DivergencePolicy("Abort").validate() + require.ErrorIs(t, err, ErrNoDivergencePolicy, "policies are exact, not case-folded") + assert.EqualError(t, err, `no divergence policy: "Abort" is not "abort" or "repair"`) +} + +// A caller loading the policy from configuration gets the typed value for +// the two names and the same refusal Check would give for anything else, +// so a wrong setting is refused before any copy runs. +func TestParseDivergencePolicy(t *testing.T) { + policy, err := ParseDivergencePolicy("abort") + require.NoError(t, err) + assert.Equal(t, DivergenceAbort, policy) + policy, err = ParseDivergencePolicy("repair") + require.NoError(t, err) + assert.Equal(t, DivergenceRepair, policy) + + policy, err = ParseDivergencePolicy("") + require.ErrorIs(t, err, ErrNoDivergencePolicy) + assert.Empty(t, policy, "a refused value yields no policy to pass on") + _, err = ParseDivergencePolicy("fix") + require.ErrorIs(t, err, ErrNoDivergencePolicy) +} + +// An Outcome nobody minted carries no proof and is not clean: a consumer +// that reads Clean, or either flag, before the error cannot be handed the +// forgeable zero value as a clean pass. +func TestOutcomeZeroValueMintsNothing(t *testing.T) { + var outcome Outcome + assert.False(t, outcome.Clean(), "no pass read anything clean") + _, minted := outcome.CleanWatermark() + assert.False(t, minted) + _, minted = outcome.VerifiedShadow() + assert.False(t, minted) +} + +// A clean pass through a partial watermark proves its prefix and nothing +// more; only the complete watermark earns the whole-shadow proof (CO-1). +func TestProveMintsTheVerifiedShadowOnlyWhenComplete(t *testing.T) { + now := time.Date(2026, time.October, 1, 12, 0, 0, 0, time.UTC) + v := &Verifier{ + shadow: fakeShadow{schema: "app", source: "orders", shadow: "_pgsprite_orders_new", sourceOID: 1, shadowOID: 2, columns: []string{"id"}}, + opts: Options{Clock: fixedClock(now)}, + } + + partial := v.prove(Report{Through: copier.NewWatermark(1500), Chunks: 2, Rows: 1500}) + clean, minted := partial.CleanWatermark() + require.True(t, minted) + assert.Equal(t, copier.NewWatermark(1500), clean.Watermark()) + _, minted = partial.VerifiedShadow() + assert.False(t, minted, "keys above 1500 were not compared") + + complete := v.prove(Report{Through: copier.NewWatermark(math.MaxInt64), Chunks: 3, Rows: 2500}) + verified, minted := complete.VerifiedShadow() + require.True(t, minted) + assert.Equal(t, "app", verified.Schema()) + assert.Equal(t, "orders", verified.Table()) + assert.Equal(t, "_pgsprite_orders_new", verified.Shadow()) + assert.Equal(t, copier.NewWatermark(math.MaxInt64), verified.Watermark()) + assert.Equal(t, now, verified.VerifiedAt()) +} + +// The two divergence errors name the chunks and counts an operator needs +// to see what differed without opening the report. +func TestDivergenceErrorMessages(t *testing.T) { + c, err := copier.NewChunk(1001, 2000) + require.NoError(t, err) + mismatch := Mismatch{Chunk: c, Source: Digest{Rows: 1000, Hash: "a"}, Shadow: Digest{Rows: 999, Hash: "b"}} + + divergence := &DivergenceError{Report: Report{Through: copier.NewWatermark(2500), Chunks: 3, Mismatches: []Mismatch{mismatch}}} + assert.EqualError(t, divergence, "source and shadow differ in 1 of 3 chunks through watermark 2500") + + repair := &RepairError{Repair: Repair{Mismatch: mismatch, Removed: 999, Inserted: 1000}, After: mismatch} + assert.EqualError(t, repair, "chunk [1001, 2000] still differs after its repair: source 1000 rows, shadow 999 rows") +} + +// fixedClock is a Clock that always reads the same instant. +type fixedClock time.Time + +func (c fixedClock) Now() time.Time { return time.Time(c) } diff --git a/pkg/checksum/proofs.go b/pkg/checksum/proofs.go index f10d46f..0760e91 100644 --- a/pkg/checksum/proofs.go +++ b/pkg/checksum/proofs.go @@ -8,12 +8,24 @@ import ( // VerifiedShadow proves a full checksum pass found source and shadow equal. // Its zero value is forgeable; consumers must reject a proof whose Table is empty. +// Only a Check that compared every key, found no difference, and repaired +// nothing mints one (CO-1, CO-2). type VerifiedShadow struct { schema, table, shadow string watermark copier.Watermark verifiedAt time.Time } +func newVerifiedShadow(shadow copier.Shadow, watermark copier.Watermark, verifiedAt time.Time) VerifiedShadow { + return VerifiedShadow{ + schema: shadow.Schema(), + table: shadow.SourceTable(), + shadow: shadow.ShadowTable(), + watermark: watermark, + verifiedAt: verifiedAt, + } +} + // Schema returns the verified source schema. func (v VerifiedShadow) Schema() string { return v.schema } @@ -31,8 +43,13 @@ func (v VerifiedShadow) VerifiedAt() time.Time { return v.verifiedAt } // CleanWatermark proves every chunk through its watermark was clean on a fresh // read and that the pass repaired nothing. Its zero value is forgeable; -// consumers must reject it when Watermark().Valid() is false. +// consumers must reject it when Watermark().Valid() is false. Only a Check +// that found no difference and repaired nothing mints one (CO-2). type CleanWatermark struct{ watermark copier.Watermark } +func newCleanWatermark(watermark copier.Watermark) CleanWatermark { + return CleanWatermark{watermark: watermark} +} + // Watermark returns the clean copied-through watermark. func (w CleanWatermark) Watermark() copier.Watermark { return w.watermark } diff --git a/pkg/checksum/repair.go b/pkg/checksum/repair.go new file mode 100644 index 0000000..312a2cb --- /dev/null +++ b/pkg/checksum/repair.go @@ -0,0 +1,124 @@ +package checksum + +import ( + "context" + "fmt" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/block/pg-sprite/pkg/copier" + "github.com/block/pg-sprite/pkg/preflight" +) + +// Repair is one chunk a pass recopied under DivergenceRepair: the +// difference it found, and what the recopy removed and put back. The recopy +// has committed; whether the chunk then read equal is the pass's result, +// not the Repair's — a chunk that did not is named by a RepairError +// alongside the Repairs that committed with it. +type Repair struct { + // Mismatch is the difference the pass found before the repair. + Mismatch Mismatch + // Removed is the number of shadow rows the repair deleted from the chunk. + Removed int64 + // Inserted is the number of source rows the repair copied back. + Inserted int64 +} + +// repairAll recopies every differing chunk in one transaction, then reads +// each one again in a fresh snapshot, in key order, and stops at the first +// chunk whose recopy did not take. The repairs it returns are the ones that +// committed: none when the recopy failed, all of them with any error the +// rereads raise. Every transaction runs under the lock session's Bind +// context, as the pass did. +func (v *Verifier) repairAll(ctx context.Context, pool *pgxpool.Pool, mismatches []Mismatch) ([]Repair, error) { + ctx, unbind := v.lock.Bind(ctx) + defer unbind() + sourceSQL, shadowSQL, err := v.digestStatements(ctx, pool) + if err != nil { + return nil, v.lostOr(err) + } + repairs, err := v.recopy(ctx, pool, mismatches) + if err != nil { + return nil, v.lostOr(err) + } + for _, repair := range repairs { + if err := v.reread(ctx, pool, sourceSQL, shadowSQL, repair); err != nil { + return repairs, v.lostOr(err) + } + } + if lost := v.lockLost(); lost != nil { + return repairs, lost + } + return repairs, nil +} + +// recopy replaces the shadow's rows for every differing chunk in one +// guarded read-write transaction: delete every chunk's rows, then run the +// copier's chunk statement for every chunk, and commit the whole set +// together. Every delete runs before any insert because the source's +// unique indexes stand on the shadow too: when the source moved a unique +// value from a row in one differing chunk to a row in another while the +// shadow stood still, recopying chunk by chunk would insert the value's new +// row while the shadow still held its old one, and the index would refuse +// the recopy on every pass. One transaction also means a half-repaired +// shadow — rows gone and not yet put back — is never visible to another +// transaction. The recopy is the copier's own statement, not a copy of it, +// so a repair puts back exactly what a copy would have. +func (v *Verifier) recopy(ctx context.Context, pool *pgxpool.Pool, mismatches []Mismatch) ([]Repair, error) { + tx, err := v.begin(ctx, pool, pgx.TxOptions{AccessMode: pgx.ReadWrite}) + if err != nil { + return nil, err + } + defer func() { + // Redundant safety closer: after a successful Commit this returns + // the guaranteed ErrTxClosed; on a failure path the server aborts + // the transaction with its session either way. + _ = tx.Rollback(context.WithoutCancel(ctx)) + }() + repairs := make([]Repair, 0, len(mismatches)) + for _, mismatch := range mismatches { + chunk := mismatch.Chunk + removed, err := tx.Exec(ctx, v.repairSQL, chunk.Lower(), chunk.Upper()) + if err != nil { + return nil, fmt.Errorf("clear chunk [%d, %d] of shadow %s.%s for repair: %w", chunk.Lower(), chunk.Upper(), v.shadow.Schema(), v.shadow.ShadowTable(), err) + } + repairs = append(repairs, Repair{Mismatch: mismatch, Removed: removed.RowsAffected()}) + } + for i := range repairs { + chunk := repairs[i].Mismatch.Chunk + inserted, err := tx.Exec(ctx, v.copySQL, chunk.Lower(), chunk.Upper()) + if err != nil { + return nil, fmt.Errorf("recopy chunk [%d, %d] of %s.%s into %s: %w", chunk.Lower(), chunk.Upper(), v.shadow.Schema(), v.shadow.SourceTable(), v.shadow.ShadowTable(), err) + } + repairs[i].Inserted = inserted.RowsAffected() + } + if err := tx.Commit(ctx); err != nil { + return nil, fmt.Errorf("commit repair of %d chunks of %s.%s: %w", len(repairs), v.target.Schema(), v.target.Table(), err) + } + return repairs, nil +} + +// reread digests the repaired chunk in a fresh snapshot and refuses a chunk +// that still differs. +func (v *Verifier) reread(ctx context.Context, pool *pgxpool.Pool, sourceSQL, shadowSQL string, repair Repair) error { + chunk := repair.Mismatch.Chunk + source, shadow, err := v.digestChunk(ctx, pool, sourceSQL, shadowSQL, chunk.Lower(), chunk.Upper()) + if err != nil { + return err + } + if source != shadow { + // INV: CO-2 + return &RepairError{Repair: repair, After: Mismatch{Chunk: chunk, Source: source, Shadow: shadow}} + } + return nil +} + +// repairSQL is the one statement every repair runs before the recopy: +// delete every shadow row whose key lies in the closed range [$1, $2]. The +// bounds are declared bigint as the copy statement's are, so the +// primary-key index serves the range. +func repairSQL(target preflight.CopySwapTarget, shadow copier.Shadow) string { + return "DELETE FROM " + pgx.Identifier{shadow.Schema(), shadow.ShadowTable()}.Sanitize() + + " WHERE " + pgx.Identifier{target.PKColumn()}.Sanitize() + " BETWEEN $1::bigint AND $2::bigint" +} diff --git a/pkg/checksum/repair_integration_test.go b/pkg/checksum/repair_integration_test.go new file mode 100644 index 0000000..a0227d8 --- /dev/null +++ b/pkg/checksum/repair_integration_test.go @@ -0,0 +1,146 @@ +package checksum_test + +import ( + "context" + "math" + "strings" + "sync" + "testing" + + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/checksum" + "github.com/block/pg-sprite/pkg/copier" +) + +// When a later chunk's repair does not take, the chunks recopied with it +// are still reported: every recopy committed in the one repair +// transaction, so the Outcome lists both repairs while the error names the +// one whose fresh read still differed. The trigger zeroes only keys above +// 1000, so the first chunk's recopy takes and the second's does not. +func TestCheckReportsCommittedRepairsWhenALaterRepairFails(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "DELETE FROM "+f.shadowName(shadow)+" WHERE id = 5") + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1500") + f.exec(t, ` + CREATE FUNCTION %s.zero_qty() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + NEW.qty := 0; + RETURN NEW; + END + $$`) + f.exec(t, "CREATE TRIGGER zero_qty BEFORE INSERT ON "+f.shadowName(shadow)+" FOR EACH ROW WHEN (NEW.id > 1000) EXECUTE FUNCTION %s.zero_qty()") + + outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + var failed *checksum.RepairError + require.ErrorAs(t, err, &failed) + assert.Equal(t, chunk(t, 1001, 2000), failed.Repair.Mismatch.Chunk) + + require.Len(t, outcome.Report.Mismatches, 2, "the comparison that led to the repairs is reported") + require.Len(t, outcome.Repairs, 2, "both recopies committed and both are reported") + assert.Equal(t, chunk(t, math.MinInt64, 1000), outcome.Repairs[0].Mismatch.Chunk) + assert.Equal(t, int64(999), outcome.Repairs[0].Removed) + assert.Equal(t, int64(1000), outcome.Repairs[0].Inserted) + assert.Equal(t, failed.Repair, outcome.Repairs[1], "the repair the error names is the second one reported") + _, exists := f.shadowRow(t, shadow, 5) + assert.True(t, exists, "the first chunk's repair committed: row 5 is back") + assertNoProofs(t, outcome) +} + +// Under DivergenceRepair the shadow converges when the source moved a +// unique value from a row in one chunk to a row in another while the +// shadow stood still, the shape slot loss leaves behind. The shadow +// carries the source's unique index, so recopying one chunk at a time +// would insert key 100 with the value key 1500 still holds in the shadow +// and the index would refuse it on every pass; deleting both chunks before +// putting either back lets both land. +func TestCheckRepairsAUniqueValueThatMovedBetweenChunks(t *testing.T) { + f := newVerifierFixture(t) + f.createOrders(t) + f.exec(t, "CREATE UNIQUE INDEX orders_qty ON %s.orders (qty)") + target := f.prove(t, "orders") + lock := f.lock(t, "orders") + shadow := f.build(t, lock, target, `ALTER TABLE %s DROP COLUMN note`) + f.copy(t, target, shadow, lock) + f.exec(t, "UPDATE %s.orders SET qty = -1 WHERE id = 100") + f.exec(t, "UPDATE %s.orders SET qty = 100 WHERE id = 1500") + f.exec(t, "UPDATE %s.orders SET qty = 1500 WHERE id = 100") + + outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + require.NoError(t, err) + require.Len(t, outcome.Repairs, 2) + assert.Equal(t, chunk(t, math.MinInt64, 1000), outcome.Repairs[0].Mismatch.Chunk) + assert.Equal(t, chunk(t, 1001, 2000), outcome.Repairs[1].Mismatch.Chunk) + qty, _ := f.shadowRow(t, shadow, 100) + assert.Equal(t, int64(1500), qty, "key 100 holds the value that moved to it") + qty, _ = f.shadowRow(t, shadow, 1500) + assert.Equal(t, int64(100), qty, "key 1500 holds the value that moved to it") + assertNoProofs(t, outcome) + + again, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + require.NoError(t, err) + assert.True(t, again.Clean(), "the fresh pass finds the shadow equal to the source") +} + +// A shadow replaced between the comparison pass and the repair is refused +// by the repair transaction's own guard before it writes (ST-6): the +// impostor that now carries the shadow's name receives no rows, the +// displaced shadow keeps the difference the pass found, and nothing is +// reported as repaired. The repair is the pass's only read-write +// transaction, so its begin is where the swap is injected. +func TestCheckRepairRefusesAReplacedShadowBeforeWriting(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1500") + name := f.shadowName(shadow) + pool := f.poolHookedBeforeStatement(t, "begin read write", func() { + f.exec(t, "ALTER TABLE "+name+" RENAME TO displaced") + f.exec(t, "CREATE TABLE "+name+" (LIKE %s.displaced INCLUDING ALL)") + }) + + outcome, err := f.check(t, pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + require.ErrorIs(t, err, checksum.ErrInvariantViolation) + assert.Contains(t, err.Error(), "(ST-6)") + assert.Empty(t, outcome.Repairs, "nothing was recopied") + var impostorRows int64 + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT count(*) FROM "+name).Scan(&impostorRows)) + assert.Zero(t, impostorRows, "the repair wrote nothing into the impostor") + var displacedQty int64 + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT qty FROM "+pgx.Identifier{f.schema, "displaced"}.Sanitize()+" WHERE id = 1500").Scan(&displacedQty)) + assert.Equal(t, int64(0), displacedQty, "the displaced shadow was not repaired either") +} + +// A source row in a differing chunk that changes inside the pass — after +// the recopy read it and before the fresh read compares it — makes the +// repaired chunk differ again, and the pass refuses it as a RepairError +// rather than repair it twice: a repair pass assumes the source rows it +// recopies hold still, which the reconciliation after slot loss arranges +// by running before change capture resumes. The recopy itself took and is +// reported; the write is injected as the recopy's insert completes. +func TestCheckRepairRefusesASourceWriteInsideThePass(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1500") + var once sync.Once + pool := f.hookedPool(t, func(_ context.Context, _ *pgx.Conn, sql string) { + if strings.HasPrefix(sql, "INSERT INTO") { + once.Do(func() { f.exec(t, "UPDATE %s.orders SET qty = qty + 1 WHERE id = 1999") }) + } + }) + + outcome, err := f.check(t, pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + var failed *checksum.RepairError + require.ErrorAs(t, err, &failed) + assert.Equal(t, chunk(t, 1001, 2000), failed.After.Chunk) + assert.Equal(t, int64(1000), failed.After.Source.Rows) + assert.Equal(t, int64(1000), failed.After.Shadow.Rows, "the rows are all there; one value differs") + assert.Equal(t, []checksum.Repair{failed.Repair}, outcome.Repairs, "the recopy committed and is reported") + qty, _ := f.shadowRow(t, shadow, 1500) + assert.Equal(t, int64(1500), qty, "the recopy took") + qty, _ = f.shadowRow(t, shadow, 1999) + assert.Equal(t, int64(1999), qty, "the shadow holds the row as the recopy read it, not the later write") + assertNoProofs(t, outcome) +} diff --git a/pkg/checksum/repair_lock_integration_test.go b/pkg/checksum/repair_lock_integration_test.go new file mode 100644 index 0000000..a19c915 --- /dev/null +++ b/pkg/checksum/repair_lock_integration_test.go @@ -0,0 +1,85 @@ +package checksum_test + +import ( + "context" + "math" + "strings" + "sync" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/checksum" + "github.com/block/pg-sprite/pkg/copier" + "github.com/block/pg-sprite/pkg/dbconn" +) + +// Losing the table lock while a repair is in flight cancels the repair's +// transaction and Check reports the loss as an invariant violation (LK-1), +// not the cancelled statement; the half-done repair rolls back with its +// transaction, so the shadow is as the pass found it. The loss is injected +// as the repair's first statement — the shadow DELETE — starts, after the +// comparison pass has already run clean through the lock. +func TestCheckAbortsWhenTheLockIsLostMidRepair(t *testing.T) { + f := newVerifierFixture(t) + f.createOrders(t) + target := f.prove(t, "orders") + buildLock := f.lock(t, "orders") + shadow := f.build(t, buildLock, target, `ALTER TABLE %s DROP COLUMN note`) + f.copy(t, target, shadow, buildLock) + require.NoError(t, buildLock.Release(t.Context())) + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1500") + + lock := f.lock(t, "orders", dbconn.WithTableLockKeepalive(100*time.Millisecond)) + pool := f.poolHookedBeforeStatement(t, "DELETE FROM", func() { + f.terminateBackend(t, lock.BackendPID()) + const lockLossDeadline = 15 * time.Second + select { + case <-lock.Done(): + case <-time.After(lockLossDeadline): + t.Errorf("lock session did not report loss within %s", lockLossDeadline) + } + }) + + _, err := f.check(t, pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + require.ErrorIs(t, err, checksum.ErrInvariantViolation) + assert.ErrorIs(t, err, lock.Err(), "the loss the session reported is the cause") + assert.Contains(t, err.Error(), "(LK-1)") + qty, _ := f.shadowRow(t, shadow, 1500) + assert.Equal(t, int64(0), qty, "the cancelled repair rolled back") +} + +// poolHookedBeforeStatement is a pool whose connections run hook once, as the first +// statement starting with prefix begins. It is a test-only pool: production +// code connects through dbconn. +func (f verifierFixture) poolHookedBeforeStatement(t *testing.T, prefix string, hook func()) *pgxpool.Pool { + t.Helper() + pc, err := pgxpool.ParseConfig(f.cfg.URL) + require.NoError(t, err) + pc.ConnConfig.Tracer = &statementHook{prefix: prefix, hook: hook} + pool, err := pgxpool.NewWithConfig(t.Context(), pc) + require.NoError(t, err) + t.Cleanup(pool.Close) + return pool +} + +// statementHook is a pgx query tracer that runs its hook once, before the +// first statement whose text starts with prefix is sent. +type statementHook struct { + prefix string + hook func() + once sync.Once +} + +func (h *statementHook) TraceQueryStart(ctx context.Context, _ *pgx.Conn, data pgx.TraceQueryStartData) context.Context { + if strings.HasPrefix(data.SQL, h.prefix) { + h.once.Do(h.hook) + } + return ctx +} + +func (h *statementHook) TraceQueryEnd(context.Context, *pgx.Conn, pgx.TraceQueryEndData) {} diff --git a/pkg/checksum/repair_sql_integration_test.go b/pkg/checksum/repair_sql_integration_test.go new file mode 100644 index 0000000..c3ba301 --- /dev/null +++ b/pkg/checksum/repair_sql_integration_test.go @@ -0,0 +1,30 @@ +package checksum + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +// The repair's clearing statement is one string: delete the shadow's rows +// over the same closed bigint-typed key range the copy statement inserts +// over, so the recopy that follows puts back exactly the keys this removed. +// Quoting goes through pgx.Identifier, so a key column named like a keyword +// survives. +func TestRepairSQLIsFrozen(t *testing.T) { + f := newProofFixture(t) + f.exec(t, ` + CREATE TABLE %s.orders ( + "select" bigint PRIMARY KEY, + qty integer NOT NULL + )`) + target := f.prove(t, "orders") + shadow := fakeShadow{ + schema: f.schema, source: "orders", shadow: "_pgsprite_orders_new", + sourceOID: 1, shadowOID: 2, + columns: []string{"select", "qty"}, + } + want := `DELETE FROM "` + f.schema + `"."_pgsprite_orders_new"` + + ` WHERE "select" BETWEEN $1::bigint AND $2::bigint` + assert.Equal(t, want, repairSQL(target, shadow)) +} diff --git a/pkg/checksum/verifier.go b/pkg/checksum/verifier.go index 34df9bb..223e2be 100644 --- a/pkg/checksum/verifier.go +++ b/pkg/checksum/verifier.go @@ -11,6 +11,7 @@ import ( "github.com/block/pg-sprite/pkg/copier" "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/internal/chunksql" "github.com/block/pg-sprite/pkg/preflight" "github.com/block/pg-sprite/pkg/progress" ) @@ -84,6 +85,11 @@ type Verifier struct { shadow copier.Shadow lock *dbconn.TableLockSession opts Options + // repairSQL clears one chunk of the shadow before its recopy and + // copySQL puts the source's rows back, the chunk insert the copier + // runs; both are frozen at construction so no pass builds SQL. + repairSQL string + copySQL string } // NewVerifier prepares verification of target against shadow. It refuses a @@ -111,7 +117,11 @@ func NewVerifier(target preflight.CopySwapTarget, shadow copier.Shadow, lock *db if _, err := copier.NewChunker(target, copier.Watermark{}, opts.Chunker); err != nil { return nil, err } - return &Verifier{target: target, shadow: shadow, lock: lock, opts: opts}, nil + return &Verifier{ + target: target, shadow: shadow, lock: lock, opts: opts, + repairSQL: repairSQL(target, shadow), + copySQL: chunksql.Insert(shadow.Schema(), shadow.SourceTable(), shadow.ShadowTable(), target.PKColumn(), shadow.CopyColumns()), + }, nil } // checkShadow refuses a shadow proof that does not describe the proven @@ -182,11 +192,34 @@ func (v *Verifier) Verify(ctx context.Context, pool *pgxpool.Pool, through copie ctx, unbind := v.lock.Bind(ctx) defer unbind() report, err := v.pass(ctx, pool, through) - if lost := v.lock.Err(); lost != nil { - // INV: LK-1 - return Report{}, fmt.Errorf("%w (LK-1): table lock was lost during verification: %w", ErrInvariantViolation, lost) + if err != nil { + return Report{}, v.lostOr(err) + } + if lost := v.lockLost(); lost != nil { + return Report{}, lost + } + return report, nil +} + +// lockLost reports the table lock's loss as the verifier's invariant +// violation, or nil while the session still holds it. +func (v *Verifier) lockLost() error { + lost := v.lock.Err() + if lost == nil { + return nil } - return report, err + // INV: LK-1 + return fmt.Errorf("%w (LK-1): table lock was lost during verification: %w", ErrInvariantViolation, lost) +} + +// lostOr returns the lock's loss when the session has reported one — a read +// cancelled by Bind is explained by the loss, not by its own error — and err +// otherwise. +func (v *Verifier) lostOr(err error) error { + if lost := v.lockLost(); lost != nil { + return lost + } + return err } func (v *Verifier) pass(ctx context.Context, pool *pgxpool.Pool, through copier.Watermark) (Report, error) { @@ -235,7 +268,7 @@ func (v *Verifier) pass(ctx context.Context, pool *pgxpool.Pool, through copier. // shadow's column types, read in a guarded transaction so the shadow whose // types they name is the shadow the proof describes. func (v *Verifier) digestStatements(ctx context.Context, pool *pgxpool.Pool) (sourceSQL, shadowSQL string, err error) { - tx, err := v.begin(ctx, pool) + tx, err := v.begin(ctx, pool, snapshotRead()) if err != nil { return "", "", err } @@ -260,7 +293,7 @@ func (v *Verifier) digestStatements(ctx context.Context, pool *pgxpool.Pool) (so // digestChunk reads both sides of the closed range [lower, upper] in one // guarded transaction, so the two digests describe one snapshot. func (v *Verifier) digestChunk(ctx context.Context, pool *pgxpool.Pool, sourceSQL, shadowSQL string, lower, upper int64) (source, shadow Digest, err error) { - tx, err := v.begin(ctx, pool) + tx, err := v.begin(ctx, pool, snapshotRead()) if err != nil { return Digest{}, Digest{}, err } diff --git a/pkg/copier/chunker.go b/pkg/copier/chunker.go index a8138ec..b5864ff 100644 --- a/pkg/copier/chunker.go +++ b/pkg/copier/chunker.go @@ -182,7 +182,7 @@ func startAfter(from Watermark) (lower int64, done bool) { if !from.Valid() { return math.MinInt64, false } - if from.Value() == math.MaxInt64 { + if from.Complete() { return 0, true } return from.Value() + 1, false diff --git a/pkg/copier/copy_chunk.go b/pkg/copier/copy_chunk.go index 478d918..e874359 100644 --- a/pkg/copier/copy_chunk.go +++ b/pkg/copier/copy_chunk.go @@ -5,13 +5,13 @@ import ( "errors" "fmt" "strconv" - "strings" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgconn" "github.com/jackc/pgx/v5/pgxpool" "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/internal/chunksql" "github.com/block/pg-sprite/pkg/preflight" ) @@ -38,7 +38,7 @@ func (c *Copier) copyChunk(ctx context.Context, pool *pgxpool.Pool, chunk Chunk) } tag, err := tx.Exec(ctx, c.sql, chunk.Lower(), chunk.Upper()) if err != nil { - return 0, fmt.Errorf("copy chunk [%d, %d] of %s.%s into %s: %w", chunk.Lower(), chunk.Upper(), c.target.Schema(), c.target.Table(), c.shadow.ShadowTable(), err) + return 0, fmt.Errorf("copy chunk [%d, %d] of %s.%s into %s: %w", chunk.Lower(), chunk.Upper(), c.shadow.Schema(), c.shadow.SourceTable(), c.shadow.ShadowTable(), err) } if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("commit chunk [%d, %d] of %s.%s: %w", chunk.Lower(), chunk.Upper(), c.target.Schema(), c.target.Table(), err) @@ -47,12 +47,13 @@ func (c *Copier) copyChunk(ctx context.Context, pool *pgxpool.Pool, chunk Chunk) } // guard prepares the transaction every write into the shadow runs in: it -// bounds it, puts it under the source owner's role, and confirms from this -// connection that the lock session's backend still holds the table and that -// the source and shadow are still the relations the proofs describe. Both -// relations are locked by name before their identity is checked, so nothing -// can replace either one between the check and the statements that follow -// it in the same transaction. +// bounds it, puts the catalog alone on its search_path, puts it under the +// source owner's role, and confirms from this connection that the lock +// session's backend still holds the table and that the source and shadow +// are still the relations the proofs describe. Both relations are locked by +// name before their identity is checked, so nothing can replace either one +// between the check and the statements that follow it in the same +// transaction. func (c *Copier) guard(ctx context.Context, tx pgx.Tx) error { if err := setCopySession(ctx, tx, c.target, c.opts); err != nil { return err @@ -66,13 +67,20 @@ func (c *Copier) guard(ctx context.Context, tx pgx.Tx) error { return c.confirmRelations(ctx, tx) } -// setCopySession bounds the transaction and puts it in the owner's shoes. -// Every identifier the copy touches is schema-qualified, so the session's -// search_path plays no part. +// setCopySession bounds the transaction, restricts its search_path to the +// catalog, and puts it in the owner's shoes. Every relation the copy names +// is schema-qualified, but the key range's BETWEEN resolves its operators +// through the search_path, and a shadow default or trigger fires under it +// too; pinning the catalog keeps a schema ahead of it on the caller's pool +// from answering either (CO-9), and gives the verifier's repair, which runs +// the same statement under the same pin, the same session as the copy. func setCopySession(ctx context.Context, tx pgx.Tx, target preflight.CopySwapTarget, opts Options) error { if err := setBudgets(ctx, tx, opts); err != nil { return err } + if _, err := tx.Exec(ctx, dbconn.LocalSearchPath("pg_catalog")); err != nil { + return fmt.Errorf("set copy search_path: %w", err) + } if _, err := tx.Exec(ctx, "SET LOCAL ROLE "+pgx.Identifier{target.OwnerRole()}.Sanitize()); err != nil { return fmt.Errorf("set owner role %s: %w", target.OwnerRole(), err) } @@ -160,10 +168,10 @@ func (c *Copier) confirmRelations(ctx context.Context, tx pgx.Tx) error { } // relationOIDsSQL resolves the source ($2) and shadow ($3) in schema $1 by -// explicit qualification; a missing relation scans as NULL. The transaction -// sets no search_path of its own, so every catalog name is qualified: a -// decoy pg_class ahead of the catalog on the session's path must not answer -// the identity check (CO-9). +// explicit qualification; a missing relation scans as NULL. Every catalog +// name in it is qualified as well as pinned by the session: a decoy +// pg_class ahead of the catalog on the caller's path must not answer the +// identity check (CO-9). const relationOIDsSQL = ` SELECT (SELECT c.oid @@ -185,29 +193,8 @@ func confirmRelation(role, schema, table string, proven uint32, found *uint32) e return nil } -// copySQL is the one statement every chunk runs: insert the shared columns -// of the source rows whose key lies in the closed range [$1, $2] into the -// shadow, skipping any key the shadow already holds — the applier always -// overwrites and the copier never does, which is what lets the two run -// concurrently (CO-4). The bounds are declared bigint whatever the key's -// integer type, as the chunker's boundary query declares them, so a bound -// outside a smaller key type's range can still be sent and the primary-key -// index still serves the range scan. The statement carries no conversion -// expression: a column whose type differs between the two tables is -// converted by the server's assignment cast, so a type change that needs a -// USING expression cannot be copied by this statement and must not be routed -// to the copier until it can carry one. +// copySQL is the one statement every chunk runs, the chunk insert shared +// with the verifier's repair, built for this target and shadow. func copySQL(target preflight.CopySwapTarget, shadow Shadow) string { - columns := make([]string, 0, len(shadow.CopyColumns())) - for _, column := range shadow.CopyColumns() { - columns = append(columns, pgx.Identifier{column}.Sanitize()) - } - list := strings.Join(columns, ", ") - key := pgx.Identifier{target.PKColumn()}.Sanitize() - return "INSERT INTO " + pgx.Identifier{shadow.Schema(), shadow.ShadowTable()}.Sanitize() + - " (" + list + ")" + - " SELECT " + list + - " FROM " + pgx.Identifier{shadow.Schema(), shadow.SourceTable()}.Sanitize() + - " WHERE " + key + " BETWEEN $1::bigint AND $2::bigint" + - " ON CONFLICT (" + key + ") DO NOTHING" + return chunksql.Insert(shadow.Schema(), shadow.SourceTable(), shadow.ShadowTable(), target.PKColumn(), shadow.CopyColumns()) } diff --git a/pkg/copier/copy_chunk_search_path_integration_test.go b/pkg/copier/copy_chunk_search_path_integration_test.go new file mode 100644 index 0000000..21723d9 --- /dev/null +++ b/pkg/copier/copy_chunk_search_path_integration_test.go @@ -0,0 +1,40 @@ +package copier_test + +import ( + "math" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/internal/testutil" + "github.com/block/pg-sprite/pkg/copier" +) + +// Every chunk transaction pins the catalog alone on its search_path (CO-9), +// so the copy statement's key range resolves BETWEEN to the catalog's +// operators whatever the caller's pool puts ahead of pg_catalog. The schema +// first on this pool's path offers a bigint <= that is never true; under the +// session's own path every chunk's range would be empty, and the copy +// would land nothing and still report the whole key space as copied. +// Functions in the statement are not at stake — it names none — so the +// operator is the one name only the transaction's search_path can pin. +func TestCopierIgnoresTheSessionSearchPath(t *testing.T) { + f := newCopierFixture(t) + const rows = 300 + target, lock, shadow := f.prepare(t, rows) + f.exec(t, `CREATE FUNCTION %s.never_le(bigint, bigint) RETURNS boolean LANGUAGE sql IMMUTABLE AS 'SELECT false'`) + f.exec(t, `CREATE OPERATOR %s.<= (LEFTARG = bigint, RIGHTARG = bigint, FUNCTION = %s.never_le)`) + shadowing := testutil.NewCatalogShadowingPool(t, f.cfg.URL, f.schema) + + c, err := copier.NewCopier(target, shadow, lock, copier.Watermark{}, copier.Options{ + Workers: 2, + Chunker: copier.ChunkerOptions{InitialRows: 100, MaxRows: 100}, + }) + require.NoError(t, err) + require.NoError(t, c.Run(t.Context(), shadowing)) + + f.assertConverged(t, shadow) + assert.Equal(t, int64(rows), c.Position().RowsInserted, "every row landed through the catalog's operators") + assert.Equal(t, copier.NewWatermark(math.MaxInt64), c.Position().Watermark) +} diff --git a/pkg/copier/types.go b/pkg/copier/types.go index 44806f2..7bbe4ea 100644 --- a/pkg/copier/types.go +++ b/pkg/copier/types.go @@ -1,6 +1,9 @@ package copier -import "fmt" +import ( + "fmt" + "math" +) // Chunk is a closed range of a single-column integer primary key. Its fields // are unexported so that lower <= upper holds for every value the chunker @@ -54,3 +57,8 @@ func (w Watermark) Valid() bool { return w.valid } // Value returns the frontier's primary key. It is zero for an invalid // watermark; callers check Valid first. func (w Watermark) Value() int64 { return w.value } + +// Complete reports whether the frontier is past every key an int64 can +// hold: the chunker's final open-above chunk has landed, so no key remains +// on the other side of the watermark. +func (w Watermark) Complete() bool { return w.valid && w.value == math.MaxInt64 } diff --git a/pkg/copier/types_test.go b/pkg/copier/types_test.go index 2816840..a334ffe 100644 --- a/pkg/copier/types_test.go +++ b/pkg/copier/types_test.go @@ -1,6 +1,7 @@ package copier import ( + "math" "testing" "github.com/stretchr/testify/assert" @@ -24,3 +25,13 @@ func TestWatermarkStates(t *testing.T) { assert.Equal(t, int64(0), w.Value()) assert.Equal(t, int64(41), NewWatermark(41).Value()) } + +// A watermark is complete only at the top of the key space, where the +// chunker's final open-above chunk ends; the zero watermark and every +// finite frontier below it leave keys on the other side. +func TestWatermarkComplete(t *testing.T) { + assert.True(t, NewWatermark(math.MaxInt64).Complete()) + assert.False(t, NewWatermark(math.MaxInt64-1).Complete(), "one key remains") + assert.False(t, NewWatermark(0).Complete()) + assert.False(t, Watermark{}.Complete(), "nothing copied is not everything copied") +} diff --git a/pkg/internal/chunksql/chunksql.go b/pkg/internal/chunksql/chunksql.go new file mode 100644 index 0000000..9bfcb4a --- /dev/null +++ b/pkg/internal/chunksql/chunksql.go @@ -0,0 +1,41 @@ +// Package chunksql holds the one statement that puts a chunk of source rows +// into the shadow. The copier runs it for every chunk it lands and the +// verifier runs it for every chunk it repairs; one text in one place is what +// makes a repair put back exactly what a copy would have, and keeps the +// statement off the public API, where a caller could run it outside the +// guard both packages wrap around it. +package chunksql + +import ( + "strings" + + "github.com/jackc/pgx/v5" +) + +// Insert is the statement that copies one chunk: insert the shared columns +// of the source rows whose key lies in the closed range [$1, $2] into the +// shadow, skipping any key the shadow already holds — the applier always +// overwrites and the copier never does, which is what lets the two run +// concurrently (CO-4). The bounds are declared bigint whatever the key's +// integer type, as the chunker's boundary query declares them, so a bound +// outside a smaller key type's range can still be sent and the primary-key +// index still serves the range scan. The statement carries no conversion +// expression: a column whose type differs between the two tables is +// converted by the server's assignment cast, so a type change that needs a +// USING expression cannot be copied by this statement and must not be +// routed to the copier until it can carry one. Every identifier is quoted +// and the relations are schema-qualified. +func Insert(schema, source, shadow, key string, columns []string) string { + quoted := make([]string, 0, len(columns)) + for _, column := range columns { + quoted = append(quoted, pgx.Identifier{column}.Sanitize()) + } + list := strings.Join(quoted, ", ") + pk := pgx.Identifier{key}.Sanitize() + return "INSERT INTO " + pgx.Identifier{schema, shadow}.Sanitize() + + " (" + list + ")" + + " SELECT " + list + + " FROM " + pgx.Identifier{schema, source}.Sanitize() + + " WHERE " + pk + " BETWEEN $1::bigint AND $2::bigint" + + " ON CONFLICT (" + pk + ") DO NOTHING" +} diff --git a/pkg/internal/chunksql/chunksql_test.go b/pkg/internal/chunksql/chunksql_test.go new file mode 100644 index 0000000..9b15c7b --- /dev/null +++ b/pkg/internal/chunksql/chunksql_test.go @@ -0,0 +1,19 @@ +package chunksql + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +// The chunk insert is the CO-4 contract in one string: the shared columns, +// a closed bigint-typed key range, and ON CONFLICT DO NOTHING so the +// statement never overwrites what the applier wrote. Quoting goes through +// pgx.Identifier, so a column named like a keyword survives. +func TestInsertIsFrozen(t *testing.T) { + want := `INSERT INTO "app"."_pgsprite_orders_new" ("id", "select", "qty")` + + ` SELECT "id", "select", "qty" FROM "app"."orders"` + + ` WHERE "id" BETWEEN $1::bigint AND $2::bigint` + + ` ON CONFLICT ("id") DO NOTHING` + assert.Equal(t, want, Insert("app", "orders", "_pgsprite_orders_new", "id", []string{"id", "select", "qty"})) +}