Skip to content

copier: parallel chunked copy step under the table lock (CO-4, LK-3) - #128

Merged
Kiran01bm merged 5 commits into
mainfrom
kiran01bm/cs5-copy-step
Sep 29, 2026
Merged

Kiran01bm merged 5 commits into
mainfrom
kiran01bm/cs5-copy-step

Conversation

@Kiran01bm

@Kiran01bm Kiran01bm commented Sep 26, 2026 •

Copy link
Copy Markdown
Collaborator

Adds copier.Copier, the parallel chunked copy from a proven source into its built shadow, and promotes the in-transaction table-lock confirmation into pkg/dbconn so the copier and the shadow builder share it.

Why

The chunker (previous PR in this stack) cuts consecutive key ranges but nothing copied them. The applier that follows needs more than a watermark: with several workers, chunks land out of order, so a captured change must be judged per key as uncut (discard), in flight (defer), or landed (apply). Only the copier knows which chunks are in flight, so it owns that answer.

The shadow builder already re-asserted the lock from inside its own transaction; the copier needs the identical check on every chunk connection. One implementation in dbconn keeps LK-1 enforced in one place.

What

  • pkg/dbconn: (*TableLockSession).Confirm(ctx, conn) — nobody holds the lock → the new ErrTableLockNotHeld naming the table; another backend holds it → *TableLockHeldError naming it. Confirm reports what the server said; each writer wraps it once as its own ErrInvariantViolation (LK-1), so the message never states the same fact twice. schemachange.confirmTableLock maps those onto its existing CauseLockUnconfirmed / CauseLockHeldElsewhere refusals (existing tests unchanged).
  • pkg/copier:
    • Shadow interface (schema, source, shadow, both OIDs, copy columns) — the shape schemachange.BuiltShadow satisfies, so pkg/schemachange (builder and orchestrator) can import the copier without a cycle. NewCopier refuses a shadow that is not the target's, or that is the source itself by name or by OID (ST-6) and a lock session that is missing, for another table, or already lost (LK-1).
    • Copier.Run(ctx, pool): N workers (default 4) under the lock session's Bind context. Each chunk runs in its own transaction: SET LOCAL lock_timeout/statement_timeout, SET LOCAL ROLE owner, lock.Confirm (only "no holder" / "another backend" is LK-1; a lookup that got no answer is reported as the connection's error), LOCK TABLE source, shadow IN ACCESS SHARE MODE (so neither relation can be replaced for the rest of the transaction; a name that no longer resolves is ST-6), re-resolve both relation OIDs against the proof, then one frozen statement INSERT INTO shadow (cols) SELECT cols FROM source WHERE pk BETWEEN $1::bigint AND $2::bigint ON CONFLICT (pk) DO NOTHING. Elapsed time from an injected progress.Clock feeds Chunker.Feedback (D12).
    • Ledger + Position: cutting and registering a chunk happen under one lock, and the ledger refuses a chunk that does not start just above the cut frontier (CO-4 fail-closed), so the frontier never runs ahead of an unregistered chunk; a chunk is registered in flight before its transaction begins and leaves the in-flight set only when it commits — a chunk whose transaction did not commit stays in InFlight, so its keys never read as landed; the watermark advances over the contiguous landed prefix; Position.Classify(key) returns KeyUncut / KeyInFlight / KeyLanded — the CO-4 rule the applier will call. Run returns only after every worker has exited, so the copier holds no chunk transaction when a caller checkpoints the watermark (LK-3, client side; a cancelled statement ends on the server when it finishes or hits statement_timeout), then fails closed if the key space is not covered.
    • Resume: a copy resumed after watermark W first deletes every shadow row above W in bounded batches (DELETE … WHERE pk IN (SELECT pk … WHERE pk >= $1 ORDER BY pk LIMIT $2), $1 being the first key the resumed chunker will cut, each batch in the same guarded transaction shape), before its first chunk is cut. An earlier run may have landed chunks above W out of order, and the applier discards changes for keys above the resumed cut frontier, so a stale shadow row left there would survive ON CONFLICT DO NOTHING. A zero watermark means nothing landed, not that the shadow is empty — a run whose first chunk never landed may have landed chunks anywhere above it — so the clear then covers the whole shadow; a first run finds it empty and its one batch removes nothing. The first clear batch takes LOCK TABLE shadow IN SHARE MODE before its delete, so a straggling chunk transaction from the earlier run — one whose INSERT is still open when the resume begins — commits (or aborts) before the clear reads the shadow, and a row it commits cannot slip in behind the delete. The fence also waits on any other open writer of the shadow and blocks new writes for the length of that one batch, bounded by lock_timeout; later batches do not fence.
  • Position.Cut is a Watermark, like Position.Watermark and the ledger's frontier, so "nothing cut" and "cut through key K" have the one spelling the applier and the checkpoint will read (Watermark.Valid / Key); CutValid is gone.
  • Not in this PR: typed refusal causes for the copier's failures (they stay ErrInvariantViolation wrapping the cause until the orchestrator PR maps them onto the schema-change refusal classes; docs/refusal-classes.md says so).
  • Not in this PR: the copy statement carries no conversion expression, so a column whose type differs between source and shadow takes the server's assignment cast. An ALTER COLUMN TYPE … USING change needs a conversion-aware copy before the planner's copy-and-swap route (still refused today) is wired to the copier; copySQL and the design package map say so.
  • Docs in the same PR: SAFETY.md copier row, invariants.md CO-4 / LK-1 / LK-3 Enforced today (LK-3 scoped to the client side), design package map (dbconn, copier) and D12, architecture.md.

