Skip to content

applier: flush drained batches column-wise with the unique-move fallback - #149

Merged
Kiran01bm merged 2 commits into
mainfrom
kiran01bm/cs8-apply
Oct 8, 2026
Merged

Kiran01bm merged 2 commits into
mainfrom
kiran01bm/cs8-apply

Conversation

@Kiran01bm

@Kiran01bm Kiran01bm commented Oct 7, 2026 •

Copy link
Copy Markdown
Collaborator

Why

pkg/applier could merge and drain but not write. Buffer.Drain hands back a Batch of landed entries — delete markers, whole images, and marker-bearing images whose unchanged-TOAST columns still need completing, with CompleteFirst naming the moved ones — and nothing consumed it. This PR adds the write half: Flusher, which applies one batch in one guarded transaction under the table lock, column-wise where it can (CO-5, CO-8) and whole-row when a unique index or exclusion constraint refuses the column-wise form (CO-6, D13), plus Buffer.Requeue for the batch a constraint still refuses and Buffer.Hold / Buffer.Release for the moved images whose completion the flush read from a row that can be newer than the stream.

The scheduling that makes a chunk copy and a backlog flush mutually exclusive, and the stream owner that calls Drain → Flush → Hold → Release in a loop, are the next PRs. This one is the SQL and its transaction.

Before / after

seats(id PRIMARY KEY, slot UNIQUE, doc EXTERNAL) holding (1,'A',…),(2,'B',…) in the shadow; the source swaps the slots in one transaction and leaves doc untouched, so both decoded images carry doc as a marker:

Before                                          After
──────                                          ─────
Buffer.Drain ──► Batch{1→'B' doc:u, 2→'A' doc:u}   Buffer.Drain ──► Batch ──► Flusher.Flush
                 (nothing consumed it)                                  │
                                                                        ├─ BEGIN; guard (owner, search_path,
                                                                        │  ACCESS SHARE ×2, lock confirmed, OIDs)
                                                                        ├─ SAVEPOINT
                                                                        │   UPDATE shadow SET slot='B' WHERE id=1   ◄── only present
                                                                        │   ✗ 23505: slot 'B' is row 2's                columns (CO-8)
                                                                        ├─ ROLLBACK TO SAVEPOINT
                                                                        │   doc for 1, 2 ◄── shadow rows FOR UPDATE   (complete)
                                                                        │   DELETE FROM shadow WHERE id = ANY({1,2})
                                                                        │   INSERT (1,'B',doc1); INSERT (2,'A',doc2)  (CO-6)
                                                                        └─ COMMIT ──► Result{Images 2, Completed 2, Fallback}

A moved image, UPDATE seats SET id = 10 WHERE id = 2 captured with doc as a marker. The old key's shadow row was copied live and the source row is live, so either can already hold a change the stream has not delivered — a different row's doc, if a later source transaction put another row at the key. The flush reads the completion but does not write it; the buffer holds the image until the stream has passed the read:

Batch{2: DeleteMarker, 10: Image{OldKey 2, slot 'B', doc:u}}   CompleteFirst → [10]
  │
  ├─ SELECT doc, pg_current_wal_insert_lsn() FROM shadow WHERE id = 2 FOR UPDATE
  │     (OldKeyReused, or no shadow row under 2 → … FROM source WHERE id = 10;
  │      neither row → DELETE FROM shadow WHERE id = 10, key in Result.Skipped)
  ├─ DELETE FROM shadow WHERE id = 2                              (deletes before images)
  └─ COMMIT ──► Result{Deletes 1, Held [{10, doc2, ReadLSN L}]}

Buffer.Hold(Result.Held)          image 10 back in the buffer, completion pending
  … stream delivers events …      an event replacing key 10, or reusing key 2, drops the completion;
                                  an UPDATE of key 10 keeps it pending behind its own values
Buffer.Release(passed ≥ L)        fills the markers still left ──► next Drain flushes image 10 whole

What

  • pkg/applier/flusher.go — Options{LockTimeout, StatementTimeout} (defaults; sub-millisecond refused as ErrInvalidOptions), Result{Images, Deletes, Completed, Held []HeldImage, Skipped []int64, Fallback}, ErrBatchDeferred. NewFlusher(target, shadow, lock, opts) runs the copier's proof and shadow checks (ST-6). Flush(ctx, pool, batch) is a no-op for an empty batch; otherwise binds the lock session's context, begins one read-write transaction with the guard, reads a completion for every CompleteFirst image into Result.Held (or writes a delete at its key into Result.Skipped), applies the rest column-wise in a savepoint — every delete marker before any image — falls back whole-row on 23505 or 23P01, commits, and reports a lost lock over any other error.
  • pkg/applier/hold.go — HeldImage{Entry, Completed, FromSource, ReadLSN}; Buffer.Hold(held) puts each image back with its completion pending (a key the buffer holds is ErrInvariantViolation, CO-5, buffer unchanged); Buffer.Release(passed) int fills the markers still left on every image whose read is at or below passed; Buffer.Held(). Any Add that replaces a held key's entry or reuses a held image's old key drops the pending completion; Drain keeps a held image buffered and counts it in Batch.Held.
  • pkg/applier/complete.go — completeFrom reads the named row's columns as text together with pg_current_wal_insert_lsn() in the same statement; the shadow read is FOR UPDATE, the source read plain.
  • pkg/applier/drain.go — CompleteFirst returns copies of the batch's entries; Batch.Held.
  • pkg/applier/flush_sql.go — the statements as free builders over sanitized identifiers: upsertSQL (present columns; a key-only image becomes DO UPDATE SET pk = EXCLUDED.pk so the row exists either way), updateSQL (present columns, WHERE pk = $1), insertSQL (whole row, no conflict clause), deleteSQL, deleteAllSQL (pk = ANY($1::bigint[])), rowSQL (each column ::text, plus the WAL position); split puts an image's columns in the shadow's column order and treats an unnamed or value-less column as a marker.
  • pkg/applier/guard.go — the package's copy of the copier's guard (owner role, catalog-only search_path, bounded timeouts, extra_float_digits = 3, ACCESS SHARE on both relations, lock confirmation, relation-OID confirmation, lock-lost reporting), now with its own tests: a shadow that is not the target's is refused at NewFlusher, a decoy = operator on the session's search_path is ignored, a connection configured at extra_float_digits = 0 still completes a float8[] marker exactly, and an empty batch takes no connection from the pool.
  • pkg/applier/requeue.go — Buffer.Requeue(batch) puts a deferred batch back as the flush left it; a key the buffer already holds is ErrInvariantViolation (CO-5) and the buffer is unchanged.
  • Docs — SAFETY.md applier row; design D13 (the hold gate replaces the "a later event overlays it" reasoning; deletes first; 23P01; a skipped image is a delete) and package map; docs/invariants.md CO-4 (what is enforced now vs. the mutual-exclusion scheduling still planned), CO-6 (Enforced), CO-8 (Enforced); pkg/applier/doc.go.

