Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion SAFETY.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
2 changes: 1 addition & 1 deletion docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
5 changes: 3 additions & 2 deletions docs/copy-and-swap-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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 |
Expand Down
19 changes: 16 additions & 3 deletions docs/invariants.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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` (`<schema>,
`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` (`<schema>,
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.

Expand Down Expand Up @@ -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
Expand Down
112 changes: 112 additions & 0 deletions pkg/checksum/digest.go
Original file line number Diff line number Diff line change
@@ -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
}
Loading
Loading