Tests (real PostgreSQL, PG 14/16/18): whole-table copy with 3 workers and 100-row chunks converges via testutil.AssertConverged while a sampler asserts every Position snapshot is consistent (lowest in-flight chunk starts at watermark+1, nothing in flight ⇒ watermark = cut); a shadow row the applier commits while the copier's insert of that chunk is waiting on it is never overwritten; resume from a watermark removes stale shadow rows above it, keeps the row below it, and copies every key above it from the source; resume from the zero watermark removes stale rows down to the smallest int64 key and converges; cancellation with one chunk pinned mid-insert by an uncommitted shadow row returns context.Canceled with exactly the pinned chunk in flight (Classify says in-flight for its keys, landed for a chunk that landed above it), every key ≤ watermark present, and a resumed copier clears the tail, picks up a source change made in it meanwhile, and converges; lock loss mid-copy returns ErrInvariantViolation wrapping the session's loss with the interrupted chunk still in flight; a dropped-and-recreated shadow or source is refused (ST-6) with zero rows written, and a dropped source fails the write transaction's own LOCK TABLE (SQLSTATE 42P01 → ST-6); a gone or rival-held lock is refused with zero rows written, and a session that already reported loss is refused by NewCopier before any connection opens; a lock confirmation that never got the server's answer is not an invariant violation; a shadow dropped and recreated between a chunk's claim and its transaction is refused per chunk (ST-6) with the impostor receiving no rows and that chunk left in flight; every row lands as the table owner (SET LOCAL ROLE, checked through a shadow column defaulting to current_user on a table owned by a NOLOGIN role); a resume whose earlier run still has a chunk transaction open waits on it at the fence (an ungranted ShareLock in pg_locks, cut frontier unchanged, nothing in flight) and reads the committed row afterwards; the copy SQL, the resume-clear SQL, and the fence SQL are frozen as exact strings (TM-2); pure ledger tests for out-of-order landing, frontier-only claims, an unlanded chunk staying in flight, resume, and every claim refused once the key space is cut through the largest key.

Before / after

Before                                   After
┌──────────┐                             ┌──────────┐  Next()   ┌──────────────────────┐
│ Chunker  │  cuts ranges, nobody        │ Chunker  │◀─────────▶│ Copier (N workers)   │
│ Next/Cut │  copies them                │          │ Feedback  │  claim → tx → land   │
└──────────┘                             └──────────┘           └──────────┬───────────┘
                                                                 per chunk │ SET LOCAL … ROLE
┌──────────────┐ confirmTableLock          ┌──────────────┐               │ lock.Confirm   (LK-1)
│ schemachange │ (own pg_locks query)      │ dbconn       │◀──────────────┤ OID check      (ST-6)
└──────────────┘                           │ Confirm()    │               │ INSERT … DO NOTHING
                                           └──────▲───────┘               ▼
                                                  │ maps to causes   ┌──────────┐
                                           ┌──────┴───────┐          │ Position │ Watermark, Cut,
                                           │ schemachange │          │ Classify │ InFlight → applier
                                           └──────────────┘          └──────────┘

References

🤖 Drafted with Amp (Claude Opus 4.6); reviewed and edited by the author.