Decisions to veto

  • A moved image's completion is held, not written, until the stream passes the read. The read's pg_current_wal_insert_lsn() bounds every commit the read could have seen; the image waits in the buffer until the stream owner has delivered through that position and calls Release. An event for the key meanwhile wins: an UPDATE overlays its own values and the completion fills only the markers left; anything that replaces the key's entry or reuses the old key drops the completion. The alternative — write the read value and let a later event overlay it — is unsound, since a later move away from the key reads the shadow row instead of overlaying it.
  • One gate, the flush's read position, for shadow and source completions alike. The shadow row under the old key was copied before the flush read it, so a stream that has passed the flush's read has passed the copy's too. This is conservative for the shadow side — the chunk's own read position would release sooner — and keeps copier.Position as it is. A per-chunk read position can be added to Position later without changing the applier's contract.
  • A moved image neither row can complete is written as a delete at Key. The source holds no row under that key, so whatever the shadow holds there is stale, and the key-moving UPDATE that landed on the key may have replaced the delete marker that was the only record of a row to remove. Deleting is right whatever the shadow holds; a later move away from the key finds no row and falls through to the source. The batch's own entry is left as it was so Requeue can put the image back.
  • Delete markers are written before images. No key holds both, so the order is always sound, and a row that moved to a lower key while keeping a unique value no longer costs the whole batch the fallback.
  • 23P01 is a collision like 23505. Preflight admits exclusion constraints, and an exclusion over a moved value fails the same way a unique index does; the fallback converges it. A DEFERRABLE INITIALLY DEFERRED unique constraint reports at COMMIT and would surface as a commit error rather than ErrBatchDeferred; tracked as an internal follow-up.
  • Result.Held and Result.Skipped carry the entries, not counts. The stream owner needs the held images to call Hold; an operator triaging a later verifier mismatch needs the skipped keys.
  • The guard is copied into pkg/applier, not shared. Same convention as pkg/checksum: each writing package carries its own guard so a change to one package's session setup is reviewed in that package, and this copy now has the tests that pin it. A shared internal package beside chunksql would save ~100 lines and couple three packages' transaction shape; tracked as an internal follow-up.
  • Completion reads as text and writes the text back. The completed value is bound as a text parameter and the server casts it to the column's type on write, the same round trip the decoder's text-mode values take. Reading typed values would need per-type scanning for no gain.
  • extra_float_digits = 3 is pinned on the flush session too. The verifier hashes the text rendering; a float the flush rendered at the session's default could land one digit off what the decoder saw and read as a mismatch. The completion read runs in the same transaction, so a held float8 value carries every digit.
  • A marker-bearing image of an unmoved row is an UPDATE, not an upsert. The marker stands for a value only a present shadow row holds; an upsert would insert the row with the omitted column invented. Zero rows updated is a CO-8 refusal.
  • The fallback completes every remaining marker from the shadow row under its own key, and refuses if absent. An unmoved image's marker refers to that row and nothing else; the fallback runs before its own deletes, so the row is still there.
  • A second collision is ErrBatchDeferred, not a retry and not a refusal. With every key in the batch deleted first, a remaining violation means the colliding row is outside the batch — a key an in-flight chunk is about to land or a deferred entry still holds. The copier moving resolves it; the caller requeues and drains later. A retry inside the flush would spin on the same state.
  • Requeue and Hold refuse a held key and run before any Add after the Drain. The stream owner's loop is drain → flush → (requeue | hold + confirm) → add; a key both the batch and the buffer hold means that order broke, and merging the two would pick an arbitrary winner.
  • No applied-LSN API. The position a caller may confirm is Batch.OldestFirstLSN of a batch whose flush committed and Buffer.OldestPending of what has not, which covers held images; the stream owner that confirms is the next PR's.
  • Nothing bounds a batch yet. A catch-up can drain a large backlog into one Flush, which is one transaction with one round trip per entry. The bound belongs with the stream owner (Drain(pos, limit) keeping pairs together) and set-based statements are the natural next step for the round trips; tracked as an internal follow-up.
  • Load-generator convergence oracle is not here. CO-6's fixed vectors run as integration tests; the load-generator form waits on the convergence harness with the stream owner. Tracked as an internal follow-up.
  • The copier's chunk insert can also meet 23505 when a flush landed a row carrying a unique value a chunk is about to copy; the copier's ON CONFLICT (pk) DO NOTHING does not cover a secondary unique index. The mutual-exclusion scheduling is where that is decided. Tracked as an internal follow-up.
  • A marker-bearing UPDATE for an unmoved key whose shadow row the copier read after a later delete or move of that key reaches the flush with no shadow row and fails closed as CO-8. The scheduling PR decides whether the stream owner holds such an image the same way or the copier's read position rules it out; tracked as an internal follow-up.

