From 55f79b013933f89f1acf83200dcf8b92d7e4e394 Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Thu, 1 Oct 2026 21:34:26 +1000 Subject: [PATCH] progress, checksum: make the verifier the tracker's work source and add the checksum counter family (format_version 5) A checksum pass was dark to an observer: the tracker could report a running step but nothing about how far the comparison or the repair had got. The verifier is now the tracker's WorkSource for the lifetime of Verify or Check (Options.Tracker), as the copier is for the copy, and reports a third counter family read from the pass's own memory: chunks_compared and rows_hashed (two-sided digests that committed and the source rows they covered; a repaired chunk's reread counts again), chunks_mismatched (chunks the comparison found differing), and chunks_repaired (chunks whose one recopy transaction committed). The counters reset when a pass starts, and Check registers once so the repair phase is covered without a second registration. Adding fields and the `checksum` operation to the snapshot bumps format_version to 5. Tests: unit tests for the counters, the reset at report, registration and release; integration tests polling the tracker after every statement of a repair pass (each chunk counted with the source's row count the moment it commits, all three mismatches counted before the repair opens, repaired moving 0 -> 3 at the one commit, the rereads counted again), of a clean pass and a second pass on the same verifier (first poll sees every counter reset), and of an abort pass (mismatches counted, nothing repaired); the progress JSON shape tests pin the version and the zero counters, and a new one pins a checksum step. Docs: progress-report.md gains the version-5 note, the checksum counter family, the checksum operation row, and a checksum-step example; the copy-and-swap package map, architecture.md and SAFETY.md name the verifier as a work source. --- SAFETY.md | 4 +- docs/architecture.md | 4 +- docs/copy-and-swap-design.md | 2 +- docs/progress-report.md | 98 +++++++++++--- pkg/checksum/check.go | 6 +- pkg/checksum/repair.go | 1 + pkg/checksum/verifier.go | 26 +++- pkg/checksum/work.go | 70 ++++++++++ pkg/checksum/work_integration_test.go | 179 ++++++++++++++++++++++++++ pkg/checksum/work_test.go | 88 +++++++++++++ pkg/progress/progress.go | 40 +++--- pkg/progress/progress_test.go | 8 +- pkg/progress/work_source.go | 3 +- pkg/progress/work_source_test.go | 64 ++++++++- 14 files changed, 548 insertions(+), 45 deletions(-) create mode 100644 pkg/checksum/work.go create mode 100644 pkg/checksum/work_integration_test.go create mode 100644 pkg/checksum/work_test.go 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) {