@Kiran01bm
Kiran01bm force-pushed the kiran01bm/cs5-copy-step branch from cc2bec2 to 3149a5b Compare September 28, 2026 10:33
Base automatically changed from kiran01bm/cs5-copier to main September 28, 2026 11:07
The chunker could cut ranges but nothing copied them; the applier needs
a per-key uncut/in-flight/landed answer that only the copier can give.
Promotes the in-transaction lock check to dbconn so copier and builder
share one LK-1 confirmation.
A resumed copy discards changes for keys above the watermark that an
earlier run may already have copied, so it now deletes that tail first
in guarded batches. Claims are frontier-ordered under one lock, an
uncommitted chunk stays in flight, and both relations are locked before
the OID check.
@Kiran01bm
Kiran01bm force-pushed the kiran01bm/cs5-copy-step branch from 3149a5b to 9119c2e Compare September 28, 2026 11:12
A zero watermark means nothing landed, not that the shadow is empty: a run
whose first chunk never landed may have landed chunks above it. Also refuse
a shadow that is the source itself and stop calling a lost connection LK-1.
@Kiran01bm
Kiran01bm marked this pull request as ready for review September 28, 2026 11:50
@chatgpt-codex-connector

Copy link
Copy Markdown

You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard.

@aparajon

Copy link
Copy Markdown
Collaborator

🤖 Review 1/2: adversarial correctness. Inspected: the whole diff at ca7af5e (copier.go, copy_chunk.go, ledger.go, position.go, resume.go, shadow.go, dbconn/table_lock_confirm.go, the schemachange/lock.go move onto Confirm, and the invariant and design doc edits). Ran go test ./pkg/copier/ ./pkg/dbconn/ ./pkg/schemachange/ against real PostgreSQL (green), then 32 mutants against those suites: 27 killed, 5 survived. Three probe tests below kill three of the survivors. The other two are defensive checks that can't be reached: the in-flight and watermark checks in finish, since every worker returns nil only after every chunk it claimed has landed.

0 blocking, 4 non-blocking.

Invariants:

  • CO-4: extended. The frontier-only claim, the in-flight-until-commit rule, and the resume clear (including the zero-watermark case) are each killed by the existing tests.
  • LK-1: extended to the copier. Every chunk transaction and every clear batch confirms from its own connection.
  • LK-3: extended on the client side. Run waits for every worker, and a chunk whose copy failed stays in InFlight.
  • ST-6: extended. OIDs are re-resolved under ACCESS SHARE, and a shadow that is the source is refused.
  • LK-2: upheld. Sub-millisecond timeouts are refused.
  • TM-2: upheld. Both SQL strings are frozen.

The code reads correct. All four findings are either guards the suite doesn't pin or an ordering edge in the resume premise.

Non-blocking

1. The per-chunk guard is unpinned. Every refusal test is actually refused by the resume clear. Run always calls clearAbove first (copier.go:194), even from the zero watermark, and each clear batch runs the same guard (resume.go:57). So TestCopierRefusesReplacedRelations and TestCopierRefusesAnUnconfirmedLock both fail in the clear, before any chunk is cut. Deleting the guard from copyChunk (copy_chunk.go:36) leaves the whole package green. The PR body's "refused with zero rows written" is true, but it's the clear doing the refusing. The probe replaces the shadow after the clear commits and before the first chunk begins, using an Options.Clock whose first Now() fires between claim and copyChunk:

Test that kills it (mutant: drop c.guard from copyChunk)
// onFirstNow runs hook the first time the copier reads the clock, which a
// worker does after claiming its first chunk and before that chunk's
// transaction begins: after the resume clear has committed.
type onFirstNow struct {
	once sync.Once
	hook func()
}

func (c *onFirstNow) Now() time.Time {
	c.once.Do(c.hook)
	return time.Now()
}

// A shadow replaced after the resume clear, before the first chunk, is
// refused by the chunk transaction's own guard.
func TestProbeChunkGuardRefusesAShadowReplacedAfterTheClear(t *testing.T) {
	f := newCopierFixture(t)
	target, lock, shadow := f.prepare(t, 100)
	shadowName := pgx.Identifier{f.schema, shadow.ShadowTable()}.Sanitize()
	clock := &onFirstNow{hook: func() {
		_, err := f.pool.Exec(context.WithoutCancel(t.Context()), "DROP TABLE "+shadowName+"; CREATE TABLE "+shadowName+" (id bigint PRIMARY KEY, qty integer NOT NULL)")
		assert.NoError(t, err)
	}}
	c, err := copier.NewCopier(target, shadow, lock, copier.Watermark{}, copier.Options{Workers: 1, Clock: clock})
	require.NoError(t, err)
	err = c.Run(t.Context(), f.pool)
	require.ErrorIs(t, err, copier.ErrInvariantViolation)
	assert.Contains(t, err.Error(), "(ST-6): shadow")
	assert.Equal(t, int64(0), f.count(t, shadow.ShadowTable(), "true"), "the impostor receives nothing")
}

