From 53ec5545f7ea5c94d01ef43ea6d6362fbe1f8c2c Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Thu, 1 Oct 2026 09:09:23 +1000 Subject: [PATCH 1/3] checksum: chunk verifier comparing source and shadow in one snapshot (CO-1) Verifier digests every chunk up to the copier's landed watermark on both sides inside one read-only REPEATABLE READ transaction, casting every column to the shadow's type (D7), under the copier's guard: owner role, catalog-only search_path, ACCESS SHARE on both relations before the snapshot, lock confirmation, relation-OID check. It reports the chunks that differ; policy, repair, and the proof constructors follow. --- SAFETY.md | 2 +- docs/architecture.md | 2 +- docs/copy-and-swap-design.md | 5 +- docs/invariants.md | 18 +- pkg/checksum/digest.go | 112 ++++++ pkg/checksum/digest_integration_test.go | 150 +++++++ pkg/checksum/doc.go | 13 +- pkg/checksum/guard.go | 156 ++++++++ pkg/checksum/report.go | 36 ++ pkg/checksum/verifier.go | 285 ++++++++++++++ pkg/checksum/verifier_integration_test.go | 451 ++++++++++++++++++++++ pkg/checksum/verifier_test.go | 59 +++ 12 files changed, 1281 insertions(+), 8 deletions(-) create mode 100644 pkg/checksum/digest.go create mode 100644 pkg/checksum/digest_integration_test.go create mode 100644 pkg/checksum/guard.go create mode 100644 pkg/checksum/report.go create mode 100644 pkg/checksum/verifier.go create mode 100644 pkg/checksum/verifier_integration_test.go create mode 100644 pkg/checksum/verifier_test.go diff --git a/SAFETY.md b/SAFETY.md index 9c8e7a3..e6d8239 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 declaration exists | ST-6, RF-1..RF-5 | | `pkg/executor` — bounded optimistic attempt; native concurrent index build with invalid-index recovery; native sequence executor for the safer idioms | ✅ core | exists (Phase 1: attempt-under-budget; Phase 3.1: concurrent index build; Phase 3.2: sequence executor) | LK-2 (attempt bound + the CONCURRENTLY wait-policy exception), CO-9 (qualified proof reads), ST-9 (create owner verified, never repaired) | -| `pkg/checksum` — chunk verifier, continuous checker, repair | ✅ core | types and proof-type declarations exist; verifier planned | CO-1, CO-2, CO-3 | +| `pkg/checksum` — chunk verifier, continuous checker, repair | ✅ core | the chunk `Verifier` exists (one read-only `REPEATABLE READ` transaction per chunk so both digests share a snapshot, every column cast to the shadow's type on both sides, the same guard as the copier — owner role, catalog-only `search_path`, `ACCESS SHARE` on both relations, lock confirmation, relation-OID check — and a `Report` of mismatched chunks up to the landed watermark); divergence policy, repair, the proof constructors, and the continuous checker planned | CO-1 (compare; the gate lands with the proof constructors), CO-2, CO-3, CO-9, LK-1 | | `pkg/copier` — shadow-table chunked copy | ✅ core | contract types and the keyset `Chunker` exist (row-count chunks over the proven key, first chunk open below and last open above, a cut frontier for the applier's discard rule, time-targeted sizing) and the parallel `Copier` (one bounded never-overwriting insert per chunk under the table lock session, frontier-ordered in-flight registry, a resume that first clears the shadow above the watermark, `Position.Classify` for the applier) exist; progress fillers planned | CO-4 (chunk coverage, copy SQL shape, in-flight registry), LK-1, LK-3 | | `pkg/applier` — change apply, buffer, flush scheduling | ✅ core | package contract exists; applier planned | CO-4, CO-5, CO-6, CO-8, LK-3 | | `pkg/decode` — logical decoding, LSN/position accounting, per-column presence | ✅ core | contract types exist; decoder planned | ST-4, CO-4, CO-8 | diff --git a/docs/architecture.md b/docs/architecture.md index 8179f4a..dc9fc8c 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -201,7 +201,7 @@ different levels of commitment: | `pkg/executor` | Native backend with stable outcome codes: the bounded optimistic attempt, the concurrent index build and its invalid-index recovery, the autocommit safer-sequence runner, the greenfield `CREATE TABLE` path, and the accepted-blocking passthrough primitive; the full `Executor` contract (`Plan`/`Execute`/`Status`/`Abort`) arrives with the copy-and-swap backend | native execution exists | | `pkg/progress` | Strategy-wide, pollable progress snapshots: native phase/elapsed time, sequence position, retry attempt, and server-reported concurrent-index work; optional copy counters are reserved for copy-and-swap | native progress exists | | `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) — there is no separate chunker package | chunker and copy loop exist; progress fillers planned | -| `pkg/checksum` | The mandatory correctness gate; continuous checker; repair primitive | Phase 5 | +| `pkg/checksum` | The mandatory correctness gate: the chunk verifier (single-snapshot per-chunk digests through the shadow's types, reporting the chunks that differ up to the landed watermark); divergence policy, repair primitive, proof constructors, and continuous checker to follow | verifier exists; gate, repair, and checker planned | | `pkg/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 31d878a..a82e038 100644 --- a/docs/copy-and-swap-design.md +++ b/docs/copy-and-swap-design.md @@ -187,7 +187,8 @@ narrowing, and precision changes instead of comparing unlike textual representat **Alternative considered → deferred.** Operation-specific checksum SQL multiplies semantic edge cases and test surfaces. -**Where enforced.** `pkg/checksum`; CO-1, CO-2. +**Where enforced.** `pkg/checksum` — the `Verifier`'s digest statement casts every copy column +to the type the shadow declares on both sides (the generated-column half is planned); CO-1, CO-2. ### D8 — Use deterministic bounded names @@ -378,7 +379,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, dependent-object, decoding, and headroom checks above); owns Tier-3 refusals. | ST-6, RF-1..RF-3 | | `pkg/copier` | Produces `Chunk` and `Watermark`; `Chunker` (built only from a `CopySwapTarget`) cuts consecutive chunks that tile the whole int64 key space — first open below, last open above — so every key a row can carry belongs to exactly one chunk and a watermark at the largest value means the copy is complete. `Copier` (built from a `CopySwapTarget`, a `Shadow` — the shape `schemachange.BuiltShadow` satisfies — and the table's `TableLockSession`) copies chunks with several workers, each in its own bounded transaction under the owner's role that confirms the lock, takes `ACCESS SHARE` on both relations, and confirms both relation OIDs before one frozen never-overwriting insert; a resumed copy first deletes every shadow row above the watermark in batches of the same guarded shape — the whole shadow after a zero watermark — so the resumed cut frontier and the shadow agree on what is uncut, and its first batch takes `SHARE MODE` on the shadow so a chunk transaction of the earlier run still committing into it ends before the clear reads; the proof check refuses a shadow that is the source by name or by OID, since the clear deletes from one and the copy reads the other; the copy statement carries no conversion expression, so a column whose type differs between the two tables takes the server's assignment cast, and an `ALTER COLUMN TYPE … USING` change needs a conversion-aware copy before the planner's copy-and-swap route is wired to the copier; `Position` snapshots the cut frontier and landed watermark as two `Watermark` values plus the in-flight chunks, and `Position.Classify` is the applier's uncut / in-flight / landed rule. 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` | Produces `VerifiedShadow` and `CleanWatermark`; their constructors are private to this package. | CO-1, CO-2, CO-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. Both sides run the identical frozen statement — row count plus `md5` of the key-ordered concatenation of each row's `md5(ROW(col::shadow_type, …)::text)` over `pk BETWEEN $1 AND $2` — so only the data can differ, and the cast on every column is D7: a converted column hashes as the value the shadow holds. Only the copy columns are compared; generated columns present on both sides are not yet hashed. A `Report` proves nothing — the divergence policy, repair, and the constructors of `VerifiedShadow` and `CleanWatermark` (private to this package) land on top of it. | CO-1, CO-2, CO-3, CO-9, LK-1 | | `pkg/decode` | Produces `ChangeEvent`, including per-column presence and `OldKey` for an UPDATE that moved the primary key. | ST-3, ST-4, CO-4, CO-8 | | `pkg/applier` | Applies presence-aware events from the per-key buffer. | CO-4, CO-5, CO-6, CO-8, LK-3 | | `pkg/checkpoint` | Produces `Checkpoint`. | ST-1, ST-2 | diff --git a/docs/invariants.md b/docs/invariants.md index 2caabd0..3de4ffd 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -36,7 +36,14 @@ several of these unrepresentable, and the in-TCB engineering rules live in A migration that cannot prove shadow == source **must refuse to cut over**. No flag, mode, or capture mechanism removes the gate; it is also the repair primitive for [slot-loss reconciliation](low-level-design.md#failover-during-migration-what-survives-and-what-doesnt). -*Enforced:* cutover entry condition. *Source:* [design-principles](design-principles.md#correctness-and-safety), +*Enforced:* cutover entry condition. *Enforced today:* `pkg/checksum` `Verifier` — the +comparison the gate will demand: every chunk up to the landed watermark digested on both sides +inside one read-only `REPEATABLE READ` transaction, with every column cast to the shadow's type, +and every differing chunk reported with its two row counts (clean-copy, changed-row, +missing-and-extra-row, converted-type, watermark-clamp, shadowing-`search_path`, +replaced-relation, and lost-lock tests). *Planned enforcement:* the `VerifiedShadow` constructor +accepts only a clean `Report`, and cutover accepts only a `VerifiedShadow`. +*Source:* [design-principles](design-principles.md#correctness-and-safety), risks-and-mitigations; Spirit's "never skip it". ### CO-2 — A persisted checksum watermark describes only chunks verified clean on a fresh read @@ -231,7 +238,9 @@ not the implicit `pg_temp` search ahead of it (pg-sprite creates no temporary ob proxy that hands the server connection to another client keeps the rewritten, stricter path. *Enforced:* `pkg/dbconn` (session hook, `LocalSearchPath`, and the test that keeps it the only `search_path` writer under `pkg/`), `pg_catalog.` qualification in `pkg/executor`, -`pkg/progress`, `pkg/schemadiff`. *Test obligation:* a shadowing `search_path` (`, +`pkg/progress`, `pkg/schemadiff`, and `pkg/checksum` (the column-type read, the relation check, +and every function in the digest statement, under a `LocalSearchPath("pg_catalog")` transaction; +decoy `md5` and `format_type` test). *Test obligation:* a shadowing `search_path` (`, pg_catalog` with decoy catalog relations and functions in the schema) yields the same answer as the default path, per read site and per pooled session. @@ -284,7 +293,10 @@ transaction that the session's backend holds the lock before the first write (ni wrong-table, reported-loss, gone-session, rival-backend, mid-build-loss, and mid-drop-loss tests); `pkg/copier` `Copier` requires the same session, runs every chunk transaction under its `Bind` context, and calls `TableLockSession.Confirm` from each chunk's own connection before -the insert (wrong-table, gone-session, rival-backend, and mid-copy-loss tests). *Planned +the insert (wrong-table, gone-session, rival-backend, and mid-copy-loss tests); `pkg/checksum` +`Verifier` requires the same session, runs every read transaction under its `Bind` context, and +confirms the lock from each transaction's own connection before the first read (wrong-table, +gone-session, reported-loss, and mid-pass-loss tests). *Planned enforcement:* cutover acquires the same session before its first write and runs under it, so loss of the lock aborts the change at every stage. *Source:* Spirit `pkg/dbconn/metadatalock.go` (stated pool invariants). This resolves the diff --git a/pkg/checksum/digest.go b/pkg/checksum/digest.go new file mode 100644 index 0000000..3c5e5e1 --- /dev/null +++ b/pkg/checksum/digest.go @@ -0,0 +1,112 @@ +package checksum + +import ( + "context" + "fmt" + "strings" + + "github.com/jackc/pgx/v5" + + "github.com/block/pg-sprite/pkg/copier" + "github.com/block/pg-sprite/pkg/preflight" +) + +// Digest is one side's summary of a chunk: how many rows the range holds +// and an order-sensitive hash of every one of them. +type Digest struct { + // Rows is the number of rows whose key lies in the chunk. + Rows int64 + // Hash is md5 over the concatenation, in key order, of the per-row md5 + // of the row's copy columns rendered as a record. It is the md5 of the + // empty string for an empty range. + Hash string +} + +// columnType pairs a copy column with the type the shadow declares for it, +// as format_type renders it: the SQL spelling a cast can name. +type columnType struct { + name string + typeName string +} + +// shadowColumnTypesSQL reads the type of every copy column as the shadow +// ($1) declares it. Every catalog name is pg_catalog-qualified so the +// session's search_path plays no part (CO-9), and the session's search_path +// is pg_catalog alone while it runs, so format_type qualifies every type +// that is not built in. +const shadowColumnTypesSQL = ` + SELECT a.attname, pg_catalog.format_type(a.atttypid, a.atttypmod) + FROM pg_catalog.pg_attribute a + WHERE a.attrelid = $1::oid + AND a.attnum > 0 + AND NOT a.attisdropped + AND a.attname::text OPERATOR(pg_catalog.=) ANY($2::text[])` + +// shadowColumnTypes returns the copy columns in the order the shadow proof +// lists them, each with the shadow's type. A copy column the shadow no +// longer has is a shadow that is not the one the proof describes. +func shadowColumnTypes(ctx context.Context, tx pgx.Tx, shadow copier.Shadow) ([]columnType, error) { + rows, err := tx.Query(ctx, shadowColumnTypesSQL, shadow.ShadowOID(), shadow.CopyColumns()) + if err != nil { + return nil, fmt.Errorf("read column types of shadow %s.%s: %w", shadow.Schema(), shadow.ShadowTable(), err) + } + found, err := pgx.CollectRows(rows, func(row pgx.CollectableRow) (columnType, error) { + var c columnType + err := row.Scan(&c.name, &c.typeName) + return c, err + }) + if err != nil { + return nil, fmt.Errorf("read column types of shadow %s.%s: %w", shadow.Schema(), shadow.ShadowTable(), err) + } + byName := make(map[string]string, len(found)) + for _, c := range found { + byName[c.name] = c.typeName + } + types := make([]columnType, 0, len(shadow.CopyColumns())) + for _, column := range shadow.CopyColumns() { + typeName, ok := byName[column] + if !ok { + // INV: ST-6 + return nil, fmt.Errorf("%w (ST-6): shadow %s.%s has no column %s the proof lists for copy", ErrInvariantViolation, shadow.Schema(), shadow.ShadowTable(), column) + } + types = append(types, columnType{name: column, typeName: typeName}) + } + return types, nil +} + +// digestSQL is the one statement run on each side of a chunk, with the +// table it reads the only difference between the two: the row count and +// the md5 of the key-ordered concatenation of every row's md5, each row +// rendered as a record of its copy columns cast to the shadow's types — +// the D7 rule of docs/copy-and-swap-design.md#d7--checksum-through-the-destination-types. +// The source side is where the casts do work: a column whose type the +// schema change widens or narrows hashes as the value the shadow holds, +// and both sides run the identical expression so nothing but the data can +// differ. The record rendering tells NULL from the empty string, and +// hashing per row before aggregating bounds the aggregate's input to one +// md5 per row. The bounds are declared bigint whatever the key's integer +// type, as the chunker's boundary query declares them, so the primary-key +// index serves the range scan. Every function is pg_catalog-qualified, and +// the guarded session's search_path is pg_catalog alone for the operators +// and casts that cannot be, so the session's own search_path plays no part +// (CO-9); COALESCE is syntax, not a function, and takes no qualifier. +func digestSQL(target preflight.CopySwapTarget, schema, table string, types []columnType) string { + cast := make([]string, 0, len(types)) + for _, c := range types { + cast = append(cast, pgx.Identifier{c.name}.Sanitize()+"::"+c.typeName) + } + key := pgx.Identifier{target.PKColumn()}.Sanitize() + return "SELECT pg_catalog.count(*)," + + " pg_catalog.md5(COALESCE(pg_catalog.string_agg(pg_catalog.md5(ROW(" + strings.Join(cast, ", ") + ")::text), '' ORDER BY " + key + "), ''))" + + " FROM " + pgx.Identifier{schema, table}.Sanitize() + + " WHERE " + key + " BETWEEN $1::bigint AND $2::bigint" +} + +// digest runs one side's statement for the closed range [lower, upper]. +func digest(ctx context.Context, tx pgx.Tx, sql string, lower, upper int64) (Digest, error) { + var d Digest + if err := tx.QueryRow(ctx, sql, lower, upper).Scan(&d.Rows, &d.Hash); err != nil { + return Digest{}, err + } + return d, nil +} diff --git a/pkg/checksum/digest_integration_test.go b/pkg/checksum/digest_integration_test.go new file mode 100644 index 0000000..7260dbc --- /dev/null +++ b/pkg/checksum/digest_integration_test.go @@ -0,0 +1,150 @@ +package checksum + +import ( + "strings" + "testing" + + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/internal/testutil" + "github.com/block/pg-sprite/pkg/copier" + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/preflight" +) + +// proofFixture is a throwaway schema on a superuser pool, enough to mint +// the copy-and-swap proof the verifier's constructor demands. The +// superuser is a SET-usable member of every role, so no provisioning is +// needed. +type proofFixture struct { + pool *pgxpool.Pool + schema string +} + +func newProofFixture(t *testing.T) proofFixture { + t.Helper() + pool, err := dbconn.NewPool(t.Context(), dbconn.Config{URL: testutil.StartPostgres(t)}) + require.NoError(t, err) + t.Cleanup(pool.Close) + return proofFixture{pool: pool, schema: testutil.NewSchema(t, pool)} +} + +// exec runs SQL with %s standing for the fixture schema. +func (f proofFixture) exec(t *testing.T, sql string) { + t.Helper() + _, err := f.pool.Exec(t.Context(), strings.ReplaceAll(sql, "%s", f.schema)) + require.NoError(t, err) +} + +// prove mints the copy-and-swap proof for table. +func (f proofFixture) prove(t *testing.T, table string) preflight.CopySwapTarget { + t.Helper() + role, err := preflight.CheckPrivileges(t.Context(), f.pool, f.schema, table, preflight.Requirement{Tier: preflight.TierCopyAndSwap}) + require.NoError(t, err) + target, err := preflight.CheckCopySwapShape(t.Context(), f.pool, f.schema, table, role) + require.NoError(t, err) + return target +} + +// fakeShadow is a Shadow minted by hand, so the constructor's proof checks +// can be exercised against every way a shadow can fail to describe the +// proven target. Passes over real tables use the shadow builder's proof. +type fakeShadow struct { + schema, source, shadow string + sourceOID, shadowOID uint32 + columns []string +} + +func (s fakeShadow) Schema() string { return s.schema } +func (s fakeShadow) SourceTable() string { return s.source } +func (s fakeShadow) ShadowTable() string { return s.shadow } +func (s fakeShadow) SourceOID() uint32 { return s.sourceOID } +func (s fakeShadow) ShadowOID() uint32 { return s.shadowOID } +func (s fakeShadow) CopyColumns() []string { return s.columns } + +// The digest statement is the D7 contract in one string: every copy column +// cast to the shadow's type inside a record, hashed per row, aggregated in +// key order, over a closed bigint-typed key range, with every function +// pg_catalog-qualified. Only the table differs between the two sides. +// Quoting goes through pgx.Identifier, so a column named like a keyword +// survives, while the type spelling is format_type's and is not quoted. +func TestDigestSQLIsFrozen(t *testing.T) { + f := newProofFixture(t) + f.exec(t, ` + CREATE TABLE %s.orders ( + id bigint PRIMARY KEY, + "select" text, + qty integer NOT NULL + )`) + target := f.prove(t, "orders") + types := []columnType{ + {name: "id", typeName: "bigint"}, + {name: "select", typeName: "character varying(20)"}, + {name: "qty", typeName: "numeric(10,2)"}, + } + want := `SELECT pg_catalog.count(*),` + + ` pg_catalog.md5(COALESCE(pg_catalog.string_agg(pg_catalog.md5(ROW("id"::bigint, "select"::character varying(20), "qty"::numeric(10,2))::text), '' ORDER BY "id"), ''))` + + ` FROM "` + f.schema + `"."_pgsprite_orders_new"` + + ` WHERE "id" BETWEEN $1::bigint AND $2::bigint` + assert.Equal(t, want, digestSQL(target, f.schema, "_pgsprite_orders_new", types)) + assert.Equal(t, strings.Replace(want, `"_pgsprite_orders_new"`, `"orders"`, 1), digestSQL(target, f.schema, "orders", types), + "the source side is the same statement over the source table") +} + +// Every way a shadow proof can fail to describe the proven target is +// refused before a connection is opened (ST-6); a good shadow still needs +// a lock session (LK-1). +func TestNewVerifierRefusesAShadowThatIsNotTheTargets(t *testing.T) { + f := newProofFixture(t) + f.exec(t, ` + CREATE TABLE %s.orders ( + id bigint PRIMARY KEY, + qty integer NOT NULL + )`) + target := f.prove(t, "orders") + good := fakeShadow{ + schema: f.schema, source: "orders", shadow: "_pgsprite_orders_new", + sourceOID: 1, shadowOID: 2, + columns: []string{"id", "qty"}, + } + cases := map[string]struct { + shadow copier.Shadow + detail string + }{ + "nil shadow": {nil, "verification requires a built shadow"}, + "zero shadow": {fakeShadow{}, "shadow proof is empty"}, + "other schema": {withSchema(good, "other"), "shadow is for other.orders, proof is for " + f.schema + ".orders"}, + "other source": {withSource(good, "invoices"), "shadow is for " + f.schema + ".invoices, proof is for " + f.schema + ".orders"}, + "shadow is the source by name": {withShadowTable(good, "orders"), "shadow of " + f.schema + ".orders is the source table itself"}, + "no source OID": {withOIDs(good, 0, 2), "shadow proof for " + f.schema + ".orders carries no relation OIDs"}, + "no shadow OID": {withOIDs(good, 1, 0), "shadow proof for " + f.schema + ".orders carries no relation OIDs"}, + "shadow is the source by relation": {withOIDs(good, 7, 7), "shadow proof for " + f.schema + ".orders names relation 7 as both source and shadow"}, + "no columns": {withColumns(good), "shadow copy columns for " + f.schema + ".orders do not include the primary key id"}, + "no primary key": {withColumns(good, "qty"), "shadow copy columns for " + f.schema + ".orders do not include the primary key id"}, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + _, err := NewVerifier(target, tc.shadow, nil, Options{}) + require.ErrorIs(t, err, ErrInvariantViolation) + assert.EqualError(t, err, "invariant violation (ST-6): "+tc.detail) + }) + } + + _, err := NewVerifier(target, good, nil, Options{}) + require.ErrorIs(t, err, ErrInvariantViolation, "a good shadow still needs a lock session") + assert.EqualError(t, err, "invariant violation (LK-1): verification requires a table lock session") +} + +func withSchema(s fakeShadow, schema string) fakeShadow { s.schema = schema; return s } +func withSource(s fakeShadow, source string) fakeShadow { s.source = source; return s } +func withShadowTable(s fakeShadow, shadow string) fakeShadow { + s.shadow = shadow + return s +} +func withOIDs(s fakeShadow, source, shadow uint32) fakeShadow { + s.sourceOID, s.shadowOID = source, shadow + return s +} +func withColumns(s fakeShadow, columns ...string) fakeShadow { s.columns = columns; return s } diff --git a/pkg/checksum/doc.go b/pkg/checksum/doc.go index a1d566b..f4ed965 100644 --- a/pkg/checksum/doc.go +++ b/pkg/checksum/doc.go @@ -1,2 +1,13 @@ -// Package checksum defines shadow-fidelity contracts enforcing CO-1, CO-2, and CO-3. +// Package checksum compares the shadow table with its source and defines +// the proofs the cutover will demand of that comparison (CO-1, CO-2, CO-3). +// A Verifier reads every key at or below the copier's landed watermark in +// chunks, digesting both tables inside one read-only REPEATABLE READ +// transaction per chunk so the two digests describe one snapshot, with +// every column cast to the type the shadow declares so a schema change that +// converts a column compares as the shadow holds it (D7). Each transaction +// runs under the table's lock session and refuses relations that are no +// longer the ones the proofs describe. A pass reports what differed; acting +// on a difference is the caller's divergence policy. VerifiedShadow and +// CleanWatermark are the proofs a clean pass will mint; their constructors +// are private to this package. package checksum diff --git a/pkg/checksum/guard.go b/pkg/checksum/guard.go new file mode 100644 index 0000000..67da12a --- /dev/null +++ b/pkg/checksum/guard.go @@ -0,0 +1,156 @@ +package checksum + +import ( + "context" + "errors" + "fmt" + "strconv" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/block/pg-sprite/pkg/dbconn" +) + +// sqlstateUndefinedTable is the SQLSTATE the server raises when a statement +// names a relation that does not exist. +const sqlstateUndefinedTable = "42P01" + +// begin opens the read-only transaction one chunk's two digests run in and +// guards it. The transaction is REPEATABLE READ so its two reads see one +// snapshot: the source and the shadow are compared as they stood at the +// same instant, and the snapshot is taken by the first query after both +// relations are locked, so no rename or drop can slip between the lock and +// the reads. +func (v *Verifier) begin(ctx context.Context, pool *pgxpool.Pool) (pgx.Tx, error) { + tx, err := pool.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly}) + if err != nil { + return nil, fmt.Errorf("begin verification of %s.%s: %w", v.target.Schema(), v.target.Table(), err) + } + if err := v.guard(ctx, tx); err != nil { + if rollbackErr := tx.Rollback(context.WithoutCancel(ctx)); rollbackErr != nil { + return nil, errors.Join(err, fmt.Errorf("roll back verification of %s.%s: %w", v.target.Schema(), v.target.Table(), rollbackErr)) + } + return nil, err + } + return tx, nil +} + +// guard prepares the transaction every read runs in: it bounds it, puts the +// catalog alone on its search_path, puts it under the source owner's role, +// locks both relations by name, and then confirms from this connection that +// the lock session's backend still holds the table and that the source and +// shadow are still the relations the proofs describe. Holding even ACCESS +// SHARE keeps a DROP or rename from completing until the transaction ends, +// so the identity check stays true for every statement after it. +func (v *Verifier) guard(ctx context.Context, tx pgx.Tx) error { + if err := setVerifySession(ctx, tx, v.target.OwnerRole(), v.opts); err != nil { + return err + } + if err := v.holdRelations(ctx, tx); err != nil { + return err + } + if err := v.confirmLock(ctx, tx); err != nil { + return err + } + return v.confirmRelations(ctx, tx) +} + +// setVerifySession bounds the transaction, restricts its search_path to the +// catalog so every unqualified operator, cast, and type name resolves there +// (CO-9), and puts it in the owner's shoes. SET LOCAL cannot take bind +// parameters; the timeouts are integer milliseconds. +func setVerifySession(ctx context.Context, tx pgx.Tx, owner string, opts Options) error { + // INV: LK-2 + budgets := "SET LOCAL lock_timeout = " + strconv.FormatInt(opts.LockTimeout.Milliseconds(), 10) + + "; SET LOCAL statement_timeout = " + strconv.FormatInt(opts.StatementTimeout.Milliseconds(), 10) + + "; " + dbconn.LocalSearchPath("pg_catalog") + if _, err := tx.Exec(ctx, budgets); err != nil { + return fmt.Errorf("set verification budgets: %w", err) + } + if _, err := tx.Exec(ctx, "SET LOCAL ROLE "+pgx.Identifier{owner}.Sanitize()); err != nil { + return fmt.Errorf("set owner role %s: %w", owner, err) + } + return nil +} + +// holdRelations takes the weakest table lock on the source and the shadow +// for the rest of the transaction. A name that no longer resolves is a +// relation replaced or gone since its proof was minted. +func (v *Verifier) holdRelations(ctx context.Context, tx pgx.Tx) error { + source := pgx.Identifier{v.shadow.Schema(), v.shadow.SourceTable()}.Sanitize() + shadow := pgx.Identifier{v.shadow.Schema(), v.shadow.ShadowTable()}.Sanitize() + _, err := tx.Exec(ctx, "LOCK TABLE "+source+", "+shadow+" IN ACCESS SHARE MODE") + if err == nil { + return nil + } + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && pgErr.Code == sqlstateUndefinedTable { + // INV: ST-6 + return fmt.Errorf("%w (ST-6): %s or its shadow %s no longer exists: %w", ErrInvariantViolation, source, shadow, err) + } + return fmt.Errorf("lock %s and %s for verification: %w", source, shadow, err) +} + +// confirmLock asks the server, on the connection about to read, whether the +// lock session's own backend holds the table. A server that answers "no one" +// or "another backend" has contradicted the session the verifier trusts, +// and that is the verifier's invariant violation; a lookup that got no +// answer is the connection's error, reported as such. +func (v *Verifier) confirmLock(ctx context.Context, conn dbconn.AdvisoryLockHolder) error { + err := v.lock.Confirm(ctx, conn) + if err == nil { + return nil + } + if lockDenied(err) { + // INV: LK-1 + return fmt.Errorf("%w (LK-1): %w", ErrInvariantViolation, err) + } + return fmt.Errorf("confirm the table lock from the verification transaction: %w", err) +} + +// lockDenied reports whether a Confirm error is the server's word that the +// session does not hold the table, in either of the two forms dbconn gives +// it. +func lockDenied(err error) bool { + if errors.Is(err, dbconn.ErrTableLockNotHeld) { + return true + } + var heldElsewhere *dbconn.TableLockHeldError + return errors.As(err, &heldElsewhere) +} + +// relationOIDsSQL resolves the source ($2) and shadow ($3) in schema $1 by +// explicit qualification; a missing relation scans as NULL. +const relationOIDsSQL = ` + SELECT + (SELECT c.oid FROM pg_catalog.pg_class c JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace WHERE n.nspname = $1 AND c.relname = $2), + (SELECT c.oid FROM pg_catalog.pg_class c JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace WHERE n.nspname = $1 AND c.relname = $3)` + +// confirmRelations refuses to read when the source or the shadow is no +// longer the relation its proof was minted for: a table dropped and +// recreated under the same name has a new OID, and a digest of it would +// verify a table nobody proved. +func (v *Verifier) confirmRelations(ctx context.Context, tx pgx.Tx) error { + var sourceOID, shadowOID *uint32 + err := tx.QueryRow(ctx, relationOIDsSQL, v.shadow.Schema(), v.shadow.SourceTable(), v.shadow.ShadowTable()).Scan(&sourceOID, &shadowOID) + if err != nil { + return fmt.Errorf("resolve %s.%s and its shadow %s: %w", v.shadow.Schema(), v.shadow.SourceTable(), v.shadow.ShadowTable(), err) + } + // INV: ST-6 + if err := confirmRelation("source", v.shadow.Schema(), v.shadow.SourceTable(), v.shadow.SourceOID(), sourceOID); err != nil { + return err + } + return confirmRelation("shadow", v.shadow.Schema(), v.shadow.ShadowTable(), v.shadow.ShadowOID(), shadowOID) +} + +func confirmRelation(role, schema, table string, proven uint32, found *uint32) error { + if found == nil { + return fmt.Errorf("%w (ST-6): %s %s.%s no longer exists", ErrInvariantViolation, role, schema, table) + } + if *found != proven { + return fmt.Errorf("%w (ST-6): %s %s.%s is relation %d, proof was minted for relation %d", ErrInvariantViolation, role, schema, table, *found, proven) + } + return nil +} diff --git a/pkg/checksum/report.go b/pkg/checksum/report.go new file mode 100644 index 0000000..8bb7b3a --- /dev/null +++ b/pkg/checksum/report.go @@ -0,0 +1,36 @@ +package checksum + +import "github.com/block/pg-sprite/pkg/copier" + +// Mismatch is one chunk whose source and shadow digests differ: the range, +// and what each side reported for it. It says that the two sides differed +// within one snapshot, not why; the divergence policy that acts on it is +// the caller's. +type Mismatch struct { + // Chunk is the closed key range the two digests cover. + Chunk copier.Chunk + // Source is the source table's digest of the chunk. + Source Digest + // Shadow is the shadow table's digest of the chunk. + Shadow Digest +} + +// Report is the outcome of one verification pass over the keys at or below +// a watermark: how much it read and every chunk that differed. A pass that +// finds no mismatch is clean, and only a clean pass can back a proof; the +// report itself proves nothing and is not one. +type Report struct { + // Through is the watermark the pass verified up to: every chunk at or + // below it was read. + Through copier.Watermark + // Chunks is the number of chunks the pass compared. + Chunks int + // Rows is the number of source rows the pass hashed. + Rows int64 + // Mismatches lists the chunks whose digests differed, in ascending key + // order. It is empty for a clean pass. + Mismatches []Mismatch +} + +// Clean reports whether every compared chunk matched. +func (r Report) Clean() bool { return len(r.Mismatches) == 0 } diff --git a/pkg/checksum/verifier.go b/pkg/checksum/verifier.go new file mode 100644 index 0000000..34df9bb --- /dev/null +++ b/pkg/checksum/verifier.go @@ -0,0 +1,285 @@ +package checksum + +import ( + "context" + "errors" + "fmt" + "slices" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/block/pg-sprite/pkg/copier" + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/preflight" + "github.com/block/pg-sprite/pkg/progress" +) + +var ( + // ErrInvariantViolation marks a fail-closed refusal; the wrapped message + // names the invariant. + ErrInvariantViolation = dbconn.ErrInvariantViolation + // ErrInvalidOptions reports verifier options that cannot bound a pass: a + // timeout below PostgreSQL's one-millisecond resolution, which the server + // would read as no timeout at all. + ErrInvalidOptions = errors.New("invalid verifier options") + // ErrNothingLanded reports a pass asked to verify up to the zero + // watermark, below which no key has been copied: there is nothing to + // compare, and a pass that compared nothing must not read as clean. + ErrNothingLanded = errors.New("nothing has landed to verify") +) + +// Options bounds a verification pass. Zero values take the defaults: the +// dbconn session timeouts for every chunk transaction, the copier's +// ChunkerOptions defaults, and the wall clock. +type Options struct { + // LockTimeout bounds every lock wait inside a chunk transaction. + LockTimeout time.Duration + // StatementTimeout bounds every statement inside a chunk transaction. + StatementTimeout time.Duration + // Chunker sizes the chunks a pass compares. They are cut independently + // of the chunks the copier copied; only the key range matters. + Chunker copier.ChunkerOptions + // 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 +} + +func (o Options) withDefaults() Options { + if o.LockTimeout == 0 { + o.LockTimeout = dbconn.DefaultLockTimeout + } + if o.StatementTimeout == 0 { + o.StatementTimeout = dbconn.DefaultStatementTimeout + } + if o.Clock == nil { + o.Clock = progress.WallClock{} + } + return o +} + +// validate runs after withDefaults, so every field is set. +func (o Options) validate() error { + for _, timeout := range []struct { + name string + value time.Duration + }{{"lock timeout", o.LockTimeout}, {"statement timeout", o.StatementTimeout}} { + if timeout.value < time.Millisecond { + // INV: LK-2 + return fmt.Errorf("%w: %s %s is below PostgreSQL's one-millisecond resolution; use zero for the default", ErrInvalidOptions, timeout.name, timeout.value) + } + } + return nil +} + +// Verifier compares a proven source table with its built shadow, chunk by +// chunk, and reports every chunk whose rows differ. Each chunk is read in +// 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. +type Verifier struct { + target preflight.CopySwapTarget + shadow copier.Shadow + lock *dbconn.TableLockSession + opts Options +} + +// NewVerifier prepares verification of target against shadow. It refuses a +// proof, shadow, or lock session that does not describe this table, and +// options that cannot bound a pass. The shadow is checked as the copier +// checks it: the verifier reads both relations the copier wrote between and +// must refuse everything the copier would have. +func NewVerifier(target preflight.CopySwapTarget, shadow copier.Shadow, lock *dbconn.TableLockSession, opts Options) (*Verifier, error) { + // INV: ST-6 + if target.Table() == "" { + return nil, fmt.Errorf("%w (ST-6): copy-and-swap target proof is empty", ErrInvariantViolation) + } + if err := checkShadow(target, shadow); err != nil { + return nil, err + } + if err := requireTableLock(lock, target); err != nil { + return nil, err + } + opts = opts.withDefaults() + if err := opts.validate(); err != nil { + return nil, err + } + // The chunker validates its own options; build one now so a pass never + // starts with options that cannot cut. + if _, err := copier.NewChunker(target, copier.Watermark{}, opts.Chunker); err != nil { + return nil, err + } + return &Verifier{target: target, shadow: shadow, lock: lock, opts: opts}, nil +} + +// checkShadow refuses a shadow proof that does not describe the proven +// target's shadow, with the copier's rules: the verifier compares the two +// relations the copier copied between, so a proof the copier would refuse +// describes nothing the verifier can vouch for. +func checkShadow(target preflight.CopySwapTarget, shadow copier.Shadow) error { + // INV: ST-6 + if shadow == nil { + return fmt.Errorf("%w (ST-6): verification requires a built shadow", ErrInvariantViolation) + } + if shadow.ShadowTable() == "" { + return fmt.Errorf("%w (ST-6): shadow proof is empty", ErrInvariantViolation) + } + if shadow.Schema() != target.Schema() || shadow.SourceTable() != target.Table() { + return fmt.Errorf("%w (ST-6): shadow is for %s.%s, proof is for %s.%s", ErrInvariantViolation, shadow.Schema(), shadow.SourceTable(), target.Schema(), target.Table()) + } + if shadow.ShadowTable() == shadow.SourceTable() { + return fmt.Errorf("%w (ST-6): shadow of %s.%s is the source table itself", ErrInvariantViolation, target.Schema(), target.Table()) + } + if shadow.SourceOID() == 0 || shadow.ShadowOID() == 0 { + return fmt.Errorf("%w (ST-6): shadow proof for %s.%s carries no relation OIDs", ErrInvariantViolation, target.Schema(), target.Table()) + } + if shadow.SourceOID() == shadow.ShadowOID() { + return fmt.Errorf("%w (ST-6): shadow proof for %s.%s names relation %d as both source and shadow", ErrInvariantViolation, target.Schema(), target.Table(), shadow.SourceOID()) + } + if !slices.Contains(shadow.CopyColumns(), target.PKColumn()) { + return fmt.Errorf("%w (ST-6): shadow copy columns for %s.%s do not include the primary key %s", ErrInvariantViolation, target.Schema(), target.Table(), target.PKColumn()) + } + return nil +} + +// requireTableLock refuses to verify without the per-table lock that keeps +// a second engine instance off the same table: the session must exist, +// carry a populated proof for the proven table, and not have reported loss. +func requireTableLock(lock *dbconn.TableLockSession, target preflight.CopySwapTarget) error { + // INV: LK-1 + if lock == nil { + return fmt.Errorf("%w (LK-1): verification requires a table lock session", ErrInvariantViolation) + } + held := lock.Lock() + if held.Table() == "" { + return fmt.Errorf("%w (LK-1): table lock proof is empty", ErrInvariantViolation) + } + if held.Schema() != target.Schema() || held.Table() != target.Table() { + return fmt.Errorf("%w (LK-1): table lock is for %s.%s, proof is for %s.%s", ErrInvariantViolation, held.Schema(), held.Table(), target.Schema(), target.Table()) + } + if err := lock.Err(); err != nil { + return fmt.Errorf("%w (LK-1): table lock was lost before verification: %w", ErrInvariantViolation, err) + } + return nil +} + +// Verify runs one pass over every key at or below through — the copier's +// landed watermark, below which every chunk has committed — and reports +// what it found. Keys above the watermark are not compared: they are in +// chunks the copier has not finished, and a difference there is expected, +// not a finding. The pass reads the shadow's column types once, then cuts +// its own chunks from the live source and digests each in its own guarded +// transaction. It stops at the first error; a report with mismatches is not +// 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. +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()) + } + ctx, unbind := v.lock.Bind(ctx) + defer unbind() + report, err := v.pass(ctx, pool, through) + if lost := v.lock.Err(); lost != nil { + // INV: LK-1 + return Report{}, fmt.Errorf("%w (LK-1): table lock was lost during verification: %w", ErrInvariantViolation, lost) + } + return report, err +} + +func (v *Verifier) pass(ctx context.Context, pool *pgxpool.Pool, through copier.Watermark) (Report, error) { + sourceSQL, shadowSQL, err := v.digestStatements(ctx, pool) + if err != nil { + return Report{}, err + } + chunker, err := copier.NewChunker(v.target, copier.Watermark{}, v.opts.Chunker) + if err != nil { + return Report{}, err + } + report := Report{Through: through} + for { + chunk, ok, err := chunker.Next(ctx, pool) + if err != nil { + return Report{}, err + } + if !ok { + return report, nil + } + lower, upper := chunk.Lower(), min(chunk.Upper(), through.Value()) + started := v.opts.Clock.Now() + source, shadow, err := v.digestChunk(ctx, pool, sourceSQL, shadowSQL, lower, upper) + if err != nil { + return Report{}, err + } + report.Chunks++ + report.Rows += source.Rows + if source != shadow { + compared, err := copier.NewChunk(lower, upper) + if err != nil { + return Report{}, fmt.Errorf("%w (CO-1): verified chunk: %w", ErrInvariantViolation, err) + } + report.Mismatches = append(report.Mismatches, Mismatch{Chunk: compared, Source: source, Shadow: shadow}) + } + if upper == through.Value() { + return report, nil + } + if err := chunker.Feedback(chunk, v.opts.Clock.Now().Sub(started)); err != nil { + return Report{}, err + } + } +} + +// digestStatements freezes the two digest statements for this pass from the +// shadow's column types, read in a guarded transaction so the shadow whose +// types they name is the shadow the proof describes. +func (v *Verifier) digestStatements(ctx context.Context, pool *pgxpool.Pool) (sourceSQL, shadowSQL string, err error) { + tx, err := v.begin(ctx, pool) + if err != nil { + return "", "", err + } + defer func() { + // Redundant safety closer: after a successful Commit this returns + // the guaranteed ErrTxClosed; on a failure path the server aborts + // the transaction with its session either way. + _ = tx.Rollback(context.WithoutCancel(ctx)) + }() + types, err := shadowColumnTypes(ctx, tx, v.shadow) + if err != nil { + return "", "", err + } + if err := tx.Commit(ctx); err != nil { + return "", "", fmt.Errorf("commit column-type read of %s.%s: %w", v.shadow.Schema(), v.shadow.ShadowTable(), err) + } + return digestSQL(v.target, v.shadow.Schema(), v.shadow.SourceTable(), types), + digestSQL(v.target, v.shadow.Schema(), v.shadow.ShadowTable(), types), + nil +} + +// digestChunk reads both sides of the closed range [lower, upper] in one +// guarded transaction, so the two digests describe one snapshot. +func (v *Verifier) digestChunk(ctx context.Context, pool *pgxpool.Pool, sourceSQL, shadowSQL string, lower, upper int64) (source, shadow Digest, err error) { + tx, err := v.begin(ctx, pool) + if err != nil { + return Digest{}, Digest{}, err + } + defer func() { + // Redundant safety closer: after a successful Commit this returns + // the guaranteed ErrTxClosed; on a failure path the server aborts + // the transaction with its session either way. + _ = tx.Rollback(context.WithoutCancel(ctx)) + }() + source, err = digest(ctx, tx, sourceSQL, lower, upper) + if err != nil { + return Digest{}, Digest{}, fmt.Errorf("digest chunk [%d, %d] of %s.%s: %w", lower, upper, v.shadow.Schema(), v.shadow.SourceTable(), err) + } + shadow, err = digest(ctx, tx, shadowSQL, lower, upper) + if err != nil { + return Digest{}, Digest{}, fmt.Errorf("digest chunk [%d, %d] of shadow %s.%s: %w", lower, upper, v.shadow.Schema(), v.shadow.ShadowTable(), err) + } + 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) + } + return source, shadow, nil +} diff --git a/pkg/checksum/verifier_integration_test.go b/pkg/checksum/verifier_integration_test.go new file mode 100644 index 0000000..e6cd3bd --- /dev/null +++ b/pkg/checksum/verifier_integration_test.go @@ -0,0 +1,451 @@ +package checksum_test + +import ( + "context" + "fmt" + "math" + "strings" + "sync" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/internal/testutil" + "github.com/block/pg-sprite/pkg/checksum" + "github.com/block/pg-sprite/pkg/copier" + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/preflight" + "github.com/block/pg-sprite/pkg/schemachange" + "github.com/block/pg-sprite/pkg/statement" +) + +// verifierFixture is a throwaway schema on a superuser pool with the real +// shadow builder and copier in front of the verifier, so every pass below +// compares a shadow the builder proved and the copier filled. The superuser +// is a SET-usable member of every role, so the copy-and-swap proof is minted +// without provisioning. +type verifierFixture struct { + cfg dbconn.Config + pool *pgxpool.Pool + schema string +} + +func newVerifierFixture(t *testing.T) verifierFixture { + t.Helper() + cfg := dbconn.Config{URL: testutil.StartPostgres(t)} + pool, err := dbconn.NewPool(t.Context(), cfg) + require.NoError(t, err) + t.Cleanup(pool.Close) + return verifierFixture{cfg: cfg, pool: pool, schema: testutil.NewSchema(t, pool)} +} + +// exec runs SQL with %s standing for the fixture schema. +func (f verifierFixture) exec(t *testing.T, sql string) { + t.Helper() + _, err := f.pool.Exec(t.Context(), strings.ReplaceAll(sql, "%s", f.schema)) + require.NoError(t, err) +} + +// rows is the size of every orders table below: keys 1..rows with qty = key. +const rows = 2500 + +// chunkRows sizes every pass below to fixed thousand-row chunks, so with +// rows keys the chunks are [MinInt64, 1000], [1001, 2000], and +// [2001, MaxInt64] whatever the chunk timing feedback says. +var chunkRows = copier.ChunkerOptions{InitialRows: 1000, MinRows: 1000, MaxRows: 1000} + +// createOrders creates the orders table every pass below compares, holding +// keys 1..rows with qty = key and a note naming the order. +func (f verifierFixture) createOrders(t *testing.T) { + t.Helper() + f.exec(t, ` + CREATE TABLE %s.orders ( + id bigint PRIMARY KEY, + qty integer NOT NULL, + note text + )`) + f.exec(t, fmt.Sprintf(` + INSERT INTO %%s.orders (id, qty, note) + SELECT n, n, 'order ' || n FROM generate_series(1, %d) AS n`, rows)) +} + +// prove mints the copy-and-swap proof for table. +func (f verifierFixture) prove(t *testing.T, table string) preflight.CopySwapTarget { + t.Helper() + role, err := preflight.CheckPrivileges(t.Context(), f.pool, f.schema, table, preflight.Requirement{Tier: preflight.TierCopyAndSwap}) + require.NoError(t, err) + target, err := preflight.CheckCopySwapShape(t.Context(), f.pool, f.schema, table, role) + require.NoError(t, err) + return target +} + +// lock acquires the per-table lock the verifier requires and releases it +// when the test ends. +func (f verifierFixture) lock(t *testing.T, table string, options ...dbconn.TableLockOption) *dbconn.TableLockSession { + t.Helper() + lock, err := dbconn.AcquireTableLock(t.Context(), f.cfg, f.schema, table, options...) + require.NoError(t, err) + t.Cleanup(func() { + // A test that deliberately loses the lock has already seen Release's + // invariant error through Err; a clean test releases cleanly. + if lock.Err() == nil { + assert.NoError(t, lock.Release(context.WithoutCancel(t.Context()))) + } + }) + return lock +} + +// build runs the shadow builder for target under lock with the given +// change, where %s stands for the schema-qualified table. +func (f verifierFixture) build(t *testing.T, lock *dbconn.TableLockSession, target preflight.CopySwapTarget, alterSQL string) schemachange.BuiltShadow { + t.Helper() + alter, err := statement.ParseOne(fmt.Sprintf(alterSQL, pgx.Identifier{f.schema, target.Table()}.Sanitize())) + require.NoError(t, err) + shadow, err := schemachange.BuildShadow(t.Context(), f.pool, lock, target, alter, schemachange.Options{}) + require.NoError(t, err) + return shadow +} + +// copy fills the shadow from the source, so the pass compares a shadow the +// copier landed. +func (f verifierFixture) copy(t *testing.T, target preflight.CopySwapTarget, shadow schemachange.BuiltShadow, lock *dbconn.TableLockSession) { + t.Helper() + c, err := copier.NewCopier(target, shadow, lock, copier.Watermark{}, copier.Options{Workers: 2, Chunker: chunkRows}) + require.NoError(t, err) + require.NoError(t, c.Run(t.Context(), f.pool)) + require.Equal(t, copier.NewWatermark(math.MaxInt64), c.Position().Watermark) +} + +// prepare creates, proves, locks, builds the shadow of, and copies an orders +// table whose schema change drops the note column, so the copy columns are +// narrower than the source. +func (f verifierFixture) prepare(t *testing.T) (preflight.CopySwapTarget, *dbconn.TableLockSession, schemachange.BuiltShadow) { + t.Helper() + f.createOrders(t) + target := f.prove(t, "orders") + lock := f.lock(t, "orders") + shadow := f.build(t, lock, target, `ALTER TABLE %s DROP COLUMN note`) + f.copy(t, target, shadow, lock) + return target, lock, shadow +} + +// shadowName is the schema-qualified, quoted shadow table. +func (f verifierFixture) shadowName(shadow schemachange.BuiltShadow) string { + return pgx.Identifier{f.schema, shadow.ShadowTable()}.Sanitize() +} + +// verify runs one pass over the whole key space with fixed chunks. +func (f verifierFixture) verify(t *testing.T, pool *pgxpool.Pool, target preflight.CopySwapTarget, shadow schemachange.BuiltShadow, lock *dbconn.TableLockSession, through copier.Watermark) (checksum.Report, error) { + t.Helper() + v, err := checksum.NewVerifier(target, shadow, lock, checksum.Options{Chunker: chunkRows}) + require.NoError(t, err) + return v.Verify(t.Context(), pool, through) +} + +func chunk(t *testing.T, lower, upper int64) copier.Chunk { + t.Helper() + c, err := copier.NewChunk(lower, upper) + require.NoError(t, err) + return c +} + +// A pass over a shadow the copier filled completely reports every chunk +// clean, with the row count the source holds, and its report carries the +// watermark it verified through (CO-1). A pass asked to verify up to the +// zero watermark has nothing to compare and refuses rather than reporting +// an empty pass as clean. +func TestVerifierReportsACompleteCopyClean(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + + report, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + assert.True(t, report.Clean()) + assert.Empty(t, report.Mismatches) + assert.Equal(t, copier.NewWatermark(math.MaxInt64), report.Through) + assert.Equal(t, 3, report.Chunks, "thousand-row chunks over 2500 keys") + assert.Equal(t, int64(rows), report.Rows) + + _, err = f.verify(t, f.pool, target, shadow, lock, copier.Watermark{}) + assert.ErrorIs(t, err, checksum.ErrNothingLanded) +} + +// One shadow row whose value differs from the source is reported as +// exactly the chunk that holds its key, with equal row counts and +// differing hashes on the two sides; the chunks around it stay clean. +func TestVerifierLocatesAChangedShadowRow(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1500") + + report, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + assert.False(t, report.Clean()) + assert.Equal(t, 3, report.Chunks) + assert.Equal(t, int64(rows), report.Rows) + require.Len(t, report.Mismatches, 1) + found := report.Mismatches[0] + assert.Equal(t, chunk(t, 1001, 2000), found.Chunk) + assert.Equal(t, int64(1000), found.Source.Rows) + assert.Equal(t, int64(1000), found.Shadow.Rows) + assert.NotEqual(t, found.Source.Hash, found.Shadow.Hash) +} + +// A row missing from the shadow and a row the shadow holds that the source +// does not are each reported in their own chunk, in key order, and the row +// counts on the two sides say which way each chunk differs. +func TestVerifierCountsMissingAndExtraShadowRows(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "DELETE FROM "+f.shadowName(shadow)+" WHERE id = 5") + f.exec(t, "INSERT INTO "+f.shadowName(shadow)+" (id, qty) VALUES (2600, 2600)") + + report, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + require.Len(t, report.Mismatches, 2) + missing, extra := report.Mismatches[0], report.Mismatches[1] + assert.Equal(t, chunk(t, math.MinInt64, 1000), missing.Chunk) + assert.Equal(t, int64(1000), missing.Source.Rows) + assert.Equal(t, int64(999), missing.Shadow.Rows) + assert.Equal(t, chunk(t, 2001, math.MaxInt64), extra.Chunk) + assert.Equal(t, int64(500), extra.Source.Rows) + assert.Equal(t, int64(501), extra.Shadow.Rows) + assert.Equal(t, int64(rows), report.Rows, "the report counts source rows") +} + +// A schema change that converts a column's type compares clean: both sides +// hash the value as the shadow's type renders it (D7). The source holds +// integer 5 and the shadow numeric 5.00; without the cast the source would +// hash "5" against the shadow's "5.00" and every chunk would differ. +func TestVerifierComparesThroughTheShadowsTypes(t *testing.T) { + f := newVerifierFixture(t) + f.createOrders(t) + target := f.prove(t, "orders") + lock := f.lock(t, "orders") + shadow := f.build(t, lock, target, `ALTER TABLE %s ALTER COLUMN qty TYPE numeric(10,2)`) + f.copy(t, target, shadow, lock) + + report, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + assert.True(t, report.Clean(), "mismatches: %+v", report.Mismatches) + assert.Equal(t, int64(rows), report.Rows) +} + +// A pass compares only the keys at or below the watermark it is given: the +// chunk straddling the watermark is clamped to it, a difference below is +// found, and a difference above — in a chunk the copier has not finished — +// is not a finding (CO-1). +func TestVerifierStopsAtTheWatermark(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id IN (1400, 1600)") + + report, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(1500)) + require.NoError(t, err) + assert.Equal(t, copier.NewWatermark(1500), report.Through) + assert.Equal(t, 2, report.Chunks) + assert.Equal(t, int64(1500), report.Rows) + require.Len(t, report.Mismatches, 1) + assert.Equal(t, chunk(t, 1001, 1500), report.Mismatches[0].Chunk, "the straddling chunk is clamped to the watermark") + assert.Equal(t, int64(500), report.Mismatches[0].Source.Rows) +} + +// The pass resolves every catalog object it names through pg_catalog, so a +// session whose search_path puts a schema of impostors first (CO-9) neither +// misreads the shadow's types nor hashes with someone else's md5: an +// impostor format_type that would make every numeric compare as text +// produces no false mismatch, and an impostor md5 that answers the same for +// every row hides no real one. +func TestVerifierIgnoresTheSessionSearchPath(t *testing.T) { + f := newVerifierFixture(t) + f.createOrders(t) + target := f.prove(t, "orders") + lock := f.lock(t, "orders") + shadow := f.build(t, lock, target, `ALTER TABLE %s ALTER COLUMN qty TYPE numeric(10,2)`) + f.copy(t, target, shadow, lock) + f.exec(t, `CREATE FUNCTION %s.md5(text) RETURNS text LANGUAGE sql IMMUTABLE AS 'SELECT ''impostor''::text'`) + f.exec(t, `CREATE FUNCTION %s.format_type(oid, integer) RETURNS text LANGUAGE sql STABLE AS 'SELECT ''text''::text'`) + shadowing := testutil.NewCatalogShadowingPool(t, f.cfg.URL, f.schema) + + report, err := f.verify(t, shadowing, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + assert.True(t, report.Clean(), "mismatches: %+v", report.Mismatches) + + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 42") + report, err = f.verify(t, shadowing, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + require.Len(t, report.Mismatches, 1) + assert.Equal(t, chunk(t, math.MinInt64, 1000), report.Mismatches[0].Chunk) +} + +// A pass refuses to read when the source, or the shadow, has been replaced +// by another relation of the same name since the proofs were minted, or is +// gone (ST-6): a same-shaped impostor is never vouched for. +func TestVerifierRefusesReplacedRelations(t *testing.T) { + t.Run("source replaced", func(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "DROP TABLE %s.orders") + f.exec(t, "CREATE TABLE %s.orders (id bigint PRIMARY KEY, qty integer NOT NULL, note text)") + + _, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.ErrorIs(t, err, checksum.ErrInvariantViolation) + assert.Contains(t, err.Error(), "(ST-6): source") + }) + t.Run("shadow replaced", func(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "DROP TABLE "+f.shadowName(shadow)) + f.exec(t, "CREATE TABLE "+f.shadowName(shadow)+" (id bigint PRIMARY KEY, qty integer NOT NULL)") + + _, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.ErrorIs(t, err, checksum.ErrInvariantViolation) + assert.Contains(t, err.Error(), "(ST-6): shadow") + }) + t.Run("shadow dropped", func(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "DROP TABLE "+f.shadowName(shadow)) + + _, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.ErrorIs(t, err, checksum.ErrInvariantViolation) + var pgErr *pgconn.PgError + require.ErrorAs(t, err, &pgErr, "the server's refusal to lock a missing relation is the cause") + assert.Equal(t, "42P01", pgErr.Code, "undefined_table") + assert.Contains(t, err.Error(), "(ST-6)") + }) +} + +// The verifier trusts the lock session only as far as the server confirms +// it from the reading connection (LK-1): a session whose backend is gone +// compares nothing. A session for a different table, or one that already +// reported loss, is refused before any connection is opened. +func TestVerifierRefusesAnUnconfirmedLock(t *testing.T) { + f := newVerifierFixture(t) + f.createOrders(t) + target := f.prove(t, "orders") + buildLock := f.lock(t, "orders") + shadow := f.build(t, buildLock, target, `ALTER TABLE %s DROP COLUMN note`) + f.copy(t, target, shadow, buildLock) + require.NoError(t, buildLock.Release(t.Context())) + + f.exec(t, `CREATE TABLE %s.other (id bigint PRIMARY KEY)`) + _, err := checksum.NewVerifier(target, shadow, f.lock(t, "other"), checksum.Options{}) + require.ErrorIs(t, err, checksum.ErrInvariantViolation, "a lock for another table") + assert.Contains(t, err.Error(), "(LK-1)") + + gone := f.goneLock(t, "orders") + _, err = f.verify(t, f.pool, target, shadow, gone, copier.NewWatermark(math.MaxInt64)) + require.ErrorIs(t, err, checksum.ErrInvariantViolation, "a gone lock session") + assert.ErrorIs(t, err, dbconn.ErrTableLockNotHeld) + + lost := f.lostLock(t, "orders") + _, err = checksum.NewVerifier(target, shadow, lost, checksum.Options{}) + require.ErrorIs(t, err, checksum.ErrInvariantViolation, "a session that already reported loss") + assert.ErrorIs(t, err, lost.Err(), "the refusal names what the session saw") + assert.Contains(t, err.Error(), "(LK-1)") +} + +// Losing the table lock mid-pass cancels the read in flight and Verify +// reports the loss as an invariant violation naming what the session saw +// (LK-1), not the cancelled statement. The loss is injected from the clock +// reading the verifier takes just before its first chunk's transaction. +func TestVerifierAbortsWhenTheLockIsLostMidPass(t *testing.T) { + f := newVerifierFixture(t) + f.createOrders(t) + target := f.prove(t, "orders") + buildLock := f.lock(t, "orders") + shadow := f.build(t, buildLock, target, `ALTER TABLE %s DROP COLUMN note`) + f.copy(t, target, shadow, buildLock) + require.NoError(t, buildLock.Release(t.Context())) + lock := f.lock(t, "orders", dbconn.WithTableLockKeepalive(100*time.Millisecond)) + clock := &hookedClock{hook: func() { + f.terminateBackend(t, lock.BackendPID()) + const lockLossDeadline = 15 * time.Second + select { + case <-lock.Done(): + case <-time.After(lockLossDeadline): + t.Errorf("lock session did not report loss within %s", lockLossDeadline) + } + }} + + v, err := checksum.NewVerifier(target, shadow, lock, checksum.Options{Chunker: chunkRows, Clock: clock}) + require.NoError(t, err) + _, err = v.Verify(t.Context(), f.pool, copier.NewWatermark(math.MaxInt64)) + require.ErrorIs(t, err, checksum.ErrInvariantViolation) + assert.ErrorIs(t, err, lock.Err(), "the loss the session reported is the cause") + assert.Contains(t, err.Error(), "(LK-1)") +} + +// goneLock acquires the table lock and then terminates the session's +// backend, so the server no longer grants the lock while the session still +// believes it holds it: the default keepalive is long enough that only an +// in-transaction confirmation can catch the loss. +func (f verifierFixture) goneLock(t *testing.T, table string) *dbconn.TableLockSession { + t.Helper() + lock, err := dbconn.AcquireTableLock(t.Context(), f.cfg, f.schema, table) + require.NoError(t, err) + t.Cleanup(func() { + const lockLossDeadline = 30 * time.Second + select { + case <-lock.Done(): + case <-time.After(lockLossDeadline): + t.Errorf("lock session did not report loss within %s", lockLossDeadline) + } + assert.Error(t, lock.Release(context.WithoutCancel(t.Context()))) + }) + f.terminateBackend(t, lock.BackendPID()) + require.NoError(t, lock.Err(), "the keepalive has not yet noticed the loss; the in-transaction check must") + return lock +} + +// lostLock acquires the table lock with a short keepalive, terminates the +// session's backend, and waits until the session itself has reported the +// loss, so Err is set before the verifier ever sees the session. +func (f verifierFixture) lostLock(t *testing.T, table string) *dbconn.TableLockSession { + t.Helper() + lock, err := dbconn.AcquireTableLock(t.Context(), f.cfg, f.schema, table, dbconn.WithTableLockKeepalive(200*time.Millisecond)) + require.NoError(t, err) + t.Cleanup(func() { assert.Error(t, lock.Release(context.WithoutCancel(t.Context()))) }) + f.terminateBackend(t, lock.BackendPID()) + const lockLossDeadline = 30 * time.Second + select { + case <-lock.Done(): + case <-time.After(lockLossDeadline): + t.Fatalf("lock session did not report loss within %s", lockLossDeadline) + } + require.Error(t, lock.Err()) + return lock +} + +func (f verifierFixture) terminateBackend(t *testing.T, pid uint32) { + t.Helper() + var terminated bool + require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT pg_terminate_backend($1)`, pid).Scan(&terminated)) + require.True(t, terminated) + const backendExitDeadline = 10 * time.Second + require.Eventually(t, func() bool { + var alive bool + require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT EXISTS (SELECT 1 FROM pg_stat_activity WHERE pid = $1)`, pid).Scan(&alive)) + return !alive + }, backendExitDeadline, 50*time.Millisecond, "terminated backend should leave pg_stat_activity") +} + +// hookedClock is a Clock whose first reading runs a hook. The verifier +// reads the clock once before each chunk's transaction begins, so the hook +// lands after the column-type read has committed and before the first +// chunk's digest transaction opens. +type hookedClock struct { + once sync.Once + hook func() +} + +func (c *hookedClock) Now() time.Time { + c.once.Do(c.hook) + return time.Now() +} diff --git a/pkg/checksum/verifier_test.go b/pkg/checksum/verifier_test.go new file mode 100644 index 0000000..1ce251a --- /dev/null +++ b/pkg/checksum/verifier_test.go @@ -0,0 +1,59 @@ +package checksum + +import ( + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/copier" + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/preflight" + "github.com/block/pg-sprite/pkg/progress" +) + +func TestNewVerifierRejectsEmptyProof(t *testing.T) { + _, err := NewVerifier(preflight.CopySwapTarget{}, nil, nil, Options{}) + require.ErrorIs(t, err, ErrInvariantViolation) + assert.EqualError(t, err, "invariant violation (ST-6): copy-and-swap target proof is empty") +} + +func TestVerifierOptionsDefaults(t *testing.T) { + opts := Options{}.withDefaults() + assert.Equal(t, dbconn.DefaultLockTimeout, opts.LockTimeout) + assert.Equal(t, dbconn.DefaultStatementTimeout, opts.StatementTimeout) + assert.Equal(t, progress.WallClock{}, opts.Clock) + require.NoError(t, opts.validate()) + + given := Options{LockTimeout: time.Second, StatementTimeout: time.Minute}.withDefaults() + assert.Equal(t, time.Second, given.LockTimeout) + assert.Equal(t, time.Minute, given.StatementTimeout) +} + +// A timeout the server would read as disabled is refused rather than +// defaulted (LK-2). +func TestVerifierOptionsRefuseUnboundedValues(t *testing.T) { + cases := map[string]Options{ + "lock timeout below a millisecond": {LockTimeout: 500 * time.Microsecond}, + "statement timeout below a millisecond": {StatementTimeout: time.Nanosecond}, + } + for name, opts := range cases { + t.Run(name, func(t *testing.T) { + assert.ErrorIs(t, opts.withDefaults().validate(), ErrInvalidOptions) + }) + } +} + +// A report is clean only when no chunk differed; how much it read does not +// enter into it. +func TestReportClean(t *testing.T) { + chunk, err := copier.NewChunk(1, 10) + require.NoError(t, err) + assert.True(t, Report{Through: copier.NewWatermark(10), Chunks: 1, Rows: 10}.Clean()) + assert.False(t, Report{Through: copier.NewWatermark(10), Chunks: 1, Rows: 10, Mismatches: []Mismatch{{ + Chunk: chunk, + Source: Digest{Rows: 10, Hash: "a"}, + Shadow: Digest{Rows: 10, Hash: "b"}, + }}}.Clean()) +} From da5303cec3a2b109695bb0a4fdd6724ea9027465 Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Thu, 1 Oct 2026 11:33:34 +1000 Subject: [PATCH 2/3] checksum: pin extra_float_digits and test every guard property The digest hashes each row's text rendering, and for float4, float8, and the geometric types that rendering follows extra_float_digits: at zero or below the server rounds to fifteen significant digits, so a database or role configured there rendered two floats that differ only in their last digits the same and the pass compared them clean. The guarded session now pins the setting to its maximum alongside the timeouts and search_path. Tests now hold every guard property a mutation could drop: a float difference under a database at extra_float_digits = 0; a decoy bigint <= operator for the local search_path, which is the one layer the qualified functions cannot cover; one snapshot for both digests, with a shadow write injected from a query hook between them; the owner-role read, against a throwaway owner whose SELECT is revoked; and the refusal of a shadow that has lost a copy column. The missing/extra-row test drops two rows so the source and shadow counts differ and the report's row count is pinned to the source. --- docs/copy-and-swap-design.md | 2 +- docs/invariants.md | 3 +- pkg/checksum/guard.go | 18 ++- pkg/checksum/verifier_integration_test.go | 172 +++++++++++++++++++++- 4 files changed, 181 insertions(+), 14 deletions(-) diff --git a/docs/copy-and-swap-design.md b/docs/copy-and-swap-design.md index 96d5625..8f3155b 100644 --- a/docs/copy-and-swap-design.md +++ b/docs/copy-and-swap-design.md @@ -379,7 +379,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, dependent-object, decoding, and headroom checks above); owns Tier-3 refusals. | ST-6, RF-1..RF-3 | | `pkg/copier` | Produces `Chunk` and `Watermark`; `Chunker` (built only from a `CopySwapTarget`) cuts consecutive chunks that tile the whole int64 key space — first open below, last open above — so every key a row can carry belongs to exactly one chunk and a watermark at the largest value means the copy is complete. `Copier` (built from a `CopySwapTarget`, a `Shadow` — the shape `schemachange.BuiltShadow` satisfies — and the table's `TableLockSession`) copies chunks with several workers, each in its own bounded transaction under the owner's role that confirms the lock, takes `ACCESS SHARE` on both relations, and confirms both relation OIDs before one frozen never-overwriting insert; a resumed copy first deletes every shadow row above the watermark in batches of the same guarded shape — the whole shadow after a zero watermark — so the resumed cut frontier and the shadow agree on what is uncut, and its first batch takes `SHARE MODE` on the shadow so a chunk transaction of the earlier run still committing into it ends before the clear reads; the proof check refuses a shadow that is the source by name or by OID, since the clear deletes from one and the copy reads the other; the copy statement carries no conversion expression, so a column whose type differs between the two tables takes the server's assignment cast, and an `ALTER COLUMN TYPE … USING` change needs a conversion-aware copy before the planner's copy-and-swap route is wired to the copier; `Position` snapshots the cut frontier and landed watermark as two `Watermark` values plus the in-flight chunks, and `Position.Classify` is the applier's uncut / in-flight / landed rule. While `Run` is in progress the `Copier` is the `progress.Tracker`'s `WorkSource` (`Options.Tracker`): `rows_copied` is the rows this run's committed chunks inserted, `rows_total` the source's catalog row count read once at the start of the run, and `bytes_copied` / `bytes_total` the shadow's and the source's `pg_table_size` measured at each poll — every counter measured, none projected; on a resumed run the row ratio is that run's share, not the copy's completion. Each measurement runs in a read-only transaction of its own under the chunk budgets, with every catalog name `pg_catalog`-qualified, so a poll on a caller-built pool is still bounded and still reads the real catalog; the copier stops being the source before `Run` returns. The copier's refusals are fail-closed `ErrInvariantViolation` values tagged with their invariant; they take a typed cause in [refusal-classes.md](refusal-classes.md#shadow-operation-refusals-keyed-on-refusalcause) when the orchestrator that runs the copy is wired, so the cause lands with its first importer. | CO-4, LK-1, LK-3 | -| `pkg/checksum` | `Verifier` (built from a `CopySwapTarget`, a `copier.Shadow`, and the table's `TableLockSession`) compares the source with its shadow up to the copier's landed watermark and returns a `Report` of the chunks that differ. It cuts its own chunks with a `copier.Chunker` and digests each in one read-only `REPEATABLE READ` transaction, so a chunk's two digests describe one snapshot and no snapshot outlives one chunk; each transaction runs the copier's guard — owner role, catalog-only `search_path`, `ACCESS SHARE` on both relations taken before the snapshot, lock confirmation, relation-OID check. Both sides run the identical frozen statement — row count plus `md5` of the key-ordered concatenation of each row's `md5(ROW(col::shadow_type, …)::text)` over `pk BETWEEN $1 AND $2` — so only the data can differ, and the cast on every column is D7: a converted column hashes as the value the shadow holds. Only the copy columns are compared; generated columns present on both sides are not yet hashed. A `Report` proves nothing — the divergence policy, repair, and the constructors of `VerifiedShadow` and `CleanWatermark` (private to this package) land on top of it. | CO-1, CO-2, CO-3, CO-9, LK-1 | +| `pkg/checksum` | `Verifier` (built from a `CopySwapTarget`, a `copier.Shadow`, and the table's `TableLockSession`) compares the source with its shadow up to the copier's landed watermark and returns a `Report` of the chunks that differ. It cuts its own chunks with a `copier.Chunker` and digests each in one read-only `REPEATABLE READ` transaction, so a chunk's two digests describe one snapshot and no snapshot outlives one chunk; each transaction runs the copier's guard — owner role, catalog-only `search_path`, `ACCESS SHARE` on both relations taken before the snapshot, lock confirmation, relation-OID check — and pins `extra_float_digits` to its maximum, since the digest hashes each row's text rendering and a database or role configured at zero or below would render two floats that differ only in their last digits the same. Both sides run the identical frozen statement — row count plus `md5` of the key-ordered concatenation of each row's `md5(ROW(col::shadow_type, …)::text)` over `pk BETWEEN $1 AND $2` — so only the data can differ, and the cast on every column is D7: a converted column hashes as the value the shadow holds. Only the copy columns are compared; generated columns present on both sides are not yet hashed. A `Report` proves nothing — the divergence policy, repair, and the constructors of `VerifiedShadow` and `CleanWatermark` (private to this package) land on top of it. | CO-1, CO-2, CO-3, CO-9, LK-1 | | `pkg/decode` | Produces `ChangeEvent`, including per-column presence and `OldKey` for an UPDATE that moved the primary key. | ST-3, ST-4, CO-4, CO-8 | | `pkg/applier` | Applies presence-aware events from the per-key buffer. | CO-4, CO-5, CO-6, CO-8, LK-3 | | `pkg/checkpoint` | Produces `Checkpoint`. | ST-1, ST-2 | diff --git a/docs/invariants.md b/docs/invariants.md index 3de4ffd..5dcf7a0 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -240,7 +240,8 @@ proxy that hands the server connection to another client keeps the rewritten, st `search_path` writer under `pkg/`), `pg_catalog.` qualification in `pkg/executor`, `pkg/progress`, `pkg/schemadiff`, and `pkg/checksum` (the column-type read, the relation check, and every function in the digest statement, under a `LocalSearchPath("pg_catalog")` transaction; -decoy `md5` and `format_type` test). *Test obligation:* a shadowing `search_path` (`, +decoy `md5`, `format_type`, and bigint `<=` operator test — the operator is what the local +`search_path` alone can pin). *Test obligation:* a shadowing `search_path` (`, pg_catalog` with decoy catalog relations and functions in the schema) yields the same answer as the default path, per read site and per pooled session. diff --git a/pkg/checksum/guard.go b/pkg/checksum/guard.go index 67da12a..5bb2394 100644 --- a/pkg/checksum/guard.go +++ b/pkg/checksum/guard.go @@ -57,14 +57,24 @@ func (v *Verifier) guard(ctx context.Context, tx pgx.Tx) error { return v.confirmRelations(ctx, tx) } -// setVerifySession bounds the transaction, restricts its search_path to the -// catalog so every unqualified operator, cast, and type name resolves there -// (CO-9), and puts it in the owner's shoes. SET LOCAL cannot take bind -// parameters; the timeouts are integer milliseconds. +// setVerifySession bounds the transaction, pins the one output setting that +// can merge two distinct values, restricts its search_path to the catalog +// so every unqualified operator, cast, and type name resolves there (CO-9), +// and puts it in the owner's shoes. The digest hashes each row's text +// rendering, and for float4, float8, and the geometric types that +// rendering follows extra_float_digits: at zero or below the server rounds +// to fifteen significant digits, so two floats that differ in their last +// digits would render, and hash, the same. The maximum of 3 renders every +// float exactly whatever the database or role configures. The other output +// settings (DateStyle, IntervalStyle, TimeZone, bytea_output) change the +// spelling of a value but never make two values spell the same, and both +// sides share the session. SET LOCAL cannot take bind parameters; the +// timeouts are integer milliseconds. func setVerifySession(ctx context.Context, tx pgx.Tx, owner string, opts Options) error { // INV: LK-2 budgets := "SET LOCAL lock_timeout = " + strconv.FormatInt(opts.LockTimeout.Milliseconds(), 10) + "; SET LOCAL statement_timeout = " + strconv.FormatInt(opts.StatementTimeout.Milliseconds(), 10) + + "; SET LOCAL extra_float_digits = 3" + "; " + dbconn.LocalSearchPath("pg_catalog") if _, err := tx.Exec(ctx, budgets); err != nil { return fmt.Errorf("set verification budgets: %w", err) diff --git a/pkg/checksum/verifier_integration_test.go b/pkg/checksum/verifier_integration_test.go index e6cd3bd..09f2d18 100644 --- a/pkg/checksum/verifier_integration_test.go +++ b/pkg/checksum/verifier_integration_test.go @@ -196,13 +196,15 @@ func TestVerifierLocatesAChangedShadowRow(t *testing.T) { assert.NotEqual(t, found.Source.Hash, found.Shadow.Hash) } -// A row missing from the shadow and a row the shadow holds that the source +// Rows missing from the shadow and a row the shadow holds that the source // does not are each reported in their own chunk, in key order, and the row -// counts on the two sides say which way each chunk differs. +// counts on the two sides say which way each chunk differs. Two rows go +// missing and one is added, so the shadow holds fewer rows than the source +// and the report's row count can only be the source's. func TestVerifierCountsMissingAndExtraShadowRows(t *testing.T) { f := newVerifierFixture(t) target, lock, shadow := f.prepare(t) - f.exec(t, "DELETE FROM "+f.shadowName(shadow)+" WHERE id = 5") + f.exec(t, "DELETE FROM "+f.shadowName(shadow)+" WHERE id IN (5, 6)") f.exec(t, "INSERT INTO "+f.shadowName(shadow)+" (id, qty) VALUES (2600, 2600)") report, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) @@ -211,7 +213,7 @@ func TestVerifierCountsMissingAndExtraShadowRows(t *testing.T) { missing, extra := report.Mismatches[0], report.Mismatches[1] assert.Equal(t, chunk(t, math.MinInt64, 1000), missing.Chunk) assert.Equal(t, int64(1000), missing.Source.Rows) - assert.Equal(t, int64(999), missing.Shadow.Rows) + assert.Equal(t, int64(998), missing.Shadow.Rows) assert.Equal(t, chunk(t, 2001, math.MaxInt64), extra.Chunk) assert.Equal(t, int64(500), extra.Source.Rows) assert.Equal(t, int64(501), extra.Shadow.Rows) @@ -236,6 +238,80 @@ func TestVerifierComparesThroughTheShadowsTypes(t *testing.T) { assert.Equal(t, int64(rows), report.Rows) } +// A float that differs from its source only beyond the fifteenth +// significant digit is a mismatch whatever extra_float_digits the database +// configures: the digest hashes each row's text rendering, and at +// extra_float_digits <= 0 that rendering rounds to fifteen digits, so 0.3 +// and 0.1 + 0.2 would spell, and hash, the same. The pass pins the setting +// itself, so a pool that connects to a database configured at 0 still finds +// the row. The database setting applies at connect, so the pass runs on a +// pool opened after it is set. +func TestVerifierHashesFloatsAtFullPrecision(t *testing.T) { + f := newVerifierFixture(t) + f.exec(t, ` + CREATE TABLE %s.readings ( + id bigint PRIMARY KEY, + v double precision NOT NULL, + note text + )`) + f.exec(t, fmt.Sprintf(` + INSERT INTO %%s.readings (id, v, note) + SELECT n, 0.3, 'reading ' || n FROM generate_series(1, %d) AS n`, rows)) + target := f.prove(t, "readings") + lock := f.lock(t, "readings") + shadow := f.build(t, lock, target, `ALTER TABLE %s DROP COLUMN note`) + f.copy(t, target, shadow, lock) + f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET v = 0.1::float8 + 0.2::float8 WHERE id = 7") + f.exec(t, `DO $$ BEGIN EXECUTE format('ALTER DATABASE %I SET extra_float_digits = 0', current_database()); END $$`) + t.Cleanup(func() { + _, err := f.pool.Exec(context.WithoutCancel(t.Context()), `DO $$ BEGIN EXECUTE format('ALTER DATABASE %I RESET extra_float_digits', current_database()); END $$`) + assert.NoError(t, err) + }) + rounding, err := dbconn.NewPool(t.Context(), f.cfg) + require.NoError(t, err) + t.Cleanup(rounding.Close) + var digits string + require.NoError(t, rounding.QueryRow(t.Context(), "SHOW extra_float_digits").Scan(&digits)) + require.Equal(t, "0", digits, "the database setting must reach the session the pass reads on") + + report, err := f.verify(t, rounding, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + require.Len(t, report.Mismatches, 1) + assert.Equal(t, chunk(t, math.MinInt64, 1000), report.Mismatches[0].Chunk) + assert.Equal(t, report.Mismatches[0].Source.Rows, report.Mismatches[0].Shadow.Rows, "the rows differ in value, not in number") +} + +// A chunk's two digests describe one snapshot (CO-1): a shadow row another +// session changes after the chunk's source digest has been read, and before +// its shadow digest is, is not seen by that chunk, so the pass compares +// clean. Under a snapshot per statement the shadow digest would see the +// change and report the chunk. The change is made from a query hook on the +// pool the pass reads on, which fires once, when the source digest of the +// chunk holding the key completes. The next pass takes new snapshots and +// finds the row. +func TestVerifierDigestsBothSidesInOneSnapshot(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + sourceDigest := " FROM " + pgx.Identifier{f.schema, "orders"}.Sanitize() + " WHERE " + var once sync.Once + hooked := f.hookedPool(t, func(sql string) { + if !strings.HasPrefix(sql, "SELECT pg_catalog.count(*)") || !strings.Contains(sql, sourceDigest) { + return + } + once.Do(func() { f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 500") }) + }) + + report, err := f.verify(t, hooked, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + assert.True(t, report.Clean(), "mismatches: %+v", report.Mismatches) + assert.Equal(t, 3, report.Chunks) + + report, err = f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + require.Len(t, report.Mismatches, 1, "the change landed; a pass with new snapshots sees it") + assert.Equal(t, chunk(t, math.MinInt64, 1000), report.Mismatches[0].Chunk) +} + // A pass compares only the keys at or below the watermark it is given: the // chunk straddling the watermark is clamped to it, a difference below is // found, and a difference above — in a chunk the copier has not finished — @@ -257,10 +333,14 @@ func TestVerifierStopsAtTheWatermark(t *testing.T) { // The pass resolves every catalog object it names through pg_catalog, so a // session whose search_path puts a schema of impostors first (CO-9) neither -// misreads the shadow's types nor hashes with someone else's md5: an -// impostor format_type that would make every numeric compare as text -// produces no false mismatch, and an impostor md5 that answers the same for -// every row hides no real one. +// misreads the shadow's types nor hashes with someone else's md5 nor +// compares keys with someone else's operator: an impostor format_type that +// would make every numeric compare as text produces no false mismatch, an +// impostor md5 that answers the same for every row hides no real one, and +// an impostor bigint <= that is never true, which would empty the key range +// on both sides and compare nothing clean, is not the <= that BETWEEN +// resolves to. Functions are qualified in the statement itself; the +// operator can only be pinned by the transaction's own search_path. func TestVerifierIgnoresTheSessionSearchPath(t *testing.T) { f := newVerifierFixture(t) f.createOrders(t) @@ -270,6 +350,8 @@ func TestVerifierIgnoresTheSessionSearchPath(t *testing.T) { f.copy(t, target, shadow, lock) f.exec(t, `CREATE FUNCTION %s.md5(text) RETURNS text LANGUAGE sql IMMUTABLE AS 'SELECT ''impostor''::text'`) f.exec(t, `CREATE FUNCTION %s.format_type(oid, integer) RETURNS text LANGUAGE sql STABLE AS 'SELECT ''text''::text'`) + f.exec(t, `CREATE FUNCTION %s.never_le(bigint, bigint) RETURNS boolean LANGUAGE sql IMMUTABLE AS 'SELECT false'`) + f.exec(t, `CREATE OPERATOR %s.<= (LEFTARG = bigint, RIGHTARG = bigint, FUNCTION = %s.never_le)`) shadowing := testutil.NewCatalogShadowingPool(t, f.cfg.URL, f.schema) report, err := f.verify(t, shadowing, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) @@ -321,6 +403,50 @@ func TestVerifierRefusesReplacedRelations(t *testing.T) { }) } +// A shadow that has lost a column the proof lists for copy is not the shadow +// the proof describes, even though it is still the same relation (ST-6): +// the pass refuses before any digest, naming the column, rather than +// comparing the columns that remain. +func TestVerifierRefusesAShadowMissingACopyColumn(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + f.exec(t, "ALTER TABLE "+f.shadowName(shadow)+" DROP COLUMN qty") + + _, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.ErrorIs(t, err, checksum.ErrInvariantViolation) + assert.Contains(t, err.Error(), "(ST-6): shadow") + assert.Contains(t, err.Error(), "no column qty") +} + +// Every read runs with the source owner's privileges, not the connected +// role's (SET LOCAL ROLE owner): a source the owner may no longer read is +// not read with the superuser's power behind the pool, and the server's +// refusal is the error. The owner is a throwaway role whose own SELECT on +// the table is revoked after the shadow is built and copied; the superuser +// pool would read it regardless. +func TestVerifierReadsAsTheTableOwner(t *testing.T) { + f := newVerifierFixture(t) + owner := testutil.NewRole(t, f.pool, "NOLOGIN") + f.exec(t, "GRANT USAGE, CREATE ON SCHEMA %s TO "+pgx.Identifier{owner}.Sanitize()) + f.createOrders(t) + f.exec(t, "ALTER TABLE %s.orders OWNER TO "+pgx.Identifier{owner}.Sanitize()) + target := f.prove(t, "orders") + require.Equal(t, owner, target.OwnerRole(), "the proof names the owner the pass reads as") + lock := f.lock(t, "orders") + shadow := f.build(t, lock, target, `ALTER TABLE %s DROP COLUMN note`) + f.copy(t, target, shadow, lock) + + report, err := f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.NoError(t, err) + assert.True(t, report.Clean(), "mismatches: %+v", report.Mismatches) + + f.exec(t, "REVOKE SELECT ON %s.orders FROM "+pgx.Identifier{owner}.Sanitize()) + _, err = f.verify(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + var pgErr *pgconn.PgError + require.ErrorAs(t, err, &pgErr) + assert.Equal(t, "42501", pgErr.Code, "insufficient_privilege: the owner, not the superuser, was refused") +} + // The verifier trusts the lock session only as far as the server confirms // it from the reading connection (LK-1): a session whose backend is gone // compares nothing. A session for a different table, or one that already @@ -449,3 +575,33 @@ func (c *hookedClock) Now() time.Time { c.once.Do(c.hook) return time.Now() } + +// hookedPool is a pool on the fixture's server whose every query, once its +// result has been read, runs hook with the query's SQL on the calling +// goroutine, so a test can act between two statements of one transaction. +func (f verifierFixture) hookedPool(t *testing.T, hook func(sql string)) *pgxpool.Pool { + t.Helper() + pc, err := pgxpool.ParseConfig(f.cfg.URL) + require.NoError(t, err) + pc.ConnConfig.Tracer = queryHook(hook) + pool, err := pgxpool.NewWithConfig(t.Context(), pc) + require.NoError(t, err) + t.Cleanup(pool.Close) + return pool +} + +// queryHook is a pgx query tracer that runs after each query completes, +// carrying the SQL from the query's start to its end through the context. +type queryHook func(sql string) + +type hookedSQLKey struct{} + +func (queryHook) TraceQueryStart(ctx context.Context, _ *pgx.Conn, data pgx.TraceQueryStartData) context.Context { + return context.WithValue(ctx, hookedSQLKey{}, data.SQL) +} + +func (h queryHook) TraceQueryEnd(ctx context.Context, _ *pgx.Conn, _ pgx.TraceQueryEndData) { + if sql, ok := ctx.Value(hookedSQLKey{}).(string); ok { + h(sql) + } +} From 95dcb2c3c7717194b501013868aebb257f097a4b Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Thu, 1 Oct 2026 12:35:19 +1000 Subject: [PATCH 3/3] checksum: test the read-only transaction and pin extra_float_digits per pool The full-precision float test set extra_float_digits with ALTER DATABASE, which reaches every later session on the server, including other packages' tests on a shared PG_DSN, and clobbers a pre-existing value. It now opens the pass's pool with the setting as a connection parameter, so the setting reaches no other pool and nothing is left behind. A new test proves the verify transaction is read-only: a write issued on the transaction's own connection between a chunk's two digests is refused with read_only_sql_transaction and the pass fails. The query hook now hands the connection it fired on so a test can act inside the transaction, not only beside it. Dropping AccessMode: pgx.ReadOnly from begin fails the test. --- pkg/checksum/verifier_integration_test.go | 90 +++++++++++++++++------ 1 file changed, 69 insertions(+), 21 deletions(-) diff --git a/pkg/checksum/verifier_integration_test.go b/pkg/checksum/verifier_integration_test.go index 09f2d18..a5223ef 100644 --- a/pkg/checksum/verifier_integration_test.go +++ b/pkg/checksum/verifier_integration_test.go @@ -6,6 +6,7 @@ import ( "math" "strings" "sync" + "sync/atomic" "testing" "time" @@ -243,9 +244,10 @@ func TestVerifierComparesThroughTheShadowsTypes(t *testing.T) { // configures: the digest hashes each row's text rendering, and at // extra_float_digits <= 0 that rendering rounds to fifteen digits, so 0.3 // and 0.1 + 0.2 would spell, and hash, the same. The pass pins the setting -// itself, so a pool that connects to a database configured at 0 still finds -// the row. The database setting applies at connect, so the pass runs on a -// pool opened after it is set. +// itself, so a pool whose every session starts at 0 still finds the row. +// The setting is a connection parameter of the pool the pass reads on, not +// a database setting, so nothing outlives the test or reaches another one +// on the same server. func TestVerifierHashesFloatsAtFullPrecision(t *testing.T) { f := newVerifierFixture(t) f.exec(t, ` @@ -262,17 +264,10 @@ func TestVerifierHashesFloatsAtFullPrecision(t *testing.T) { shadow := f.build(t, lock, target, `ALTER TABLE %s DROP COLUMN note`) f.copy(t, target, shadow, lock) f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET v = 0.1::float8 + 0.2::float8 WHERE id = 7") - f.exec(t, `DO $$ BEGIN EXECUTE format('ALTER DATABASE %I SET extra_float_digits = 0', current_database()); END $$`) - t.Cleanup(func() { - _, err := f.pool.Exec(context.WithoutCancel(t.Context()), `DO $$ BEGIN EXECUTE format('ALTER DATABASE %I RESET extra_float_digits', current_database()); END $$`) - assert.NoError(t, err) - }) - rounding, err := dbconn.NewPool(t.Context(), f.cfg) - require.NoError(t, err) - t.Cleanup(rounding.Close) + rounding := f.poolWithRuntimeParam(t, "extra_float_digits", "0") var digits string require.NoError(t, rounding.QueryRow(t.Context(), "SHOW extra_float_digits").Scan(&digits)) - require.Equal(t, "0", digits, "the database setting must reach the session the pass reads on") + require.Equal(t, "0", digits, "the connection parameter must reach the session the pass reads on") report, err := f.verify(t, rounding, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) require.NoError(t, err) @@ -292,10 +287,9 @@ func TestVerifierHashesFloatsAtFullPrecision(t *testing.T) { func TestVerifierDigestsBothSidesInOneSnapshot(t *testing.T) { f := newVerifierFixture(t) target, lock, shadow := f.prepare(t) - sourceDigest := " FROM " + pgx.Identifier{f.schema, "orders"}.Sanitize() + " WHERE " var once sync.Once - hooked := f.hookedPool(t, func(sql string) { - if !strings.HasPrefix(sql, "SELECT pg_catalog.count(*)") || !strings.Contains(sql, sourceDigest) { + hooked := f.hookedPool(t, func(_ context.Context, _ *pgx.Conn, sql string) { + if !f.isSourceDigest(sql) { return } once.Do(func() { f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 500") }) @@ -312,6 +306,44 @@ func TestVerifierDigestsBothSidesInOneSnapshot(t *testing.T) { assert.Equal(t, chunk(t, math.MinInt64, 1000), report.Mismatches[0].Chunk) } +// The transaction a chunk's digests run in is read-only (CO-9): a write +// issued on that transaction's own connection, between its two digests, is +// refused by the server with read_only_sql_transaction rather than landing +// on the shadow, and the pass fails because the transaction is aborted. The +// write is attempted from a query hook on the pool the pass reads on, on +// the connection the hook is handed, so it runs inside the verify +// transaction and not beside it. +func TestVerifierTransactionRefusesWrites(t *testing.T) { + f := newVerifierFixture(t) + target, lock, shadow := f.prepare(t) + var attempted atomic.Bool + var writeErr error + hooked := f.hookedPool(t, func(ctx context.Context, conn *pgx.Conn, sql string) { + if !f.isSourceDigest(sql) || !attempted.CompareAndSwap(false, true) { + return + } + _, writeErr = conn.Exec(ctx, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 500") + }) + + _, err := f.verify(t, hooked, target, shadow, lock, copier.NewWatermark(math.MaxInt64)) + require.Error(t, err, "a transaction with a refused statement in it cannot finish the pass") + require.True(t, attempted.Load(), "the write must have been attempted inside the verify transaction") + var pgErr *pgconn.PgError + require.ErrorAs(t, writeErr, &pgErr) + assert.Equal(t, "25006", pgErr.Code, "read_only_sql_transaction") + + var qty int + require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT qty FROM "+f.shadowName(shadow)+" WHERE id = 500").Scan(&qty)) + assert.NotEqual(t, 0, qty, "the refused write left the shadow row as the copier wrote it") +} + +// isSourceDigest reports whether sql is the statement that digests a chunk +// of the fixture's orders table, the first of a chunk's two reads. +func (f verifierFixture) isSourceDigest(sql string) bool { + sourceDigest := " FROM " + pgx.Identifier{f.schema, "orders"}.Sanitize() + " WHERE " + return strings.HasPrefix(sql, "SELECT pg_catalog.count(*)") && strings.Contains(sql, sourceDigest) +} + // A pass compares only the keys at or below the watermark it is given: the // chunk straddling the watermark is clamped to it, a difference below is // found, and a difference above — in a chunk the copier has not finished — @@ -576,10 +608,26 @@ func (c *hookedClock) Now() time.Time { return time.Now() } +// poolWithRuntimeParam is a pool on the fixture's server whose every session +// starts with the named setting at value, as a connection parameter: the +// setting reaches no other pool and nothing is left behind on the server. +func (f verifierFixture) poolWithRuntimeParam(t *testing.T, name, value string) *pgxpool.Pool { + t.Helper() + pc, err := pgxpool.ParseConfig(f.cfg.URL) + require.NoError(t, err) + pc.ConnConfig.RuntimeParams[name] = value + pool, err := pgxpool.NewWithConfig(t.Context(), pc) + require.NoError(t, err) + t.Cleanup(pool.Close) + return pool +} + // hookedPool is a pool on the fixture's server whose every query, once its -// result has been read, runs hook with the query's SQL on the calling -// goroutine, so a test can act between two statements of one transaction. -func (f verifierFixture) hookedPool(t *testing.T, hook func(sql string)) *pgxpool.Pool { +// result has been read, runs hook with the query's SQL and the connection +// it ran on, on the calling goroutine, so a test can act between two +// statements of one transaction, on that transaction's own connection or +// on another. +func (f verifierFixture) hookedPool(t *testing.T, hook func(ctx context.Context, conn *pgx.Conn, sql string)) *pgxpool.Pool { t.Helper() pc, err := pgxpool.ParseConfig(f.cfg.URL) require.NoError(t, err) @@ -592,7 +640,7 @@ func (f verifierFixture) hookedPool(t *testing.T, hook func(sql string)) *pgxpoo // queryHook is a pgx query tracer that runs after each query completes, // carrying the SQL from the query's start to its end through the context. -type queryHook func(sql string) +type queryHook func(ctx context.Context, conn *pgx.Conn, sql string) type hookedSQLKey struct{} @@ -600,8 +648,8 @@ func (queryHook) TraceQueryStart(ctx context.Context, _ *pgx.Conn, data pgx.Trac return context.WithValue(ctx, hookedSQLKey{}, data.SQL) } -func (h queryHook) TraceQueryEnd(ctx context.Context, _ *pgx.Conn, _ pgx.TraceQueryEndData) { +func (h queryHook) TraceQueryEnd(ctx context.Context, conn *pgx.Conn, _ pgx.TraceQueryEndData) { if sql, ok := ctx.Value(hookedSQLKey{}).(string); ok { - h(sql) + h(ctx, conn, sql) } }