diff --git a/SAFETY.md b/SAFETY.md index a83f8c2..39695e8 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, and the copy step's `progress.WorkSource` — rows from committed chunks, the source's catalog row count, both tables' measured sizes) exist | CO-4 (chunk coverage, copy SQL shape, in-flight registry), LK-1, LK-3 | | `pkg/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 d96e894..85e9f95 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -201,7 +201,7 @@ different levels of commitment: | `pkg/executor` | Native backend with stable outcome codes: the bounded optimistic attempt, the concurrent index build and its invalid-index recovery, the autocommit safer-sequence runner, the greenfield `CREATE TABLE` path, and the accepted-blocking passthrough primitive; the full `Executor` contract (`Plan`/`Execute`/`Status`/`Abort`) arrives with the copy-and-swap backend | native execution exists | | `pkg/progress` | Strategy-wide, pollable progress snapshots: native phase/elapsed time, sequence position, retry attempt, and server-reported concurrent-index work, and engine-measured copy work polled from the copier (`WorkSource`) | native progress and the copy work source exist | | `pkg/copier` | PK-range chunker over one integer-family primary key with dynamic time-based sizing (produces `Chunk` and `Watermark`; composite keys refused in v1), and the parallel chunked copy into the shadow table (never overwrites; reports the cut frontier, in-flight chunks, and landed watermark for the applier, and fills the progress tracker's copy counters while it runs) — there is no separate chunker package | chunker, copy loop, and progress fillers exist | -| `pkg/checksum` | The mandatory correctness gate; 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 1db030a..8f3155b 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. 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` | 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 — 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 2caabd0..5dcf7a0 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,10 @@ 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`, `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. @@ -284,7 +294,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..5bb2394 --- /dev/null +++ b/pkg/checksum/guard.go @@ -0,0 +1,166 @@ +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, 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) + } + 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..a5223ef --- /dev/null +++ b/pkg/checksum/verifier_integration_test.go @@ -0,0 +1,655 @@ +package checksum_test + +import ( + "context" + "fmt" + "math" + "strings" + "sync" + "sync/atomic" + "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) +} + +// 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. 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 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)) + 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(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) + 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 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 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, ` + 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") + 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 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) + 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) + var once sync.Once + 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") }) + }) + + 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) +} + +// 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 — +// 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 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) + 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'`) + 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)) + 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)") + }) +} + +// 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 +// 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() +} + +// 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 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) + 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(ctx context.Context, conn *pgx.Conn, 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, conn *pgx.Conn, _ pgx.TraceQueryEndData) { + if sql, ok := ctx.Value(hookedSQLKey{}).(string); ok { + h(ctx, conn, sql) + } +} 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()) +}