On the mutant, the impostor shadow takes the copy and Run returns nil:

--- FAIL: TestProbeChunkGuardRefusesAShadowReplacedAfterTheClear (1.63s)
    zz_probe128_test.go:45:
        	Error:      	Expected error with "invariant violation" in chain but got nil.

2. SET LOCAL ROLE <owner> is unpinned (copy_chunk.go:80). The fixture pool is a superuser and the fixture tables are owned by that same superuser, so switching the role changes nothing the tests can see. Replacing the statement with SELECT 1 leaves the package green. The role is observable in practice: anything in the shadow that reads current_user, such as an audit column's DEFAULT current_user, records the tool's login rather than the owner, and privilege checks run as the wrong role. The probe below moves the table to a separate owner and checks who the copy wrote as:

Test that kills it (mutant: SET LOCAL ROLE → SELECT 1)
// The copy writes as the source's owner: a shadow column whose default
// records current_user names the owner, not the tool's login.
func TestProbeCopyWritesAsTheOwner(t *testing.T) {
	f := newCopierFixture(t)
	owner := "own_" + f.schema
	f.exec(t, "CREATE ROLE "+owner)
	t.Cleanup(func() {
		_, err := f.pool.Exec(context.Background(), "DROP OWNED BY "+owner+" CASCADE; DROP ROLE "+owner)
		assert.NoError(t, err)
	})
	f.createOrders(t, 100)
	f.exec(t, "GRANT USAGE, CREATE ON SCHEMA %s TO "+owner)
	f.exec(t, "ALTER TABLE %s.orders OWNER TO "+owner)
	target := f.prove(t, "orders")
	require.Equal(t, owner, target.OwnerRole())
	lock := f.lock(t, "orders")
	alter, err := statement.ParseOne(fmt.Sprintf(`ALTER TABLE %s ADD COLUMN writer name NOT NULL DEFAULT current_user`, pgx.Identifier{f.schema, "orders"}.Sanitize()))
	require.NoError(t, err)
	shadow, err := schemachange.BuildShadow(t.Context(), f.pool, lock, target, alter, schemachange.Options{})
	require.NoError(t, err)

	c, err := copier.NewCopier(target, shadow, lock, copier.Watermark{}, copier.Options{Workers: 2})
	require.NoError(t, err)
	require.NoError(t, c.Run(t.Context(), f.pool))
	assert.Equal(t, int64(100), f.count(t, shadow.ShadowTable(), "writer = '"+owner+"'"), "every copied row was written as the owner")
}
--- FAIL: TestProbeCopyWritesAsTheOwner (1.64s)
    zz_probe128_test.go:74:
        	Error:      	Not equal:
        	            	expected: 100
        	            	actual  : 0
        	Messages:   	every copied row was written as the owner

3. The cut == MaxInt64 frontier guard is unpinned (ledger.go:66). Without it, l.cut+1 wraps to MinInt64, so a chunk [MinInt64, 10] claimed after the whole space is cut is accepted and the cut moves back to 10. From then on, Classify reports every landed key above 10 as KeyUncut, and the applier would discard changes for rows already copied. The chunker never produces that chunk today, which is why nothing fails. But this guard is the ledger's own CO-4 defense, and a unit test costs nothing:

Test that kills it (mutant: drop the cut == MaxInt64 branch)
// Once the whole key space is cut, no chunk is at the frontier: the
// frontier's successor does not wrap to the smallest key.
func TestProbeLedgerRefusesAClaimAfterTheWholeSpaceIsCut(t *testing.T) {
	l := newLedger(Watermark{})
	require.True(t, l.claim(mustChunk(t, math.MinInt64, math.MaxInt64)))
	assert.False(t, l.claim(mustChunk(t, math.MinInt64, 10)))
	assert.Equal(t, int64(math.MaxInt64), l.position().Cut, "the frontier never moves back")
}
--- FAIL: TestProbeLedgerRefusesAClaimAfterTheWholeSpaceIsCut (0.00s)
    zz_probe128_ledger_test.go:16:
        	Error:      	Should be false
    zz_probe128_ledger_test.go:17:
        	Error:      	Not equal:
        	            	expected: 9223372036854775807
        	            	actual  : 10
        	Messages:   	the frontier never moves back