Verification

  • make lint 0 issues; go vet ./pkg/applier/; go build ./...; SKIP_INTEGRATION=1 go test -race ./pkg/applier/ green.
  • go test -race -count=1 ./pkg/applier/ green on PostgreSQL 14, 16, and 18 (~50s each); scripts/test-flaky.sh 8/8 on the hold, skip, and empty-batch tests.
  • Integration, column-wise and fallback: a batch of one insert, one partial update, and one delete lands exactly those changes; an omitted doc (SET STORAGE EXTERNAL, ToastBytes() > 0 proven) is left in place byte-for-byte; a marker for a row the shadow lacks is ErrInvariantViolation (CO-8); the seats exchange {1→'B', 2→'A'} converges in one flush with Fallback set, and so does the same exchange under EXCLUDE USING btree (slot WITH =); a move to a lower key keeping its slot stays column-wise; the fallback completes markers from the shadow; a batch colliding with a row outside it returns ErrBatchDeferred, Requeue takes it back, and the next flush converges once the colliding row is in the batch; an image the fallback cannot complete is refused.
  • Integration, hold gate: a move whose old key landed is held with the old key's doc and converges after Release; a move whose old key never landed is held with the source's; a reused old key completes from the source; a move whose row is gone writes a delete and names the key in Skipped; a row that moves onto a deleted key and on again keeps its own doc (the skipped delete removes the predecessor); a row that moves 3→10→11 while another row is inserted at 10 before the first flush reads the source — the first hold carries the other row's doc, the stream's move away from 10 drops it, and key 11 ends with row 3's own; a held shadow completion is dropped when the stream delivers a reuse of the old key; an UPDATE of a held key keeps the completion pending and Release fills only the markers it left.
  • Integration, guard: NewFlusher refuses a shadow that is not the target's (nil lock, other table's lock); a decoy = operator on the session's search_path is ignored; a connection at extra_float_digits = 0 still completes a float8[] marker holding 0.1 + 0.2 with every digit; an empty batch takes no connection; the writes run as the table owner; a replaced shadow is refused (ST-6); a gone lock is reported as lost (LK-1). Removing the search_path pin or the extra_float_digits pin from guard.go fails the corresponding test.
  • Unit: Hold keeps the image until the stream passes the read, releases only the completions at or below passed, refuses a held key, fills only markers, drops on delete/move/reuse, copies the image and its completion, no-op on empty; CompleteFirst returns copies; the SQL builders against literal expected statements; split follows the shadow's column order; Requeue round-trips a batch, refuses a held key, no-op on empty; option defaults and refusals; 23505 and 23P01 matched by SQLSTATE only.

Add Flusher, the write half of pkg/applier. Flush applies one drained Batch
in one bounded read-write transaction under the table lock session, guarded
as the copier's chunk inserts are (owner role, catalog-only search_path,
ACCESS SHARE on both relations, lock confirmation, relation-OID check).

The flush first completes every image Batch.CompleteFirst names: from the
old key's shadow row, read FOR UPDATE so the batch's own delete of that key
cannot race it; from the source row under Key when the old key was reused or
its shadow row is absent; an image the source has since removed is skipped
(D13). It then writes column-wise inside a savepoint: a delete marker as a
DELETE, a whole image as an upsert, a marker-bearing image of an unmoved row
as an UPDATE of only its present columns, where zero rows updated is a CO-8
refusal rather than an insert that invents the omitted value.

On SQLSTATE 23505 the flush rolls back to the savepoint and reapplies the
batch whole-row (CO-6): every remaining marker completed from the shadow row
under its own key, every key in the batch deleted in one statement, every
surviving image inserted plain. A second 23505 returns ErrBatchDeferred with
the shadow untouched; Buffer.Requeue puts the batch back for a later drain
and refuses, fail closed, if any of its keys is held.

The flush pins extra_float_digits so a float lands as the value the
verifier will hash. Result counts images, deletes, completions by source,
and skips, and reports whether the fallback ran.

Integration tests cover the column-wise path, the omitted-column UPDATE, the
seats unique exchange converging in one flush, completion from the shadow
and from the source (including a reused old key), a moved row that is gone,
a batch a unique index still refuses, the owner role, a replaced shadow, and
a lost lock. SAFETY.md, the design's package map, and CO-4/CO-6/CO-8 now
describe the flush.
@Kiran01bm
Kiran01bm marked this pull request as ready for review October 7, 2026 05:23
@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

aparajon commented Oct 7, 2026

Copy link
Copy Markdown
Collaborator

🤖 1/2: adversarial correctness review of e93aa1d. I read the whole diff against the buffer and drain from #148, the copier's guard it copies, Position.Classify, and CO-4, CO-5, CO-6 and CO-8 with D13 as amended. I ran the flusher across several drains with the real Buffer in front of it, since what a flush leaves in the shadow is what the next flush's completion reads. That covered chains of moves that cross a flush boundary, old keys reused before or after the copier reads them, moves toward lower keys on a table with a unique index, and constraint shapes that preflight admits. Everything ran on PostgreSQL 16.

2 blocking, 3 non-blocking.

The column-wise path, the savepoint fallback and Requeue do what the description says, and the CO-6 vectors converge. The guard matches the copier's step for step, apart from one equivalent reordering. 21 of 30 mutants are killed. Of the nine survivors, two are equivalent today and seven are guard branches with no flusher test (N3). Both blocking findings come from the same place: completing a moved image from the source. The flush reads whatever row sits under Key when it runs, and the next flush trusts that shadow row as the moved row's own. D13's reason for the source read is that "any later change arrives as a later event that overlays it" (design L377-L381). A later move away from Key does not overlay that value. It reads it. Nothing calls Flush yet, but CO-6 and CO-8 name Flusher in their Enforced lines from this merge on, and the two tests below fail against it. On invariants: CO-6 and CO-8 are extended, with their first applier-side enforcement, though with B1 and B2 open CO-6's convergence claim is not yet true of this code. CO-4 and CO-5 are upheld (Requeue keeps the buffer disjoint). ST-6, LK-1 and LK-2 are upheld by the copied guard. CO-9 is upheld but untested here (N3).

Blocking

