checksum: run the verifier under an explicit divergence policy and mint the proofs (CO-2, CO-3) - #136
Conversation
…nt the proofs (CO-2, CO-3) Check runs one pass under a DivergencePolicy the caller states every time; the zero value is refused. A clean pass mints a CleanWatermark, and a VerifiedShadow when the watermark is complete. Under abort a difference is a DivergenceError with the shadow untouched. Under repair every differing chunk is replaced in one guarded transaction — delete the shadow's rows over the chunk, then the copier's own chunk statement — and read again in a fresh snapshot; a chunk that still differs is a RepairError, and a pass with repairs mints no proof. copier exports InsertChunk so a repair recopies with the statement the copy used, and Watermark.Complete names the top-of-key-space test.
16eaf22 to
4d39e3e
Compare
|
You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard. |
|
🤖 1/2: adversarial correctness review of 0 blocking, 5 non-blocking. I looked for any way the policy turns a real difference into a pass or a quiet repair, and found none. The zero and unknown policies are refused before anything is read. Non-blocking1. Every error path returns Test that fails on 4d39e3e and passes with the fix// A pass that aborted on a difference does not hand back an Outcome that
// answers Clean() == true.
func TestCheckAbortOutcomeIsNotClean(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")
outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceAbort)
var divergence *checksum.DivergenceError
require.ErrorAs(t, err, &divergence)
assert.False(t, outcome.Clean(), "a pass that found a difference is not clean")
}
2. When a later chunk's repair fails, the result drops the repairs that already committed. Each chunk commits its own repair. Test that fails on 4d39e3e and passes with the fix// When a later chunk's repair does not take, the chunks already recopied
// and committed are still reported.
func TestCheckReportsCommittedRepairsWhenALaterRepairFails(t *testing.T) {
f := newVerifierFixture(t)
target, lock, shadow := f.prepare(t)
f.exec(t, "DELETE FROM "+f.shadowName(shadow)+" WHERE id = 5")
f.exec(t, "UPDATE "+f.shadowName(shadow)+" SET qty = 0 WHERE id = 1500")
f.exec(t, `
CREATE FUNCTION %s.zero_qty() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
NEW.qty := 0;
RETURN NEW;
END
$$`)
f.exec(t, "CREATE TRIGGER zero_qty BEFORE INSERT ON "+f.shadowName(shadow)+" FOR EACH ROW WHEN (NEW.id > 1000) EXECUTE FUNCTION %s.zero_qty()")
outcome, err := f.check(t, f.pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair)
var failed *checksum.RepairError
require.ErrorAs(t, err, &failed)
assert.Equal(t, chunk(t, 1001, 2000), failed.Repair.Mismatch.Chunk)
_, exists := f.shadowRow(t, shadow, 5)
require.True(t, exists, "the first chunk's repair committed: row 5 is back")
require.Len(t, outcome.Repairs, 1, "the committed repair of [MinInt64, 1000] is reported")
assert.Equal(t, chunk(t, math.MinInt64, 1000), outcome.Repairs[0].Mismatch.Chunk)
}
3. No test pins that the repair transaction runs the guard. Replacing Test that passes on 4d39e3e and fails with the repair's guard removed// A shadow replaced between the comparison pass and the repair is refused
// by the repair transaction's own guard before it writes: the impostor
// that now carries the shadow's name is left empty (ST-6, LK-1).
func TestCheckRepairRefusesAReplacedShadowBeforeWriting(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")
name := f.shadowName(shadow)
// The repair is the pass's only read-write transaction.
pool := f.poolAfter(t, &afterHook{prefix: "begin read write", hook: func() {
f.exec(t, "ALTER TABLE "+name+" RENAME TO displaced")
f.exec(t, "CREATE TABLE "+name+" (LIKE %s.displaced INCLUDING ALL)")
}})
_, err := f.check(t, pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair)
require.ErrorIs(t, err, checksum.ErrInvariantViolation)
assert.Contains(t, err.Error(), "(ST-6)")
var n int64
require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT count(*) FROM "+name).Scan(&n))
assert.Zero(t, n, "the repair wrote nothing into the impostor")
}
// afterHook runs hook once, before the first statement starting with prefix
// that begins after `after` statements containing marker were sent.
type afterHook struct {
marker, prefix string
after int
hook func()
mu sync.Mutex
seen int
once sync.Once
}
func (h *afterHook) TraceQueryStart(ctx context.Context, _ *pgx.Conn, data pgx.TraceQueryStartData) context.Context {
h.mu.Lock()
if strings.Contains(data.SQL, h.marker) {
h.seen++
}
ready := h.seen >= h.after
h.mu.Unlock()
if ready && strings.HasPrefix(data.SQL, h.prefix) {
h.once.Do(h.hook)
}
return ctx
}
func (h *afterHook) TraceQueryEnd(context.Context, *pgx.Conn, pgx.TraceQueryEndData) {}
func (f verifierFixture) poolAfter(t *testing.T, h *afterHook) *pgxpool.Pool {
t.Helper()
pc, err := pgxpool.ParseConfig(f.cfg.URL)
require.NoError(t, err)
pc.ConnConfig.Tracer = h
pool, err := pgxpool.NewWithConfig(t.Context(), pc)
require.NoError(t, err)
t.Cleanup(pool.Close)
return pool
}
4. On a table taking writes, the reread treats a repair that worked as one that did not, so the reconciliation pass that This fails closed, so it is not a false pass. The concern is that CO-3's self-heal cannot be reached. The design runs the repair pass after slot loss, before CDC resumes from the new slot ( CO-2 already withholds the proof until the next fresh pass. So the reread does not need to be the final word on a transient difference. It does need to be the final word on a persistent one. One shape that keeps both: report a reread mismatch as repaired but unverified, and refuse with Test that shows it on 4d39e3e (passes: one source write after the recopy becomes a RepairError)// One application write to a repaired chunk between its recopy and its
// reread turns a repair that took into a RepairError for the whole pass.
func TestCheckRepairOneSourceWriteAfterTheRecopy(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")
// After the repair's INSERT, the next transaction to begin is the reread.
pool := f.poolAfter(t, &afterHook{marker: "INSERT INTO", prefix: "SET LOCAL lock_timeout", after: 1, hook: func() {
f.exec(t, "UPDATE %s.orders SET qty = qty + 1 WHERE id = 1999")
}})
_, err := f.check(t, pool, target, shadow, lock, copier.NewWatermark(math.MaxInt64), checksum.DivergenceRepair)
var failed *checksum.RepairError
require.ErrorAs(t, err, &failed)
qty, _ := f.shadowRow(t, shadow, 1500)
assert.Equal(t, int64(1500), qty, "the recopy did take")
}
5. LK-1's enforcement says " Verified
This review was generated by Claude Code (claude-opus-5). |
|
🤖 2/2: OSS adoption and integration ease, at The shape is easy to adopt. A string-typed 1. Let importers validate a policy when they load config. 2. State on 3. Keep 4. Progress is still dark during 5. Doc accuracy: the verifier's guard is not "the same guard as the copier" on This review was generated by Claude Code (claude-opus-5). |
aparajon
left a comment
There was a problem hiding this comment.
🤖 Approving 4d39e3e with 0 blocking findings. The 1/2 comment has five non-blocking findings. Three of them come with tests: a refused Check returns an outcome whose Clean() is true, repairs that already committed are dropped from a failed result, and nothing pins the repair transaction's guard. The 2/2 comment has five non-blocking integration notes.
This stamp was left by Claude Code (claude-opus-5).
…mitted repairs, share the chunk insert
Review follow-up for the divergence-policy pass.
- Repair recopies every differing chunk in one guarded read-write
transaction: all deletes, then all inserts, one commit. The shadow carries
the source's unique indexes, so a unique value that moved between two
differing chunks made the chunk-by-chunk recopy fail with unique_violation
on every pass; one transaction also keeps a half-repaired shadow invisible.
Rereads still run per chunk in fresh snapshots.
- Check returns Outcome{Report, Repairs} alongside DivergenceError,
RepairError, and LK-1, so the repairs that committed before a later
failure are reported; every refusal is wrapped as "verify schema.table".
- Outcome.Clean() means "this pass minted the clean watermark"; a refused,
aborted, or repaired pass reads not clean.
- A repair pass states its precondition on Check and in doc.go: no other
shadow writer, and source rows in differing chunks hold still until the
rereads. A source write inside the pass fails it closed. The low-level
design's reconciliation bullet now agrees that the pass runs before the
new slot's stream is applied.
- ParseDivergencePolicy lets an importer refuse a bad setting at config
load with the same ErrNoDivergencePolicy.
- pkg/internal/chunksql holds the one chunk-insert statement the copier and
the repair both build from; copier.InsertChunk is removed. depguard allows
pkg/internal for core files; SAFETY.md gains the package row.
- The copier pins search_path to pg_catalog in the copy transaction, so a
repair runs the shared statement under the same session as the copy and
the "same guard" claim in SAFETY.md and the design is exact.
Tests on PostgreSQL: committed repairs reported when a later chunk's repair
fails; a unique value moved between chunks converges; a shadow replaced as
the repair transaction begins is refused as ST-6 with the impostor empty; a
source write inside the pass is refused as RepairError; the copier ignores
the session search_path. Unit: the frozen chunk-insert statement,
ParseDivergencePolicy, the zero Outcome is not clean.
|
🤖 Adversarial review response — created by Kiran's code review agent (Amp, Claude Opus 4.6) — pull/136, follow-up commit Verdict: no blocking findings; all five correctness findings are fixed (two of them as the reviewer proposed, one by a repair-pass precondition rather than a retry); of the five integration-lens notes, four are taken (policy parser, godoc precondition,
Tests added or changed in the follow-up commit: Source: block/pg-sprite#136, review comments 5924708298 and 5924709068 at head |
…policy * origin/main: preflight: finish copy-and-swap preflight — dependents and environment (#135) # Conflicts: # docs/copy-and-swap-design.md
Adds
Verifier.Check: the chunk verifier run under an explicit divergence policy, with the repair primitive and the private constructors ofVerifiedShadowandCleanWatermark, so a pass now says what a difference means and what it proves.Why
#134 landed the comparison and a
Reportthat proves nothing. The gate needs three things on top of it: a policy the caller states for every pass, never a default (CO-3 — in the steady state a divergence is a defect and aborts; only reconciliation after a lost slot repairs); a repair that recopies the differing chunks with exactly the statement the copy used, so the two cannot drift; and proofs that only a pass which found nothing and repaired nothing can mint (CO-2 — a repaired chunk was recopied, not verified; the proof comes from the next fresh pass).What
pkg/checksum:DivergencePolicy(DivergenceAbort=abort,DivergenceRepair=repair); the zero value and any other string are refused withErrNoDivergencePolicybefore anything is read.ParseDivergencePolicy(string)gives an importer the same refusal at config load.Check(ctx, pool, through, policy) (Outcome, error)runsVerifyand acts on the report. Clean →Outcomecarrying aCleanWatermarkatthrough, plus aVerifiedShadowwhenthrough.Complete(); a partial clean pass mints only the watermark. Abort + differences →*DivergenceError{Report}, shadow untouched. Repair → every differing chunk is replaced in one guarded read-write transaction — allDELETE FROM shadow WHERE pk BETWEEN $1 AND $2, then all chunk inserts, one commit — and each chunk is then digested again in a fresh snapshot; a chunk that still differs is*RepairError{Repair, After}. One transaction because the shadow carries the source's unique indexes: a unique value that moved between two differing chunks would make a chunk-by-chunk recopy fail withunique_violationon every pass, and a half-repaired shadow is never visible. Every refusalCheckreturns is wrappedverify <schema>.<table>: ….Outcome{Report, Repairs}comes back with the error too, so the repairs that committed before aRepairErroror a lost lock are reported.Outcome.Clean()means "this pass minted the clean watermark";CleanWatermark() (CleanWatermark, bool)andVerifiedShadow() (VerifiedShadow, bool)gate on their flags. The proof constructors are private; the zero values stay forgeable and consumers check the flag.Checkand indoc.go: it assumes no other shadow writer and that source rows in differing chunks hold still until the rereads — the pass after slot loss runs before change capture resumes, and the resumed capture carries the application's writes. A write that lands inside the pass makes the repaired chunk read different again and is refused asRepairErrorrather than repaired twice. The live transient-difference protocol (applier-LSN feed, retries) is its own plan item.begintakespgx.TxOptions, so the repair transaction runs the same guard as the digest transactions (owner role, catalog-onlysearch_path,ACCESS SHAREon both relations, lock confirmation, relation-OID check); a lost lock during a repair cancels it and is reported as LK-1.pkg/internal/chunksql:Insert(schema, source, shadow, key, columns)is the one chunk-insert statement;copier.copySQLand the verifier's repair both build from it, so the statement cannot drift and nothing outside the module can call it. SAFETY.md and.golangci.yml(depguard) carry the package.pkg/copier: the copy transaction now pinssearch_pathtopg_catalog, so the repair runs the shared statement under the same session as the copy and "the same guard as the copier" is exact in the docs.Watermark.Complete()names the top-of-key-space test the chunker already made.invariants.mdCO-1 / CO-2 / CO-3 "enforced today",copy-and-swap-design.mdpackage map and D14 "where enforced",low-level-design.mdreconciliation bullet (the pass runs before the new slot's stream is applied),SAFETY.mdandarchitecture.mdpkg/checksumrows.BEFORE INSERTtrigger that zeroesqtymakes the repair not take and is refused asRepairError, with the committed recopy left as the trigger wrote it; the same trigger on keys above 1000 only, with two differing chunks, reports both committed repairs and names the second; a unique value swapped between ids 100 and 1500 on the source converges; a shadow renamed away and replaced by an impostor as the repair transaction begins is refused as ST-6 with the impostor empty; a source write between the recopy and the reread is refused asRepairError; the zero and an unknown policy are refused with no pass run; losing the lock as the repair'sDELETEstarts (injected through a pool tracer, deterministic) is reported as LK-1 and the repair rolls back; the copier ignores a sessionsearch_paththat shadows the bigint<=operator. Unit: policy validation and parsing, the zeroOutcomeis not clean and mints nothing,provemints the shadow proof only at the complete watermark, the frozen repair and chunk-insert statements,Watermark.Completeat both sides of the boundary. Each new test fails under the mutation it exists to catch (chunk-by-chunk recopy, unguarded repair transaction, repairs dropped on error,Clean()reading the report, copier without the pin).Before / after
Deferred (follow-up)
A checksum
progress.WorkSourcefiller needs anOperationChecksumand counters for chunks compared, rows hashed, and chunks repaired in the progress contract (rows_copiedis the wrong name); it lands with that contract change, before an orchestrator wiresCheck.References
docs/invariants.mdCO-1, CO-2, CO-3;docs/copy-and-swap-design.mdD14.🤖 Drafted with Amp (Claude Opus 4.6); reviewed and edited by the author.