4. An ambiguous chunk commit from the earlier run can become visible after the resume clear. This is reasoned from PostgreSQL's commit ordering; I have not reproduced it. The scenario:

  • A chunk's Commit (copy_chunk.go:43) fails on the client, for example because losing the lock cancelled the context while COMMIT was in flight. The chunk rightly stays in flight, and the checkpointed watermark stays below it.
  • With synchronous_commit waiting on a standby, the server has committed locally but only publishes the transaction (ProcArrayEndTransaction) after the standby acknowledges it. A backend whose client has gone keeps waiting, because nothing checks the socket during that wait unless client_connection_check_interval is set.
  • A resumed copier's clearAbove (resume.go:26) can therefore run and finish before those rows appear.
  • Once they appear, they sit above the new watermark carrying the old run's snapshot values. The resumed chunk's ON CONFLICT DO NOTHING keeps them.

The applier discarded changes to those keys while they were uncut, so the shadow diverges. That breaks the premise the resume clear states, that every key above W is genuinely uncut. The earlier transaction holds ROW EXCLUSIVE on the shadow until it is published, so a cheap closure is to take LOCK TABLE <shadow> IN SHARE MODE in the first clear batch. That lock waits on any straggling writer from the earlier run, bounded by lock_timeout, and fails closed. Another option is to refuse the resume while any other backend holds a lock on the shadow. The copier isn't wired to a caller yet, and the planned checksum pass would catch the divergence before cutover, so I'm not blocking on this. It's worth settling before the orchestrator starts resuming.

This review was generated by Claude Code (claude-opus-5-5).

@aparajon

Copy link
Copy Markdown
Collaborator

🤖 Review 2/2: OSS adoption and integration ease. Inspected at ca7af5e: the exported surface of pkg/copier (go doc -all ./pkg/copier, sha1 135c2fd6) and the additions to pkg/dbconn (TableLockSession.Confirm, ErrTableLockNotHeld). I read both as a first-time pkg.go.dev reader would, and as an importer such as schemabot or an orchestrator would call them.

0 blocking, 3 non-blocking.

What works well for importers:

  • The dbconn change is purely additive; LookupTableLockHolder is untouched.
  • Confirm gives both failure forms typed shapes (ErrTableLockNotHeld, *TableLockHeldError), so schemachange and the copier now branch on the same two values.
  • KeyState is the right thing to export. An applier needs only Position().Classify(key) and never has to re-derive the uncut, in-flight, and landed boundaries from Cut and Watermark itself.
  • The Run doc is candid that cancellation does not mean the server has finished with the chunk. An importer needs to know that before it acts on Run's return.

Non-blocking

1. Copier refusals carry their cause only in message text (copier.go:115–241, and copy_chunk.go/shadow.go). Every refusal is ErrInvariantViolation, with (LK-1), (ST-6), (CO-4) or (LK-3) in the string. An importer needs a different recovery for each:

  • lock lost: re-acquire the lock and resume from the watermark;
  • source or shadow replaced: rebuild the shadow;
  • ledger contradiction: a bug, so stop.

Only lock loss can be told apart structurally, via errors.Is(err, lock.Err()) and the wrapped dbconn values. The next package over sets the opposite contract. docs/copy-and-swap-design.md says every schemachange refusal is a *RefusalError with a RefusalCause "so an importer routes on the cause rather than on message text". Giving the copier's refusals a cause type, or extending docs/refusal-classes.md to cover them, before a caller exists would stop the first importer from writing strings.Contains(err.Error(), "(ST-6)").

2. Three spellings of "an optional key". Watermark is a type with Valid()/Value(), Position uses a Cut int64 + CutValid bool pair, and Chunker.Cut() returns (upper int64, ok bool). An importer serializing a Position for a checkpoint handles each one differently, and a zero Cut with CutValid forgotten reads as "cut through key 0". Reusing Watermark's shape for the cut frontier (for example Cut Watermark) would give one rule for all three.

3. Options.Clock cites (D12) with no link (copier.go:45). On pkg.go.dev that's an opaque token. It refers to docs/copy-and-swap-design.md#d12--throttle-by-chunk-time-and-slot-lag; a short phrase such as "chunk-time throttling" or the link would make it readable. The same token appears in chunker.go from #127, so it's worth fixing both together.