B1. A skipped image leaves the deleted predecessor's row under Key, and the next move away from Key completes from it. flusher.go:236-256

When a key-moving UPDATE lands on a key whose row was deleted earlier in the same window, addKeyMove replaces that delete marker with the image. So the image is now the only record that the shadow row under Key has to go. If the flush then finds neither the old key's shadow row nor a source row under Key, it drops the image, and the delete goes with it. The deleted row stays in the shadow under Key. The later event the skip relies on is a move away from Key. That event completes from the shadow row under Key, which is the deleted row.

seats holds rows 1, 2 and 3. The source deletes row 2, moves row 3 to key 2 (the copier reads key 3's chunk after this, so there is no shadow row 3), and then moves it on to key 10. doc stays unchanged throughout.

window 1  buffer {2: image(OldKey 3, doc: u), 3: delete}
          complete key 2: shadow[3] absent, source[2] absent  -> skip
          shadow[2] is still the deleted row 2
window 2  buffer {2: delete, 10: image(OldKey 2, doc: u)}
          complete key 10 from shadow[2] FOR UPDATE           -> row 2's doc
          shadow[10].doc = row 2's document; source[10].doc = row 3's

The fix is local. Write a skipped image as a delete at Key, not as nothing. The source holds no row there, and deleting is right whatever the shadow holds under Key: a later move away finds no row and falls through to the source. Build the delete in completeMoved's filtered slice. Do not rewrite the batch's entry, because Requeue must still put the image back. The buffer would refuse a later legal move away from Key as "a move from a buffered delete".

Test that fails on e93aa1d and passes with the fix
func everyKeyLanded() copier.Position {
	return copier.Position{Watermark: copier.NewWatermark(math.MaxInt64), Cut: copier.NewWatermark(math.MaxInt64)}
}

func moveEvent(lsn decode.LSN, from, to int64, cols ...decode.Column) decode.ChangeEvent {
	return decode.ChangeEvent{Kind: decode.Update, LSN: lsn, Key: to, OldKey: &from, Columns: cols}
}

// Row 3 moves onto key 2 after the row there was deleted, then on to key
// 10; the copier read key 3's chunk after the first move. The first flush
// has nothing to complete key 2 from, and the second still completes the
// move to key 10 with row 3's own document, not the deleted row's.
func TestFlushSkippedMoveLeavesNoPredecessorRow(t *testing.T) {
	f := newFlushFixture(t)
	f.createSeats(t, 3)
	p := f.prepare(t, "seats")
	doc3, _ := f.text(t, "seats", "doc", 3)
	f.exec(t, `DELETE FROM %s.`+p.shadow.ShadowTable()+` WHERE id = 3`)
	f.exec(t, `DELETE FROM %s.seats WHERE id = 2`)
	f.exec(t, `UPDATE %s.seats SET id = 2 WHERE id = 3`)
	f.exec(t, `UPDATE %s.seats SET id = 10 WHERE id = 2`)

	buffer := applier.NewBuffer()
	require.NoError(t, buffer.Add(decode.ChangeEvent{Kind: decode.Delete, LSN: 1, Key: 2}))
	require.NoError(t, buffer.Add(moveEvent(2, 3, 2, col("slot", "C"), marker("doc"), col("note", "seat 3"))))
	result, err := p.flush.Flush(t.Context(), f.pool, buffer.Drain(everyKeyLanded()))
	require.NoError(t, err)
	assert.Equal(t, 1, result.Skipped)

	require.NoError(t, buffer.Add(moveEvent(3, 2, 10, col("slot", "C"), marker("doc"), col("note", "seat 3"))))
	_, err = p.flush.Flush(t.Context(), f.pool, buffer.Drain(everyKeyLanded()))
	require.NoError(t, err)

	got, _ := f.text(t, p.shadow.ShadowTable(), "doc", 10)
	assert.Equal(t, *doc3, *got, "the moved row keeps its own document")
	f.assertConverged(t, p.shadow)
}

--- FAIL: TestFlushSkippedMoveLeavesNoPredecessorRow (0.09s) on e93aa1d: key 10 holds row 2's document, and AssertConverged reports key 10 in both directions. It passes once the skipped key is written as a delete in completeMoved's filtered entries.

B2. A source completion can take another row's value, and a later flush spreads it. flusher.go:236-243, complete.go:22-30

The source row under Key is read at flush time. The stream can be behind that. If, before the read, the moved row left Key and another row arrived there, the read returns the other row's document. The flush writes it into the moved row's image. The moved row's later move away from Key then completes from that shadow row and carries the wrong document on to the new key. The other row's own INSERT, when it arrives, fixes only Key. Here row 3 moves 3→10, then 10→11, and then a new row is inserted at 10, all before window 1's flush reads the source:

window 1  {3: delete, 10: image(OldKey 3, doc: u)}  -> source[10] is the NEW row: doc = its document
window 2  {10: delete, 11: image(OldKey 10, doc: u)} -> shadow[10] -> shadow[11].doc = the new row's
window 3  {10: insert}                               -> key 10 fixed; key 11 stays wrong

The reuse flag has the same blind spot on the shadow side. D13 relies on "the buffer sees the reuse" (L373). But the copier reads chunks live. It can copy the old key's chunk after another row has moved in, while the stream has not yet delivered that row's event. OldKeyReused is then false, and the completion reads the other row from the shadow. B2 is not run here, because the test layer has no stream to lag. The cause is the same, though: a completion read can be newer than the buffer, and newer can mean a different row.

This is a D13 decision, so I am not offering a patch. The gate needs this property: a completion read counts only once the buffer holds every event the read could have seen. One shape is to read pg_current_wal_insert_lsn() in the same statement as the source read. Every commit visible to that statement's snapshot has its commit record below that LSN. The image then stays buffered until the stream has added events through that LSN, and is written only if no event touched the key meanwhile. The old key's chunk would need a read position in Position for the same check. Whatever gate you choose, the "overlays it" sentence in D13 and the matching item in the description's veto list should change in this PR. As written, the registry claims convergence that this test disproves.

Test that fails on e93aa1d (fix is a D13 decision)
// Row 3 moves to key 10 and on to 11, and a new row is inserted at key 10,
// before the first flush reads the source. Every key ends holding its own
// row's document.
func TestFlushSourceCompletionNeverTakesAnotherRowsValue(t *testing.T) {
	f := newFlushFixture(t)
	f.createSeats(t, 3)
	p := f.prepare(t, "seats")
	doc3, _ := f.text(t, "seats", "doc", 3)
	f.exec(t, `DELETE FROM %s.`+p.shadow.ShadowTable()+` WHERE id = 3`)
	f.exec(t, `UPDATE %s.seats SET id = 10 WHERE id = 3`)
	f.exec(t, `UPDATE %s.seats SET id = 11 WHERE id = 10`)
	f.exec(t, `INSERT INTO %s.seats (id, slot, doc, note) VALUES (10, 'Z', 'the other row''s document', 'seat 10')`)

	buffer := applier.NewBuffer()
	require.NoError(t, buffer.Add(moveEvent(1, 3, 10, col("slot", "C"), marker("doc"), col("note", "seat 3"))))
	_, err := p.flush.Flush(t.Context(), f.pool, buffer.Drain(everyKeyLanded()))
	require.NoError(t, err)
	require.NoError(t, buffer.Add(moveEvent(2, 10, 11, col("slot", "C"), marker("doc"), col("note", "seat 3"))))
	_, err = p.flush.Flush(t.Context(), f.pool, buffer.Drain(everyKeyLanded()))
	require.NoError(t, err)
	require.NoError(t, buffer.Add(decode.ChangeEvent{Kind: decode.Insert, LSN: 3, Key: 10,
		Columns: []decode.Column{col("slot", "Z"), col("doc", "the other row's document"), col("note", "seat 10")}}))
	_, err = p.flush.Flush(t.Context(), f.pool, buffer.Drain(everyKeyLanded()))
	require.NoError(t, err)

	got, _ := f.text(t, p.shadow.ShadowTable(), "doc", 11)
	assert.Equal(t, *doc3, *got, "the moved row keeps its own document")
	f.assertConverged(t, p.shadow)
}

--- FAIL: TestFlushSourceCompletionNeverTakesAnotherRowsValue (0.09s) on e93aa1d: key 11 holds "the other row's document". The B1 fix does not change it.

Non-blocking

N1. Every move to a lower key on a table with a unique index takes the fallback. flusher.go:267-285

The column-wise pass writes in key order. A row that moves 3→0 keeps its slot, so the upsert at 0 collides with the row still under 3, which the delete marker at 3 has not removed yet. That costs the whole batch the fallback: a completion read FOR UPDATE of every marker-bearing image, the delete-all, and a whole-row rewrite of every image, out-of-line values included. Writing the batch's delete markers before its images avoids it. That order is always safe, because no key holds both an image and a delete marker. With deletes first, TestFlushFallbackCompletesMarkersFromTheShadow also stops needing the fallback, since its collision was a key-order artifact rather than a cycle. Its vector would need a real exchange to keep covering the fallback.

Test that fails on e93aa1d and passes with deletes first
func TestFlushMoveToALowerKeyStaysColumnWise(t *testing.T) {
	f := newFlushFixture(t)
	f.createSeats(t, 3)
	p := f.prepare(t, "seats")
	doc3, _ := f.text(t, "seats", "doc", 3)
	f.exec(t, `UPDATE %s.seats SET id = 0 WHERE id = 3`)

	result, err := p.flush.Flush(t.Context(), f.pool, batch(
		moved(3, 0, col("slot", "C"), col("doc", *doc3)),
		deleted(3),
	))

	require.NoError(t, err)
	assert.False(t, result.Fallback, "only the row's own old key held its slot")
	f.assertConverged(t, p.shadow)
}

--- FAIL: TestFlushMoveToALowerKeyStaysColumnWise (0.52s) on e93aa1d: Fallback is true. It passes when applyColumnWise writes delete markers first.

N2. An exclusion constraint has the CO-6 problem but gets neither the fallback nor ErrBatchDeferred. flusher.go:362-365

Preflight admits a table with EXCLUDE USING btree (slot WITH =). The seats exchange on it fails 23P01 and comes back as a plain upsert error. That fails closed, but the change can never converge. Matching 23P01 alongside 23505 makes the same exchange converge through the fallback (measured). The alternative is to refuse exclusion constraints in preflight. A DEFERRABLE INITIALLY DEFERRED unique constraint is the same gap in another form. PostgreSQL reports it at COMMIT, so a collision with a row outside the batch surfaces as a commit error rather than ErrBatchDeferred (not run here).

Test that fails on e93aa1d and passes with 23P01 in the fallback
func TestFlushConvergesAnExclusionExchange(t *testing.T) {
	f := newFlushFixture(t)
	f.exec(t, `CREATE TABLE %s.seats (id bigint PRIMARY KEY, slot text, doc text, note text, EXCLUDE USING btree (slot WITH =))`)
	f.exec(t, `INSERT INTO %s.seats VALUES (1, 'A', 'd1', 'n1'), (2, 'B', 'd2', 'n2')`)
	p := f.prepare(t, "seats")
	f.exec(t, `UPDATE %s.seats SET slot = NULL WHERE id = 1`)
	f.exec(t, `UPDATE %s.seats SET slot = 'A' WHERE id = 2`)
	f.exec(t, `UPDATE %s.seats SET slot = 'B' WHERE id = 1`)

	result, err := p.flush.Flush(t.Context(), f.pool, batch(
		image(1, col("slot", "B"), col("doc", "d1")),
		image(2, col("slot", "A"), col("doc", "d2")),
	))
	require.NoError(t, err)
	assert.True(t, result.Fallback)
	f.assertConverged(t, p.shadow)
}

--- FAIL: TestFlushConvergesAnExclusionExchange (0.07s) on e93aa1d: conflicting key value violates exclusion constraint … (SQLSTATE 23P01).

N3. Several branches of the copied guard have no flusher test, so nothing pins them in this package's copy. guard.go:70-101, flusher.go:120-174

Each of these mutants survives the whole pkg/applier suite:

  • dropping the catalog-only search_path (CO-9; the copier has TestCopierIgnoresTheSessionSearchPath for its copy);
  • dropping extra_float_digits = 3 (a float8[] or geometric marker under a role set to 0 would pin it);
  • dropping the ACCESS SHARE hold;
  • dropping NewFlusher's checkShadow and requireTableLock;
  • dropping the lock-loss check after commit;
  • dropping the empty-batch short-circuit (TestFlushOfAnEmptyBatchIsANoOp never observes whether a transaction opened);
  • dropping FOR UPDATE on the completion read.

The veto list keeps the guard per package so that each copy is reviewed where it lives. That works only if each copy carries its own tests. Porting the copier's search_path and replaced-shadow-after-claim tests would close most of the list.

Verified

  • go test -count=1 ./pkg/applier/ and SKIP_INTEGRATION=1 go test -race ./pkg/applier/ are green at e93aa1d. CI is green on PostgreSQL 14 to 18.
  • Checked by reading, and not findings:
    • the savepoint rollback leaves completions in place and the shadow pristine for the fallback's completion reads;
    • a rollback-to-savepoint failure joined to a 23505 drives the fallback into an aborted transaction, which then fails closed;
    • Requeue refuses before writing anything, and a batch is disjoint from what the drain deferred;
    • a move away and back (OldKey == Key) completes from its own row;
    • NULL completes as NULL;
    • the guard's step order differs from the copier's only by confirming the lock after the hold, which is equivalent.

Mutants ran against ./pkg/applier/ (unit and integration). Every edit was restored with git checkout.

Mutant Caught by
reused image completes from the shadow TestFlushCompletesAReusedOldKeyFromTheSource
never complete from the shadow TestFlushFallbackCompletesMarkersFromTheShadow
no skip TestFlushSkipsAMovedImageWhoseRowIsGone
no rollback to the savepoint 4 fallback tests
uniqueViolation never / always 5 tests / TestUniqueViolationMatchesBySQLSTATE only
marker image upserted TestFlushRefusesAMarkerForARowTheShadowLacks
zero-row UPDATE accepted TestFlushRefusesAMarkerForARowTheShadowLacks
fallback skips completion TestFlushFallbackCompletesMarkersFromTheShadow
fallback deletes only image keys TestFlushFallbackCompletesMarkersFromTheShadow
second 23505 not ErrBatchDeferred TestFlushDefersABatchAUniqueIndexStillRefuses
fallback keeps the column-wise Images count 3 fallback tests
copy columns include the key 8 tests
source completion reads the shadow TestFlushCompletesAMoveWhoseOldKeyNeverLandedFromTheSource
split treats a marker as present 3 tests
UPDATE binds off by one 5 tests
no SET LOCAL ROLE TestFlushWritesAsTheTableOwner
no lock confirmation TestFlushRefusesWhenTheTableLockIsGone
no relation-OID check TestFlushRefusesAReplacedShadow
Requeue without the held check / writes before checking TestRequeueRefusesAHeldKey / 4 tests
no search_path pin, no extra_float_digits, no ACCESS SHARE hold survive (N3)
no checkShadow, no requireTableLock, no post-commit lock check, no empty-batch return survive (N3)
no FOR UPDATE on the completion read survives (N3); equivalent until a concurrent writer can touch the shadow
NULL completes as '' survives; equivalent today, since an unchanged-TOAST marker is never NULL and a later event overlays a source value that went NULL

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

@aparajon

aparajon commented Oct 7, 2026

Copy link
Copy Markdown
Collaborator

🤖 2/2: OSS adoption and integration ease, at e93aa1d. These are lenses, not correctness findings. 0 blocking, 3 non-blocking.

For adoption, this is the first piece of pg-sprite that writes change data, and it reads well. Each SQL shape is a small builder pinned against its literal statement. The savepoint fallback is ten lines that a reader can check against CO-6. The integration tests state each source change as SQL next to the batch the buffer would drain, so a newcomer can see what every vector means. The "Decisions to veto" list made this review faster: B1 and B2 in 1/2 are both answers to one of its items.

For SchemaBot and other importers, nothing calls Flush yet. These notes are about what the stream owner, and a host rendering its progress, will need from this surface. Each is cheaper to shape now than after the stream owner depends on it.

1. Nothing bounds a batch, and a batch is one transaction with one round trip per entry. flusher.go:160-214

Options bounds each statement and each lock wait. It does not bound the batch. A catch-up after a long copy can drain hundreds of thousands of keys into one Flush, which means one statement per entry, a completion read per marker image, and on fallback a whole-row rewrite of every image, all in one transaction. That transaction holds ACCESS SHARE on both relations and row locks on the shadow, and it holds back the vacuum horizon on the source's primary for as long as it runs. The stream owner and the ST-3 lag ceiling both need to choose a flush size. The engine should own that choice, as Drain(pos, limit) or Options.MaxEntries, rather than leave each host to slice Batch.Entries. Slicing by hand would break the pair-travels-together rule CO-4 relies on. Set-based statements (unnest of key and value arrays) are the natural next step for the round trips, and the builders are already shaped for that.

2. The completions a host should trust least are reported as counts, without keys. flusher.go:79-90

CompletedFromSource and Skipped are the two paths where the flush inferred a value instead of reading the row the marker refers to. Today the only signal is a count. When a later verifier pass reports a mismatched chunk, an operator wants to know which keys were completed from the source or skipped, and when. A host like SchemaBot renders what the engine reports rather than reconstructing it. Returning the keys in Result (bounded, with a total), or emitting one structured log line per such completion, would make those paths triageable. Whatever gate B2 settles on would also hold these same entries back, so it is worth naming them now.

3. The guard now has three copies, and this one is the least tested. guard.go:46-57

The veto list keeps a guard per package, so that a change to one package's session setup is reviewed in that package. For someone maintaining a fork, three hand-kept copies of a security-relevant sequence (role, catalog-only search_path, lock confirmation, OID check) means three places to find and patch, and N3 in 1/2 shows the cost: this copy's search_path pin and ACCESS SHARE hold can be deleted with every test still green. The repo already has the shape for sharing that stays out of reach of callers: the unexported chunksql package holds the copier's statement for both of its users. A guard package beside it, tested once against the server, would let each package keep only what is its own, such as the flush's extra_float_digits pin.

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.

🤖 Stamping with comments: 2 blocking, 6 non-blocking (see the two review comments above). Please address B1 (write a skipped moved image as a delete at its key) and B2 (gate source completion so it cannot take another row's value, and correct D13's "overlays it" claim) before merging.

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.

Approving. The write half is carefully ordered, and the two orderings that matter most are both right: completing a moved image from the old key's shadow row before anything in the batch deletes that key, and deleting every key in one statement before inserting any of them so the batch cannot collide with itself on a unique index. Refusing a zero-row UPDATE as a CO-8 violation rather than converting it to an insert is the right instinct — inventing a value for an omitted TOAST column is exactly the failure that would be invisible until a verifier run.

Verified rather than assumed:

  • extra_float_digits = 3 is the one rendering setting that genuinely needs pinning, and the others genuinely do not. rowSQL reads columns ::text and writes the text straight back, so output lossiness matters — and float output is the only case where the output function can lose information the input function cannot recover. DateStyle, IntervalStyle and bytea_output change the form of the text but not its information content, and output and input happen in the same session with the same settings, so they round-trip exactly. Pinning one GUC and not the others looks arbitrary at first read; it is correct.
  • The $1::bigint[] cast in deleteAllSQL matches the copier's established convention, not a new assumption. pkg/copier/chunker.go declares its parameters bigint "whatever the key's integer type" with the operator-family reasoning written out, so the applier is consistent with it.
  • The fallback completes markers before it deletes. Reading FOR UPDATE under the table lock makes the read-then-delete sequence safe without depending on the isolation level, which the transaction does not set.
  • Requeue refusing when the buffer already holds one of the batch's keys is the right invariant: it makes "nothing may Add between Drain and the flush's outcome" checkable rather than conventional.

A deferred batch has no bound and no way to tell "waiting on the copier" from "stuck forever"

ErrBatchDeferred → Requeue → drain again. The reasoning given is that the caller "drains again once the copier has landed the chunk the rows wait on", which is a sound account of the expected cause and the only cause the code can recover from. Nothing distinguishes that case from one that will never resolve.

Concretely: the shadow holds a row under key 3 with slot = 'C' from an earlier chunk copy. A batch wants key 1 to take slot = 'C'. The column-wise UPDATE hits 23505; the whole-row fallback deletes only the batch's own keys, so row 3 survives and the insert hits 23505 again; the batch defers. If the event that would have freed 'C' was skipped — and the flush has a Skipped counter, so skipping is a real outcome — or if the shadow has drifted for any other reason, no later chunk and no later batch will clear it. The stream stops advancing and nothing raises an error. Buffer.OldestPending quietly stops moving.

The scheduling loop is explicitly a later PR, and that is the right place for the detection, so this is not a change I would hold this PR for. But it should be named now while the reasoning is fresh: a deferral count on the batch, or a check that OldestPending advances within some bound, turns a silent stall into a reportable condition. Without one, the difference between "the copier is behind" and "the shadow has drifted" is invisible from the outside, and those need opposite responses.

Non-blocking

split treats an unnamed column and a value-less column as the same thing, and they are not. A value-less column is a legitimate unchanged-TOAST marker: the correct handling is to leave that column alone. An unnamed column is a decoder defect. Folding the second into the first means a decoding bug does not surface as an error — it surfaces as a column silently retaining its old value in the shadow, which is precisely the drift that the deferral above cannot recover from, discovered later by a verifier with no way to trace it back. Given how much care the rest of this package takes to refuse rather than guess (zero-row UPDATE, absent shadow row in the fallback, a missing source row on completion), an unnamed column reaching split reads like it belongs in the same category: refuse it.

deleteAllSQL's comment says "the bigint array" without the reasoning the copier's equivalent carries. pkg/copier/chunker.go explains that parameters are declared bigint whatever the key's integer type and why the operator family makes that safe. A reader of pkg/applier alone sees a hardcoded ::bigint[] next to four other builders that bind $1 with no cast at all, and cannot tell whether that is a deliberate widening or an assumption nobody checked. A one-line cross-reference would settle it.

The key-only upsert writes a real row version. DO UPDATE SET pk = EXCLUDED.pk makes the row exist either way, which is the point, but on the conflict branch it is a genuine update: a new tuple, WAL, index maintenance, and any triggers on the shadow. For a key-only image whose row already exists with the right key, nothing needed to change. DO NOTHING would not be equivalent if the intent is to lock or touch the row, so if the write is deliberate it is worth saying so in the comment on upsertSQL; if it is not, this is a free saving on a path the applier takes often.

…read

A moved, marker-bearing image is completed from the old key's shadow row
or the source row under its key. Neither row is known to be the stream's:
both are live, so a later source transaction can already have put another
row's value there. The flush no longer writes such a completion. It reads
the value together with pg_current_wal_insert_lsn() in the same statement
and returns the image in Result.Held; Buffer.Hold keeps it buffered with
the completion pending, any event that replaces the key or reuses its old
key drops the completion, an UPDATE keeps it pending behind its own
values, and Buffer.Release(passed) fills the markers still left once the
stream has delivered every change committed through the read. The next
Drain flushes the whole image (D13, CO-8).

A moved image neither row can complete is written as a delete marker at
its key, named in Result.Skipped, instead of being left out: the source
holds no row under that key, so a row the shadow may hold there is stale.

The column-wise pass writes every delete marker before any image, so a
row that moved to a lower key while keeping a unique value finds its old
row gone before its upsert. The fallback now also catches SQLSTATE 23P01,
since an exclusion constraint over a moved value fails like a unique
index.

CompleteFirst returns copies of the batch's entries, never pointers into
them, and the guard tests prove the shadow proof check, the catalog-only
search_path, the extra_float_digits pin on a completion read, and that an
empty batch opens no transaction.

Docs: D13, the pkg/applier package-map row, CO-4, CO-6, CO-8, doc.go.
@Kiran01bm

Copy link
Copy Markdown
Collaborator Author

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

Both blocking findings are fixed by one change to D13: a moved image's completion is held in the buffer with the WAL position of its read and written only once the stream has passed it; a skipped image is a delete at its key. N1, N2, and L2 are fixed; N3 is fixed for the branches that matter to correctness and the rest, with L1 and L3, are tracked as internal follow-ups.

# Finding Status Explanation
B1 A skipped image leaves the deleted predecessor's row under Key, and the next move away from Key completes from it fixed completeMoved now writes a delete marker at Key for an image neither row can complete, in its filtered slice, and names the key in Result.Skipped; the batch's entry is untouched so Requeue still puts the image back. The suggested test is in as TestFlushSkippedMoveLeavesNoPredecessorRow, with the assertion that the second flush completes key 10 with row 3's own document.
B2 A source completion can take another row's value, and a later flush spreads it fixed D13 amended as suggested: the completion read (shadow row under the old key or source row under Key) carries pg_current_wal_insert_lsn() in the same statement, the flush returns the image unwritten in Result.Held, Buffer.Hold keeps it buffered with the completion pending, any event that replaces the key or reuses the old key drops the completion, an UPDATE keeps it pending behind its own values, and Buffer.Release(passed) fills the remaining markers once the stream has delivered through the read. One gate covers both reads: the shadow row was copied before the flush read it, so Position needs no per-chunk read position (conservative on the shadow side; recorded as a veto item). The suggested test is in as TestFlushSourceCompletionNeverTakesAnotherRowsValue and asserts the first hold carries the other row's document, the move away from 10 drops it, and key 11 ends with row 3's own; TestFlushHeldShadowCompletionIsDroppedWhenTheOldKeyIsReused covers the shadow-side blind spot. The "overlays it" sentence in D13, the matching veto item, and the CO-6 / CO-8 Enforced lines are rewritten.
N1 Every move to a lower key on a table with a unique index takes the fallback fixed applyColumnWise writes every delete marker before any image. TestFlushMoveToALowerKeyStaysColumnWise is in; TestFlushFallbackCompletesMarkersFromTheShadow now uses a real exchange so it still covers the fallback.
N2 An exclusion constraint has the CO-6 problem but gets neither the fallback nor ErrBatchDeferred fixed constraintCollision matches 23P01 alongside 23505 by SQLSTATE; TestFlushConvergesAnExclusionExchange is in and converges through the fallback. The DEFERRABLE INITIALLY DEFERRED unique form, which reports at COMMIT, is tracked as an internal follow-up and recorded as a veto item.
N3 Several branches of the copied guard have no flusher test fixed New guard_integration_test.go: a decoy = ahead of pg_catalog on the connection's search_path is ignored (CO-9); a connection at extra_float_digits = 0 still completes a float8[] marker holding 0.1 + 0.2 with every digit; NewFlusher refuses every malformed shadow proof, a nil lock, and a lock for another table; an empty batch takes no connection from the pool. Removing the search_path or extra_float_digits pin fails the corresponding test. The ACCESS SHARE hold, the post-commit lock-loss check, and FOR UPDATE on the completion read are not pinned by a test in this package yet; tracked as an internal follow-up.
L1 Nothing bounds a batch, and a batch is one transaction with one round trip per entry deferred The bound belongs with the stream owner's drain (Drain(pos, limit) keeping move pairs together), and set-based statements follow it; both land with the scheduling PR. Recorded as a veto item; tracked as an internal follow-up.
L2 The completions a host should trust least are reported as counts, without keys fixed Result.Held carries each held image with its completion, source, and read position — the stream owner needs them for Hold — and Result.Skipped carries the keys written as deletes.
L3 The guard now has three copies, and this one is the least tested deferred This copy now has its own tests (N3). Sharing the guard as an internal package beside chunksql is a cross-package change; recorded as a veto item and tracked as an internal follow-up.
Verified Column-wise path, savepoint fallback, Requeue, guard parity, CO-4 / CO-5 / ST-6 / LK-1 / LK-2 upheld no action —

Source: block/pg-sprite#149, review comments 6031869738 and 6031870515 and review 5438186038 at head e93aa1d; fixes in the follow-up commit

@Kiran01bm
Kiran01bm enabled auto-merge (squash) October 8, 2026 00:31
@Kiran01bm
Kiran01bm merged commit 5551a15 into main Oct 8, 2026
16 checks passed
@Kiran01bm
Kiran01bm deleted the kiran01bm/cs8-apply branch October 8, 2026 00:36
Kiran01bm added a commit that referenced this pull request Oct 8, 2026
* origin/main:
  migrate: let only the caller's cancellation end the unknown-outcome test's lock wait (#145)
  applier: flush drained batches column-wise with the unique-move fallback (#149)

# Conflicts:
#	docs/copy-and-swap-design.md
Kiran01bm added a commit that referenced this pull request Oct 8, 2026
…bm/cs7-pgoutput

* origin/kiran01bm/cs7-slot:
  decode: prove the server, verify the publication's shape, type the refusals
  schemachange: re-derive the swapped-table proof from the catalog for the post-swap resume (#146)
  migrate: let only the caller's cancellation end the unknown-outcome test's lock wait (#145)
  applier: flush drained batches column-wise with the unique-move fallback (#149)

# Conflicts:
#	SAFETY.md
#	docs/copy-and-swap-design.md
#	docs/invariants.md
#	pkg/decode/doc.go
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