diff --git a/SAFETY.md b/SAFETY.md index 881d565..0632c39 100644 --- a/SAFETY.md +++ b/SAFETY.md @@ -20,7 +20,7 @@ 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, 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/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 pass is the progress tracker's `WorkSource` while it runs (chunks compared, rows hashed, chunks mismatched, chunks repaired, all from memory); 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 | @@ -34,7 +34,7 @@ The invariant registry (invariant IDs referenced below) lives in | `pkg/diffplan` — desired schema → routed convergence plan, the declarative front door as a library (the CLI `diff` and embedding orchestrators share it) | ❌ periphery | exists | — | | `pkg/migrate` — one gated statement → resolve, classify, route, execute → one verdict; the imperative front door as a library (the CLI `migrate` and embedding orchestrators share it), plus the desired-state execution loop (`RunDesired`: derive the convergence plan, admit it as a whole, run each planned statement back through the same pipeline) | ❌ periphery² | exists | — | | `internal/cli` — CLI, flags, help, prompts | ❌ periphery | `migrate`, `pull`, `diff`, `fmt`, `lint`, `suggest`, `capabilities`, and `status` exist | — | -| `pkg/progress` — strategy-wide progress snapshots; the executors' observation seam (core imports it, so its locking discipline is core-critical); the `WorkSource` seam for engine-measured steps such as the copy | ✅ core | native progress and the work-source seam exist; the copier fills it | — | +| `pkg/progress` — strategy-wide progress snapshots; the executors' observation seam (core imports it, so its locking discipline is core-critical); the `WorkSource` seam for engine-measured steps such as the copy and the checksum pass | ✅ core | native progress and the work-source seam exist; the copier and the checksum verifier fill it | — | | orchestrator adapter | ❌ periphery | planned (Phase 11) | OC-* hold *at* the boundary | | `internal/testutil` | ❌ test-only | exists | — | diff --git a/docs/architecture.md b/docs/architecture.md index fa547fb..b5a9047 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -199,9 +199,9 @@ different levels of commitment: | `pkg/migrate` | The imperative front door as a library: one parsed statement in — gate, resolve, classify, route, execute — one `verdict.Verdict` out; the CLI `migrate` and embedding orchestrators share this one pipeline. Also the desired-state execution loop: `RunDesired` derives the convergence plan (`diffplan.Plan`), admits it as a whole (existence, destructive guard, dispositions, optional fingerprint pin), and runs each planned statement back through `Run` — per-statement verdicts, committed-prefix semantics | exists | | `pkg/router` | Route classified statements to native / copy-and-swap / refuse dispositions; copy-and-swap reports unavailable until that backend lands | exists (Phase 2.4) | | `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/progress` | Strategy-wide, pollable progress snapshots: native phase/elapsed time, sequence position, retry attempt, and server-reported concurrent-index work, and engine-measured work polled from the copier and the checksum verifier (`WorkSource`) | native progress and the copy and checksum work sources 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), `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/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; the pass reports its chunks compared, rows hashed, mismatches, and repairs to the progress tracker as its `WorkSource`; continuous checker to follow | verifier, policy, repair, proofs, and progress counters 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 738d4bb..34663fb 100644 --- a/docs/copy-and-swap-design.md +++ b/docs/copy-and-swap-design.md @@ -387,7 +387,7 @@ 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 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/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. While `Verify` or `Check` runs — through the repair phase as well — the verifier is the tracker's `progress.WorkSource` (`Options.Tracker`): `chunks_compared` and `rows_hashed` count the two-sided digests that have committed and the source rows they covered (a repaired chunk's reread counts again), `chunks_mismatched` the chunks the comparison found differing, and `chunks_repaired` the chunks whose one recopy transaction has committed — every counter read from the pass's memory, none from the database, and all reset when a pass starts. | 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/progress-report.md b/docs/progress-report.md index de21a62..b950246 100644 --- a/docs/progress-report.md +++ b/docs/progress-report.md @@ -4,9 +4,9 @@ The progress snapshot is the machine-readable observation a caller receives when running schema change through the `*WithProgress` executor entry points. It is the one JSON shape an operator or orchestrator consumes to display or act on execution progress. This document is the contract: the fields, the closed vocabularies, and the behavior required of -a consumer. The Go source of truth is `pkg/progress`; `TestSnapshotJSONShape` and -`TestSnapshotJSONShapeForACopyStep` pin the exact keys, including the two examples at the end -of this page. +a consumer. The Go source of truth is `pkg/progress`; `TestSnapshotJSONShape`, +`TestSnapshotJSONShapeForACopyStep` and `TestSnapshotJSONShapeForAChecksumStep` pin the exact +keys, including the three examples at the end of this page. ## Versioning: `format_version` @@ -17,11 +17,13 @@ phase or operation value is a contract change and bumps `format_version`, even i is added or renamed. Adding a field bumps `format_version` so a strict consumer can detect the new shape from the -version. The current version is **4**: version 2 added `detail.statement`; version 3 added +version. The current version is **5**: version 2 added `detail.statement`; version 3 added `detail.current_locker_pid`, `work.lockers_total`, and `work.lockers_done`; version 4 added the `copy` operation and split `work` into two counter families — the server-observed build counters and the engine-measured copy counters (`rows_*`, `bytes_*`) — selected by -`detail.operation`, so `work` is no longer a signal that a concurrent index build is running. +`detail.operation`, so `work` is no longer a signal that a concurrent index build is running; +version 5 added the `checksum` operation and its engine-measured counter family +(`chunks_compared`, `rows_hashed`, `chunks_mismatched`, `chunks_repaired`). The [plan report](plan-report.md), [lint report](lint-report.md), and [suggest report](suggest-report.md) are separate contracts with their own `format_version`; @@ -56,23 +58,25 @@ licenses a consumer to intervene in the change itself. | `active` | bool | always | Whether an operation is executing now. `false` with `phase: "running"` means a concurrent build's progress row has left the server view. | | `attempt` | int | bounded retries only | The current attempt number when the executor is inside its bounded retry loop. | | `current_locker_pid` | int | while waiting on a locker | PostgreSQL backend PID currently blocking the concurrent build; omitted when none is published. | -| `work` | object | measured work only | Present exactly when something measured the step's work — the server published a progress row for a concurrent build, or the engine reported the copy step's counters; then **every** counter below is present, so a fresh build or an empty copy reports honest zeros rather than an empty object. Which counters mean anything is decided by `operation`, not by `work` being present — see [Work counters](#work-counters). | +| `work` | object | measured work only | Present exactly when something measured the step's work — the server published a progress row for a concurrent build, or the engine reported the copy or checksum step's counters; then **every** counter below is present, so a fresh build or an empty copy reports honest zeros rather than an empty object. Which counters mean anything is decided by `operation`, not by `work` being present — see [Work counters](#work-counters). | `statement` is the SQL the engine is running for the step: for a native operation the submitter's statement after qualification and canonicalization, for a `copy` step the engine's own frozen chunk insert (the template with its `$1`/`$2` key bounds, not a chunk's rendered values). Either way it is real SQL that reached the server, so a consumer rendering -it into a shared surface must clamp and escape it. +it into a shared surface must clamp and escape it. A `checksum` step runs several frozen +statements per chunk, so the orchestrator that owns the step chooses what, if anything, to +put in `statement`. ### Work counters -Two operations publish `work`, and each measures only its own counters; the other -operation's counters are `0`, never estimated. A consumer selects the counter family from +Three operations publish `work`, and each measures only its own counters; the other +operations' counters are `0`, never estimated. A consumer selects the counter family from `detail.operation`, never from the presence of `work`: a present `work` says only that something measured the step, and a consumer that reads its presence as "a concurrent index build is running" will show a `copy` step as a build with zero blocks. Render the build -counters for `concurrent-index-build`, the copy counters for `copy`, and nothing from `work` -for an operation you do not recognize. +counters for `concurrent-index-build`, the copy counters for `copy`, the checksum counters +for `checksum`, and nothing from `work` for an operation you do not recognize. **Server-observed** (`concurrent-index-build`): `blocks_done` / `blocks_total`, `tuples_done` / `tuples_total`, and `lockers_done` / `lockers_total` come from @@ -105,6 +109,26 @@ The size read runs in a read-only transaction of the copy's own, under the copy' that timeout whatever session defaults the caller's pool carries, so an observer never holds the copy's stop path open. +**Engine-measured** (`checksum`): `chunks_compared`, `rows_hashed`, `chunks_mismatched` and +`chunks_repaired` come from the checksum pass itself, which cuts its own chunks and digests +each one on both tables inside one snapshot. The counters cover a whole `Check`: the +comparison that finds differing chunks and, under the `repair` policy, the recopy and the +reread that follow it, so a pass that spends most of its time repairing still shows movement. +They are read from the pass's memory — no catalog read, nothing for a poll to wait on — and +reset to zero when a pass starts. Native and `copy` operations report none of them. + +| Counter | Meaning | +| --- | --- | +| `chunks_compared` | Two-sided chunk digests this pass has completed: every chunk of the comparison once, and every repaired chunk once more when its reread completes. It therefore passes the chunk count of the comparison during the repair phase. | +| `rows_hashed` | Source rows those digests covered, summed the same way: a repaired chunk's rows count again at its reread. | +| `chunks_mismatched` | Chunks the comparison found differing. Under `abort` this is the finding the pass returns with; under `repair` it is the number of chunks the recopy covers. | +| `chunks_repaired` | Chunks whose recopy from the source has committed. Every differing chunk is recopied in one transaction, so this moves from `0` to `chunks_mismatched` at that commit and `chunks_compared` then advances as each repaired chunk is reread. A chunk still differing at its reread stops the pass; it stays counted here because its recopy did commit. | + +There is no chunk total: a pass sizes its chunks from the time each one takes, so the count +of chunks is known only when the pass ends. A consumer that wants a completion figure for the +comparison phase has the copy step's `rows_total` from the previous step and this step's +`rows_hashed`. + ## Phases | Value | Meaning | @@ -127,13 +151,14 @@ returns the identical snapshot, elapsed values included. | `validate-constraint` | A constraint-validation scan. | | `concurrent-index-build` | A concurrent index build (`work` is server-observed). | | `copy` | The copy-and-swap row copy from the source table into its shadow (`work` is engine-measured). | +| `checksum` | The copy-and-swap checksum pass comparing the source table with its shadow and, under the `repair` policy, recopying differing chunks (`work` is engine-measured). | ## Polling semantics The tracker is caller-owned and has no goroutines or timers: polling lifetime is exactly the caller's context. A poll during an active concurrent index build performs one read of the -server's progress view over the executor's reserved session; a poll during a `copy` step -asks the engine's work source once; every other poll is pure memory. On a query or source +server's progress view over the executor's reserved session; a poll during a `copy` or +`checksum` step asks the engine's work source once; every other poll is pure memory. On a query or source error the returned snapshot still carries the last-known tracker state — `phase` is never empty — with the error returned alongside for the caller to classify. Pollers serialize against each other, so the reserved session and the work source each see one observation @@ -175,7 +200,7 @@ A poll during step 2 of a 3-step sequence, mid concurrent index build: ```json { - "format_version": 4, + "format_version": 5, "phase": "running", "step": 2, "total_steps": 3, @@ -193,6 +218,10 @@ A poll during step 2 of a 3-step sequence, mid concurrent index build: "rows_total": 0, "bytes_copied": 0, "bytes_total": 0, + "chunks_compared": 0, + "rows_hashed": 0, + "chunks_mismatched": 0, + "chunks_repaired": 0, "blocks_done": 11, "blocks_total": 40, "tuples_done": 7, @@ -209,7 +238,7 @@ pins it): ```json { - "format_version": 4, + "format_version": 5, "phase": "running", "step": 2, "total_steps": 4, @@ -224,6 +253,45 @@ pins it): "rows_total": 5000, "bytes_copied": 98304, "bytes_total": 409600, + "chunks_compared": 0, + "rows_hashed": 0, + "chunks_mismatched": 0, + "chunks_repaired": 0, + "blocks_done": 0, + "blocks_total": 0, + "tuples_done": 0, + "tuples_total": 0, + "lockers_total": 0, + "lockers_done": 0 + } + } +} +``` + +A poll during step 3 of a 4-step copy-and-swap, mid checksum pass under the `repair` policy — +three chunks compared, two found differing and recopied, one of the two reread so far +(`TestSnapshotJSONShapeForAChecksumStep` pins it): + +```json +{ + "format_version": 5, + "phase": "running", + "step": 3, + "total_steps": 4, + "elapsed_ns": 2750000000, + "step_elapsed_ns": 750000000, + "detail": { + "operation": "checksum", + "active": true, + "work": { + "rows_copied": 0, + "rows_total": 0, + "bytes_copied": 0, + "bytes_total": 0, + "chunks_compared": 4, + "rows_hashed": 3500, + "chunks_mismatched": 2, + "chunks_repaired": 2, "blocks_done": 0, "blocks_total": 0, "tuples_done": 0, diff --git a/pkg/checksum/check.go b/pkg/checksum/check.go index d5c2b27..cc6d2c0 100644 --- a/pkg/checksum/check.go +++ b/pkg/checksum/check.go @@ -68,11 +68,15 @@ func (o Outcome) VerifiedShadow() (VerifiedShadow, bool) { // 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. +// +// While it runs, through the repair phase as well, the verifier is the +// tracker's work source (Options.Tracker). 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) + defer v.report()() + report, err := v.verify(ctx, pool, through) if err != nil { return Outcome{}, err } diff --git a/pkg/checksum/repair.go b/pkg/checksum/repair.go index 312a2cb..688254b 100644 --- a/pkg/checksum/repair.go +++ b/pkg/checksum/repair.go @@ -96,6 +96,7 @@ func (v *Verifier) recopy(ctx context.Context, pool *pgxpool.Pool, mismatches [] 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) } + v.countRepaired(len(repairs)) return repairs, nil } diff --git a/pkg/checksum/verifier.go b/pkg/checksum/verifier.go index 223e2be..fa627d4 100644 --- a/pkg/checksum/verifier.go +++ b/pkg/checksum/verifier.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "slices" + "sync" "time" "github.com/jackc/pgx/v5/pgxpool" @@ -44,6 +45,12 @@ type Options struct { // Clock times each chunk for the chunker's chunk-time sizing feedback // (docs/copy-and-swap-design.md#d12--throttle-by-chunk-time-and-slot-lag). Clock progress.Clock + // Tracker, when set, is told the pass's work for the lifetime of Verify + // or Check: the verifier is its progress.WorkSource from before the + // first chunk until just before the call returns, through the repair + // phase as well. The caller owns the tracker's steps; the verifier only + // fills the current step's counters. + Tracker *progress.Tracker } func (o Options) withDefaults() Options { @@ -78,8 +85,9 @@ func (o Options) validate() error { // its own short read-only transaction whose one snapshot covers both // tables, so a chunk's two digests describe the same instant and no pass // holds a snapshot open for longer than one chunk. It runs only under the -// table's lock session. A Verifier holds no state between passes; the same -// one can run a pass after every repair. +// table's lock session. The only state a Verifier keeps between passes is +// the counters of the pass in progress, which the next pass resets, so the +// same one runs a pass after every repair, one pass at a time. type Verifier struct { target preflight.CopySwapTarget shadow copier.Shadow @@ -90,6 +98,10 @@ type Verifier struct { // runs; both are frozen at construction so no pass builds SQL. repairSQL string copySQL string + + // mu guards work, the counters Work reports for the pass in progress. + mu sync.Mutex + work progress.Work } // NewVerifier prepares verification of target against shadow. It refuses a @@ -185,7 +197,15 @@ func requireTableLock(lock *dbconn.TableLockSession, target preflight.CopySwapTa // an error, and the caller's divergence policy decides what to do with it. // Every transaction runs under the lock session's Bind context, so losing // the table lock cancels the read in flight and Verify reports the loss. +// While it runs the verifier is the tracker's work source (Options.Tracker). func (v *Verifier) Verify(ctx context.Context, pool *pgxpool.Pool, through copier.Watermark) (Report, error) { + defer v.report()() + return v.verify(ctx, pool, through) +} + +// verify is the pass Verify and Check share, run with the work source +// already registered by the caller. +func (v *Verifier) verify(ctx context.Context, pool *pgxpool.Pool, through copier.Watermark) (Report, error) { if !through.Valid() { return Report{}, fmt.Errorf("%w: %s.%s", ErrNothingLanded, v.target.Schema(), v.target.Table()) } @@ -254,6 +274,7 @@ func (v *Verifier) pass(ctx context.Context, pool *pgxpool.Pool, through copier. return Report{}, fmt.Errorf("%w (CO-1): verified chunk: %w", ErrInvariantViolation, err) } report.Mismatches = append(report.Mismatches, Mismatch{Chunk: compared, Source: source, Shadow: shadow}) + v.countMismatch() } if upper == through.Value() { return report, nil @@ -314,5 +335,6 @@ func (v *Verifier) digestChunk(ctx context.Context, pool *pgxpool.Pool, sourceSQ if err := tx.Commit(ctx); err != nil { return Digest{}, Digest{}, fmt.Errorf("commit digest of chunk [%d, %d] of %s.%s: %w", lower, upper, v.target.Schema(), v.target.Table(), err) } + v.countCompared(source.Rows) return source, shadow, nil } diff --git a/pkg/checksum/work.go b/pkg/checksum/work.go new file mode 100644 index 0000000..5ca6c6e --- /dev/null +++ b/pkg/checksum/work.go @@ -0,0 +1,70 @@ +package checksum + +import ( + "context" + + "github.com/block/pg-sprite/pkg/progress" +) + +// Work reports the pass's counters for the progress tracker: it is the +// progress.WorkSource the verifier registers for the lifetime of Verify or +// Check. Every counter is read from memory — the pass knows what it has +// digested and recopied — so a poll never waits on the database, and the +// counters reset to zero when a pass starts. +// +// - chunks_compared is the number of two-sided chunk digests the pass has +// completed: every chunk of the comparison once, and every repaired +// chunk once more when its reread completes. +// - rows_hashed is the number of source rows those digests covered, summed +// the same way. +// - chunks_mismatched is the number of chunks the comparison found +// differing. +// - chunks_repaired is the number of chunks whose recopy has committed. +// Every differing chunk is recopied in one transaction, so it moves from +// zero to chunks_mismatched at that commit. +func (v *Verifier) Work(context.Context) (progress.Work, error) { + v.mu.Lock() + defer v.mu.Unlock() + return v.work, nil +} + +// report resets the counters for a new pass and registers the verifier with +// the tracker, when there is one. The returned stop is the fence before the +// pass returns: it waits for an in-flight poll, so no poll that began while +// the pass ran completes against a verifier whose caller has moved on. +func (v *Verifier) report() (stop func()) { + v.mu.Lock() + v.work = progress.Work{} + v.mu.Unlock() + if v.opts.Tracker != nil { + v.opts.Tracker.SetWorkSource(v) + } + return func() { + if v.opts.Tracker != nil { + v.opts.Tracker.StopWorkSource() + } + } +} + +// countCompared records one completed two-sided digest covering rows +// source rows. +func (v *Verifier) countCompared(rows int64) { + v.mu.Lock() + defer v.mu.Unlock() + v.work.ChunksCompared++ + v.work.RowsHashed += uint64(rows) +} + +// countMismatch records one chunk the comparison found differing. +func (v *Verifier) countMismatch() { + v.mu.Lock() + defer v.mu.Unlock() + v.work.ChunksMismatched++ +} + +// countRepaired records chunks whose recopy has committed. +func (v *Verifier) countRepaired(chunks int) { + v.mu.Lock() + defer v.mu.Unlock() + v.work.ChunksRepaired += uint64(chunks) +} diff --git a/pkg/checksum/work_integration_test.go b/pkg/checksum/work_integration_test.go new file mode 100644 index 0000000..804a250 --- /dev/null +++ b/pkg/checksum/work_integration_test.go @@ -0,0 +1,179 @@ +package checksum_test + +import ( + "context" + "math" + "sync" + "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/progress" +) + +// checksumTracker is a tracker inside a running checksum step, the state in +// which an orchestrator polls a pass. +func checksumTracker(t *testing.T) *progress.Tracker { + t.Helper() + tracker, err := progress.NewTracker(progress.WallClock{}) + require.NoError(t, err) + tracker.Start(1, progress.OperationChecksum) + tracker.StartStep(1, progress.OperationChecksum, "") + return tracker +} + +// workObserver polls the tracker after every statement a pass runs and +// keeps each poll's counters in order, so a test can ask what an observer +// saw at each point of the pass rather than only at its end. +type workObserver struct { + tracker *progress.Tracker + mu sync.Mutex + seen []progress.Work +} + +// pool is a pool whose every statement is followed by one poll of the +// tracker, on the pass's own goroutine. +func (o *workObserver) pool(t *testing.T, f verifierFixture) *pgxpool.Pool { + t.Helper() + return f.hookedPool(t, func(ctx context.Context, _ *pgx.Conn, _ string) { + snapshot, err := o.tracker.Progress(ctx) + require.NoError(t, err) + require.NotNil(t, snapshot.Detail.Work, "every statement of a pass runs while the verifier is the work source") + o.mu.Lock() + defer o.mu.Unlock() + o.seen = append(o.seen, *snapshot.Detail.Work) + }) +} + +// snapshots returns every poll so far, oldest first. +func (o *workObserver) snapshots() []progress.Work { + o.mu.Lock() + defer o.mu.Unlock() + return append([]progress.Work(nil), o.seen...) +} + +// forget drops the polls recorded so far, so the next pass is observed on +// its own. +func (o *workObserver) forget() { + o.mu.Lock() + defer o.mu.Unlock() + o.seen = nil +} + +// first returns the earliest poll that satisfies want, and whether one did. +func (o *workObserver) first(want func(progress.Work) bool) (progress.Work, bool) { + for _, w := range o.snapshots() { + if want(w) { + return w, true + } + } + return progress.Work{}, false +} + +// A repair pass over a shadow that differs in every chunk is visible to a +// poller stage by stage: each chunk's digest counts the source's rows the +// moment it commits, the three mismatches are counted before the repair +// transaction opens, the repair moves chunks_repaired to three at its one +// commit, and the rereads count as compared chunks again. Once Check +// returns the tracker no longer asks the verifier, whose own counters +// still hold the whole pass. +func TestCheckReportsEachStageOfARepairPass(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.corrupt(t, shadow) + observer := &workObserver{tracker: checksumTracker(t)} + pool := observer.pool(t, f) + + v, err := checksum.NewVerifier(target, shadow, lock, checksum.Options{Chunker: chunkRows, Tracker: observer.tracker}) + require.NoError(t, err) + outcome, err := v.Check(t.Context(), pool, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair) + require.NoError(t, err) + require.Len(t, outcome.Repairs, 3) + + firstChunk, ok := observer.first(func(w progress.Work) bool { return w.ChunksCompared == 1 }) + require.True(t, ok, "a poll lands between the first chunk's commit and the second's first statement") + assert.Equal(t, uint64(1000), firstChunk.RowsHashed, "the first chunk's digest covers the source's 1000 keys, not the shadow's 999") + assert.Equal(t, uint64(1), firstChunk.ChunksMismatched, "the first chunk's missing row is counted with its digest") + assert.Zero(t, firstChunk.ChunksRepaired) + + compared, ok := observer.first(func(w progress.Work) bool { return w.ChunksCompared == 3 }) + require.True(t, ok, "the repair transaction's first statement is polled with the comparison complete") + assert.Equal(t, progress.Work{ChunksCompared: 3, RowsHashed: 2500, ChunksMismatched: 3}, compared, "every chunk differs and none has been repaired yet") + + repaired, ok := observer.first(func(w progress.Work) bool { return w.ChunksRepaired > 0 }) + require.True(t, ok, "the rereads are polled after the repair committed") + assert.Equal(t, progress.Work{ChunksCompared: 3, RowsHashed: 2500, ChunksMismatched: 3, ChunksRepaired: 3}, repaired, "one transaction repairs all three chunks, so the count moves from zero to three at once") + + after, err := observer.tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Nil(t, after.Detail.Work, "a returned pass is no longer the tracker's work source") + direct, err := v.Work(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.Work{ChunksCompared: 6, RowsHashed: 5000, ChunksMismatched: 3, ChunksRepaired: 3}, direct, "three chunks compared, then the same three reread after their repair") +} + +// A clean Verify pass counts its three chunks and the source's rows and +// nothing else, and a second pass on the same verifier starts from zero: +// the first poll of the second pass sees no chunk compared, so a poller +// never reads the previous pass's total as this pass's progress. The +// pass's final figure is read from the verifier after it returns: the last +// statement of a pass is the last chunk's commit, polled before that chunk +// is counted. +func TestVerifyCountsAPassAndTheNextPassStartsFromZero(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + observer := &workObserver{tracker: checksumTracker(t)} + pool := observer.pool(t, f) + + v, err := checksum.NewVerifier(target, shadow, lock, checksum.Options{Chunker: chunkRows, Tracker: observer.tracker}) + require.NoError(t, err) + report, err := v.Verify(t.Context(), pool, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + require.True(t, report.Clean()) + + direct, err := v.Work(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.Work{ChunksCompared: 3, RowsHashed: 2500}, direct) + after, err := observer.tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Nil(t, after.Detail.Work, "a returned pass is no longer the tracker's work source") + + observer.forget() + _, err = v.Verify(t.Context(), pool, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + seen := observer.snapshots() + require.NotEmpty(t, seen) + assert.Equal(t, progress.Work{}, seen[0], "the second pass's first statement is polled with every counter reset") + direct, err = v.Work(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.Work{ChunksCompared: 3, RowsHashed: 2500}, direct, "the second pass ends with only its own chunks counted") +} + +// Under DivergenceAbort the pass counts what it found and repairs nothing: +// chunks_mismatched is the only divergence signal an observer gets, so it +// must move even when the policy refuses the repair. +func TestCheckUnderAbortCountsMismatchesAndNoRepairs(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.corrupt(t, shadow) + observer := &workObserver{tracker: checksumTracker(t)} + pool := observer.pool(t, f) + + v, err := checksum.NewVerifier(target, shadow, lock, checksum.Options{Chunker: chunkRows, Tracker: observer.tracker}) + require.NoError(t, err) + _, err = v.Check(t.Context(), pool, copier.NewWatermark(math.MaxInt64), checksum.DivergenceAbort) + var divergence *checksum.DivergenceError + require.ErrorAs(t, err, &divergence) + + direct, err := v.Work(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.Work{ChunksCompared: 3, RowsHashed: 2500, ChunksMismatched: 3}, direct, "the comparison ran to the end; the policy stopped the pass before any repair") + after, err := observer.tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Nil(t, after.Detail.Work, "a refused pass has released the tracker too") +} diff --git a/pkg/checksum/work_test.go b/pkg/checksum/work_test.go new file mode 100644 index 0000000..59ede3b --- /dev/null +++ b/pkg/checksum/work_test.go @@ -0,0 +1,88 @@ +package checksum + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/progress" +) + +// runningTracker is a tracker inside a running checksum step, the state in +// which an orchestrator polls a pass. +func runningTracker(t *testing.T) *progress.Tracker { + t.Helper() + tracker, err := progress.NewTracker(progress.WallClock{}) + require.NoError(t, err) + tracker.Start(1, progress.OperationChecksum) + tracker.StartStep(1, progress.OperationChecksum, "") + return tracker +} + +// The counters accumulate what the pass tells them: compared chunks carry +// their source row counts, a mismatch and a repair count chunks, and +// nothing else moves. A fresh verifier reports all zeros. +func TestWorkAccumulatesTheCountersAPassReports(t *testing.T) { + v := &Verifier{} + work, err := v.Work(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.Work{}, work, "a verifier that has run no pass has counted nothing") + + v.countCompared(1000) + v.countCompared(500) + v.countMismatch() + v.countRepaired(2) + + work, err = v.Work(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.Work{ + ChunksCompared: 2, + RowsHashed: 1500, + ChunksMismatched: 1, + ChunksRepaired: 2, + }, work) +} + +// report starts a pass: it zeroes the previous pass's counters and makes +// the verifier the tracker's work source until stop runs, after which a +// poll carries no engine work and the counters stay readable from the +// verifier itself. +func TestReportResetsTheCountersAndRegistersForThePass(t *testing.T) { + tracker := runningTracker(t) + v := &Verifier{opts: Options{Tracker: tracker}} + v.countCompared(1000) + v.countMismatch() + + stop := v.report() + v.countCompared(250) + + polled, err := tracker.Progress(t.Context()) + require.NoError(t, err) + require.NotNil(t, polled.Detail.Work, "while a pass runs a poll asks the verifier") + assert.Equal(t, progress.Work{ChunksCompared: 1, RowsHashed: 250}, *polled.Detail.Work, "the previous pass's chunk and mismatch were reset at report") + + stop() + after, err := tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Nil(t, after.Detail.Work, "a pass that returned is no longer the tracker's work source") + direct, err := v.Work(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.Work{ChunksCompared: 1, RowsHashed: 250}, direct, "the finished pass's counters survive until the next report") +} + +// A verifier with no tracker still counts: report registers nowhere and +// stop has nothing to fence, so a library caller that does not observe +// progress pays nothing for it. +func TestReportWithoutATrackerOnlyResetsTheCounters(t *testing.T) { + v := &Verifier{} + v.countRepaired(3) + + stop := v.report() + v.countCompared(10) + stop() + + work, err := v.Work(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.Work{ChunksCompared: 1, RowsHashed: 10}, work) +} diff --git a/pkg/progress/progress.go b/pkg/progress/progress.go index 26a7bae..943799a 100644 --- a/pkg/progress/progress.go +++ b/pkg/progress/progress.go @@ -1,8 +1,9 @@ // Package progress defines the strategy-wide, machine-readable execution // progress contract. One snapshot shape serves every operation: a concurrent // index build's counters are read from the server's progress view, the -// copy-and-swap row copy's counters come from the engine's own WorkSource, -// and each operation leaves the other's counters at zero. +// copy-and-swap row copy's and checksum pass's counters come from the +// engine's own WorkSource, and each operation leaves the others' counters +// at zero. package progress import ( @@ -48,7 +49,7 @@ const ( // field semantics. Adding a phase or operation value is a contract change // and bumps this version, even when no field is added or renamed. Adding a // field also bumps this version so strict consumers can detect the new shape. -const FormatVersion = 4 +const FormatVersion = 5 // Operation is the current operation's execution class. type Operation string @@ -68,6 +69,10 @@ const ( // OperationCopy is the copy-and-swap row copy from the source table into // its shadow. OperationCopy Operation = "copy" + // OperationChecksum is the copy-and-swap checksum pass comparing the + // source table with its shadow chunk by chunk, and repairing differing + // chunks when its policy says so. + OperationChecksum Operation = "checksum" ) // Work reports observed work. It is present only when something measured @@ -75,19 +80,24 @@ const ( // engine's own WorkSource reported the step's counters — and then every // counter marshals explicitly — a fresh build reports honest zeros, never an // empty object a consumer must guess at. Rows and bytes belong to the -// copy-and-swap row copy; blocks, tuples and lockers to a concurrent index -// build. Neither operation fabricates the other's counters. +// copy-and-swap row copy; chunks compared, rows hashed, chunks mismatched +// and chunks repaired to the checksum pass; blocks, tuples and lockers to a +// concurrent index build. No operation fabricates another's counters. type Work struct { - RowsCopied uint64 `json:"rows_copied"` - RowsTotal uint64 `json:"rows_total"` - BytesCopied uint64 `json:"bytes_copied"` - BytesTotal uint64 `json:"bytes_total"` - BlocksDone uint64 `json:"blocks_done"` - BlocksTotal uint64 `json:"blocks_total"` - TuplesDone uint64 `json:"tuples_done"` - TuplesTotal uint64 `json:"tuples_total"` - LockersTotal uint64 `json:"lockers_total"` - LockersDone uint64 `json:"lockers_done"` + RowsCopied uint64 `json:"rows_copied"` + RowsTotal uint64 `json:"rows_total"` + BytesCopied uint64 `json:"bytes_copied"` + BytesTotal uint64 `json:"bytes_total"` + ChunksCompared uint64 `json:"chunks_compared"` + RowsHashed uint64 `json:"rows_hashed"` + ChunksMismatched uint64 `json:"chunks_mismatched"` + ChunksRepaired uint64 `json:"chunks_repaired"` + BlocksDone uint64 `json:"blocks_done"` + BlocksTotal uint64 `json:"blocks_total"` + TuplesDone uint64 `json:"tuples_done"` + TuplesTotal uint64 `json:"tuples_total"` + LockersTotal uint64 `json:"lockers_total"` + LockersDone uint64 `json:"lockers_done"` } // Detail describes the operation currently executing. diff --git a/pkg/progress/progress_test.go b/pkg/progress/progress_test.go index c8de4e5..1f1c2d1 100644 --- a/pkg/progress/progress_test.go +++ b/pkg/progress/progress_test.go @@ -207,7 +207,7 @@ func TestSnapshotJSONShape(t *testing.T) { raw, err := json.Marshal(snapshot) require.NoError(t, err) assert.JSONEq(t, `{ - "format_version": 4, + "format_version": 5, "phase": "running", "step": 2, "total_steps": 3, @@ -225,6 +225,10 @@ func TestSnapshotJSONShape(t *testing.T) { "rows_total": 0, "bytes_copied": 0, "bytes_total": 0, + "chunks_compared": 0, + "rows_hashed": 0, + "chunks_mismatched": 0, + "chunks_repaired": 0, "blocks_done": 11, "blocks_total": 40, "tuples_done": 7, @@ -266,7 +270,7 @@ func TestSnapshotJSONOmitsUnsetOptionalFields(t *testing.T) { raw, err := json.Marshal(snapshot) require.NoError(t, err) assert.JSONEq(t, `{ - "format_version": 4, + "format_version": 5, "phase": "pending", "elapsed_ns": 0, "step_elapsed_ns": 0, diff --git a/pkg/progress/work_source.go b/pkg/progress/work_source.go index 18e4418..c957a88 100644 --- a/pkg/progress/work_source.go +++ b/pkg/progress/work_source.go @@ -5,7 +5,8 @@ import "context" // WorkSource reports engine-derived work for the current step: the counters // of an operation whose progress PostgreSQL publishes no view for, such as // the copy-and-swap row copy, which knows its own rows copied from the -// chunks it has landed. The tracker polls it inside Progress, on the +// chunks it has landed, or the checksum pass, which knows the chunks it has +// compared and repaired. The tracker polls it inside Progress, on the // observer's context, so the source must be safe to call from any goroutine // while the step runs; the tracker itself never calls it from more than one // goroutine at a time. Every counter the source does not measure stays zero diff --git a/pkg/progress/work_source_test.go b/pkg/progress/work_source_test.go index 155118c..3a3c8ec 100644 --- a/pkg/progress/work_source_test.go +++ b/pkg/progress/work_source_test.go @@ -54,9 +54,9 @@ func runningTrackerWithSource(t *testing.T, source progress.WorkSource) *progres return tracker } -// The copy step's JSON is the adapter-facing contract for format_version 4: -// the copy operation value, and work carrying the engine's rows and bytes -// with the build counters at honest zero. +// The copy step's JSON is the adapter-facing contract for the copy +// operation: its operation value, and work carrying the engine's rows and +// bytes with the checksum and build counters at honest zero. func TestSnapshotJSONShapeForACopyStep(t *testing.T) { var polls atomic.Int32 clock := &fakeClock{now: time.Unix(100, 0)} @@ -74,7 +74,7 @@ func TestSnapshotJSONShapeForACopyStep(t *testing.T) { raw, err := json.Marshal(snapshot) require.NoError(t, err) assert.JSONEq(t, `{ - "format_version": 4, + "format_version": 5, "phase": "running", "step": 2, "total_steps": 4, @@ -89,6 +89,10 @@ func TestSnapshotJSONShapeForACopyStep(t *testing.T) { "rows_total": 5000, "bytes_copied": 98304, "bytes_total": 409600, + "chunks_compared": 0, + "rows_hashed": 0, + "chunks_mismatched": 0, + "chunks_repaired": 0, "blocks_done": 0, "blocks_total": 0, "tuples_done": 0, @@ -101,6 +105,58 @@ func TestSnapshotJSONShapeForACopyStep(t *testing.T) { assert.Equal(t, int32(1), polls.Load(), "one poll asks the source once") } +// checksumCounters are distinct per field so a swapped pair cannot pass: +// a repair pass that compared three chunks, found two differing, recopied +// both, and has reread one of them so far. +var checksumCounters = progress.Work{ChunksCompared: 4, RowsHashed: 3500, ChunksMismatched: 2, ChunksRepaired: 2} + +// The checksum step's JSON is the adapter-facing contract for the checksum +// operation: its operation value, and work carrying the engine's chunk and +// row counters with the copy and build counters at honest zero. +func TestSnapshotJSONShapeForAChecksumStep(t *testing.T) { + clock := &fakeClock{now: time.Unix(100, 0)} + tracker, err := progress.NewTracker(clock) + require.NoError(t, err) + tracker.Start(4, progress.OperationAdmitting) + clock.now = clock.now.Add(2 * time.Second) + tracker.StartStep(3, progress.OperationChecksum, "") + tracker.SetWorkSource(fakeSource{work: func(context.Context) (progress.Work, error) { return checksumCounters, nil }}) + clock.now = clock.now.Add(750 * time.Millisecond) + + snapshot, err := tracker.Progress(t.Context()) + require.NoError(t, err) + raw, err := json.Marshal(snapshot) + require.NoError(t, err) + assert.JSONEq(t, `{ + "format_version": 5, + "phase": "running", + "step": 3, + "total_steps": 4, + "elapsed_ns": 2750000000, + "step_elapsed_ns": 750000000, + "detail": { + "operation": "checksum", + "active": true, + "work": { + "rows_copied": 0, + "rows_total": 0, + "bytes_copied": 0, + "bytes_total": 0, + "chunks_compared": 4, + "rows_hashed": 3500, + "chunks_mismatched": 2, + "chunks_repaired": 2, + "blocks_done": 0, + "blocks_total": 0, + "tuples_done": 0, + "tuples_total": 0, + "lockers_total": 0, + "lockers_done": 0 + } + } + }`, string(raw)) +} + // Each poll asks the source afresh, so the counters a consumer sees are the // source's current observation, not a value captured at registration. func TestProgressPollsTheWorkSourceOnEveryObservation(t *testing.T) {