This review was generated by Claude Code (claude-opus-5-5).

@aparajon aparajon left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤖 Stamped: 0 blocking, 7 non-blocking across the two review comments above.

This stamp was left by Claude Code (claude-opus-5-5).

@morgo morgo left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤖 Automated adversarial review, posted on Morgan Tocker's behalf.

The structure here is right, and the part I went looking hardest at — whether Position can ever lie to the applier — holds up. Cutting and registering under one lock is what makes "the frontier never runs ahead of an unregistered chunk" a property rather than a hope, and keeping the unlanded chunk in flight rather than rolling the frontier back is the correct shape: it fails toward deferral, which is the safe direction for CO-4.

The Chunker.Next mutex note I left on #127 was acted on, and more thoroughly than I suggested — the chunker now splits nextMu (held across the boundary query) from mu (taken only to read the cursor and publish the result), so Cut and Feedback no longer park behind a live query, and Copier.claim mirrors the same split with claimMu/mu. Lock order is claimMu → mu everywhere, never the reverse. Thank you.

Verified rather than assumed:

  • Run's clean exit is not turned into a cancellation by its own defer. return c.finish(context.Cause(ctx)) evaluates the argument before defer stop(nil) fires, so a copy that covered the key space reports nil rather than context.Canceled. That ordering is load-bearing and invisible; it is worth a comment on the defer, because reordering it into cause := context.Cause(ctx) after the defers would break every successful run.
  • A chunk whose Commit failed after the server committed is handled. copyChunk returns the error, the chunk never lands, it stays in flight, and the watermark stops below it — so the committed rows sit above the watermark and the resume clear removes them before they can be read as landed. This is the case that justifies clearAbove existing at all, and it works.
  • Clearing under a live applier is safe. Position already reports Cut as the resume watermark (or invalid) from NewCopier onward, before Run is called, so every key above the watermark classifies as KeyUncut for the entire duration of the clear and the applier discards rather than writing there. Nothing can land in the range being deleted while it is being deleted.
  • A resumed ledger cannot accept a chunk overlapping the landed prefix. newLedger seeds cut from the watermark, so startsAtFrontier demands lower == watermark+1 for the very first claim, and refuses everything once cut == math.MaxInt64.
  • The clear loop terminates. It relies on c.chunker.Rows() being positive — a zero would make removed < batch permanently false and spin forever — and ChunkerOptions.validate rejects a non-positive row count after withDefaults, so Rows() is at least MinRows. Sound, but it is an unstated dependency of resume.go on the chunker's validation.
  • guard's ordering does what its comment claims. LOCK TABLE resolves both names before the OID comparison, and ACCESS SHARE is sufficient because the things that would invalidate the proof — DROP, rename, most rewriting DDL — need ACCESS EXCLUSIVE.

One finding.

The resume clear re-scans from the same lower bound on every batch

batch := c.chunker.Rows()
for {
	removed, err := c.clearBatch(ctx, pool, lower, batch)
	...
	if removed < batch {
		return nil
	}
}

lower is computed once and never moves. Each batch runs

DELETE FROM shadow WHERE key IN (
  SELECT key FROM shadow WHERE key >= $1 ORDER BY key LIMIT $2)

with the same $1, so batch k must walk past everything the previous k-1 batches deleted before it reaches a live row. Deleted index entries are not removed by the delete itself — only by vacuum — so the scan traverses them. The work is quadratic in the number of rows cleared.

What makes this worth fixing rather than noting is that clearAboveSQL's own comment states the property that makes the fix trivially correct:

The subquery orders by the key so each batch removes a contiguous stretch off the primary-key index rather than an arbitrary sample

A contiguous stretch has a highest key. The cursor can advance to it. The chunker sitting next to this code does exactly that — advance(upper) sets next = upper + 1 — and the clear does not.

Failure scenario. A copy of a 50M-row table fails after the first chunk was pinned mid-insert and cancelled while later chunks committed (the situation TestCopierCancellationKeepsTheCancelledChunkInFlight constructs deliberately). Nothing landed contiguously, so the checkpointed watermark is zero, and per this PR's own reasoning the resume must clear the whole shadow. batch is Rows() on a freshly constructed chunker, i.e. DefaultInitialChunkRows = 1000, so that is 50,000 batches. With a bigint key at roughly 500 entries per btree leaf page, batch k traverses about 2k leaf pages of dead entries before finding a live one; summed, ~2.5×10⁹ buffer accesses instead of the ~100,000 an advancing cursor would need.

Estimate, not measurement: at sub-microsecond per cached page access that is on the order of tens of minutes of pure CPU, versus seconds. Two things soften it and neither bounds it — btree LP_DEAD hinting makes repeat traversals cheap per entry but does not remove them, and autovacuum may reclaim some of the tail concurrently, which makes the runtime depend on autovacuum timing rather than on the data. A resume that makes no visible progress for an unpredictable stretch, with every batch committing so there is no long transaction to alert on, is a bad failure mode for the path an operator reaches only after something already went wrong.

Fix, returning the batch's own high-water mark:

tag, err := tx.Exec(ctx, c.clearSQL, lower, limit)   // → RETURNING the key

Have clearBatch scan DELETE … RETURNING key, report the maximum, and set lower = maxKey + 1 for the next batch — with a maxKey == math.MaxInt64 check to end the loop, since that is the one case where exactly batch rows were removed and there is genuinely nothing above them (incrementing would wrap).

Narrowing the range like that is safe here for the reason in the third bullet above: every key above the watermark reads as KeyUncut for the whole clear, so nothing can insert into a range a previous batch already passed. That is worth saying in the comment, because it is the non-obvious precondition the current fixed-$1 form does not need and the advancing form does.

No test would catch this. TestCopierResumesFromTheZeroWatermark clears 201 rows with InitialRows: 40, so about six batches — enough to prove correctness, far too few for the shape to show. A t.Log of the batch count, or an assertion that clearing N rows takes ⌈N/batch⌉ batches rather than more, would pin it without needing a large fixture.

Notes

land's invariant marker and its error tag disagree.

// INV: CO-4
if !c.ledger.land(chunk, inserted) {
	return fmt.Errorf("%w (LK-3): landed chunk [%d, %d] was not in flight", ...)
}

claim right above it has // INV: CO-4 over a (CO-4) error. These markers look greppable against docs/invariants.md, so one of the two spellings here is wrong — LK-3 is the defensible one for "was not in flight", which makes the comment the thing to change.

batch := c.chunker.Rows() couples the delete batch to the copy chunk size, inertly. It reads once, before any Feedback, so it is always InitialRows — the D12 sizing loop never influences it. That is fine, but the expression reads as if the clear adapts, and a future reader moving the read inside the loop would make delete batches follow insert timings, which is not a relationship anyone wants. A named constant, or opts.Chunker.InitialRows with a one-line reason, says what is actually meant.

checkShadow does not verify the precondition ON CONFLICT (key) needs. The doc says "what the shape cannot promise, checkShadow verifies before a connection is opened", and it covers the column list carrying the primary key — but ON CONFLICT (key) DO NOTHING additionally requires a unique index on that column in the shadow. A builder that defers index creation until after the copy (a normal optimisation, and one this design invites since the copy is the expensive part) produces a shadow that passes every check here and fails on the first chunk with SQLSTATE 42P10, after the lock session, the guard, and the clear have all run. The OID query in confirmRelations is already round-tripping to the catalog and could confirm the index in the same statement.

A transient error anywhere costs the whole copy back to the watermark. There is no retry: one lock_timeout on LOCK TABLE in guard, on any of the workers, returns from work, cancels the rest, and ends the run. That is a defensible first cut — resume exists — but combined with the zero-watermark case above it means a single transient failure before the first chunk lands discards every committed chunk and pays the clear to delete them. Worth stating in SAFETY.md as a known operational cost rather than leaving it to be discovered, since the obvious mitigation (checkpoint the watermark, retry the transient classes inside work) is a separate change.

The first resume-clear batch takes SHARE MODE on the shadow so a chunk
transaction still open from the earlier run cannot commit a row behind
the delete; Position.Cut becomes a Watermark; new tests pin the per-chunk
OID guard, SET LOCAL ROLE, the ledger cut-through-MaxInt64 refusal, and
the fence.
@Kiran01bm

Copy link
Copy Markdown
Collaborator Author

🤖 Adversarial review response — created by Kiran's code review agent (Amp, Claude Opus 4.6) — pull/128, follow-up commit

Six of the seven findings are fixed in follow-up commit d67e2a2 (history preserved; nothing rewritten); the typed-refusal finding is deferred to the orchestrator PR and recorded in the docs so no importer meets it unannounced.

# Finding Status Explanation
C1-F1 The per-chunk guard in copyChunk is unpinned: every refusal test is refused by the resume clear's guard, so deleting the chunk guard leaves the package green fixed TestCopierRefusesAShadowReplacedAfterTheClaim (copy_chunk_guard_integration_test.go) drops and recreates the shadow at the copier's first Clock.Now() — after the clear has committed and the first chunk is claimed, before its transaction begins — and asserts ErrInvariantViolation naming (ST-6): shadow, zero rows in the impostor, and that one chunk in flight. It fails with c.guard removed from copyChunk.
C1-F2 SET LOCAL ROLE <owner> is unpinned because fixture tables are owned by the superuser the pool logs in as fixed TestCopierWritesAsTheTableOwner moves orders to a NOLOGIN role (newOwnedCopierFixture), builds the shadow with ADD COLUMN writer name NOT NULL DEFAULT current_user, and asserts all 100 rows carry writer = <owner>. It fails with the role statement replaced by SELECT 1.
C1-F3 The ledger's cut == MaxInt64 branch is unpinned; without it cut+1 wraps and a [MinInt64, 10] claim moves the frontier back fixed TestLedgerRefusesEveryClaimOnceTheKeySpaceIsCut (ledger_test.go) claims [MinInt64, MaxInt64], then asserts the wrapped claim is refused and the frontier stays at the largest key. It fails with the branch removed.
C1-F4 A straggling chunk transaction from the earlier run (client-side COMMIT failure, server still publishing) can become visible after the resume clear finishes, leaving stale rows above the new watermark that ON CONFLICT DO NOTHING keeps fixed The suggested closure: the first clear batch runs LOCK TABLE <shadow> IN SHARE MODE (fenceStragglers, resume.go) after its guard and before its delete, so it waits — bounded by lock_timeout, fail-closed — on every open shadow writer; later batches do not fence. TestCopierResumeWaitsForAStragglingChunkTransaction (resume_fence_integration_test.go) holds a shadow insert open, resumes from watermark 150, observes the ungranted ShareLock in pg_locks with the cut frontier unchanged and nothing in flight, commits the straggler, and asserts the resumed copy reads the committed value. The fence SQL is frozen (TestFenceStragglersSQLIsFrozen). The fence also waits on any concurrent shadow writer for that one batch; the three existing tests that opened a shadow write before Run now place it at the copier's first Clock.Now(), which models a concurrent applier write rather than a pre-resume straggler. invariants.md CO-4 Enforced today and the design package map record the fence.
C2-F1 Copier refusals carry their cause (LK-1, ST-6, CO-4, LK-3) only in message text, while schemachange promises a RefusalCause importers route on deferred Agreed on the destination; the mapping lands with the orchestrator PR, which is the first caller and owns the RefusalCause values the copier's failures map onto — adding a parallel cause type here would give importers two taxonomies to reconcile. Until then docs/refusal-classes.md (new paragraph after the shadow-refusal table) and the design package map state that copier failures are ErrInvariantViolation wrapping the cause, with lock loss the one structurally distinguishable case via errors.Is(err, lock.Err()), so no importer is invited to match on text.
C2-F2 Three spellings of "an optional key": Watermark, Cut int64 + CutValid bool, and Chunker.Cut() (int64, bool) fixed Position.Cut is a Watermark and Chunker.Cut() returns one; CutValid is gone and the ledger stores its frontier as a Watermark. The Watermark doc comment now describes a frontier in the key space rather than only the landed prefix. A checkpoint serializes both fields with one rule.
C2-F3 Options.Clock (and chunker.go) cite (D12) with no link fixed Both packages now say "chunk-time throttle" and link docs/copy-and-swap-design.md#d12--throttle-by-chunk-time-and-slot-lag (Options.Clock, TargetChunkTime, Feedback).

The mutation run (27 of 32 killed) and the invariant checks on CO-4, LK-1, LK-3, ST-6, LK-2 and TM-2 confirm behaviour that is unchanged — no action. The two surviving defensive mutants in finish (in-flight and watermark checks after every worker has returned nil) stay as fail-closed guards rather than being pinned or removed.

Source: block/pg-sprite#128, review comments 5869943436 and 5869943838 and review 5338688433 at head ca7af5e; fixes in follow-up commit d67e2a2.

@Kiran01bm
Kiran01bm enabled auto-merge (squash) September 29, 2026 04:14
@Kiran01bm
Kiran01bm merged commit b24a8b9 into main Sep 29, 2026
16 checks passed
@Kiran01bm
Kiran01bm deleted the kiran01bm/cs5-copy-step branch September 29, 2026 04:17
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants