Repository navigation
Conversation
Add the slot lifecycle to pkg/decode, the first use of the pinned jackc/pglogrepl dependency SAFETY.md records for this package. CreateSlot takes the copy-and-swap target preflight minted and creates both halves of the route's logical-decoding state under the derived name: the single-table publication on the caller's pool first, so a role that may create a slot but lacks CREATE on the database is refused with nothing to reap, then the logical pgoutput slot with EXPORT_SNAPSHOT on a dedicated replication connection. The connection is proven to be on the target's database (IDENTIFY_SYSTEM) before the slot exists anywhere; the pool is proven the same way before the publication. The returned Slot carries the exported snapshot name and the consistent point and keeps the replication connection open, because the snapshot lives only as long as that walsender's transaction; Close ends it and the slot persists. A logical slot of the derived name already in the target's database is the route's own earlier slot and is reported as *SlotExistsError for the caller to resume on or drop. A publication of the name that publishes anything but the target is ErrForeignDecodingState, never adopted. DropSlot drops only a logical slot of the pool's own database — any other slot of the name is ErrForeignDecodingState and a name outside the engine's pgsprite_<8hex> shape is ErrInvariantViolation — with DROP_REPLICATION_SLOT ... WAIT on its own replication connection, so a walsender still streaming is waited for under the caller's context alone; a wait the context ends is that error, never a report that the slot is gone. An absent slot is success, so the drop is idempotent. The publication goes with it. InspectSlot reads the slot's pg_replication_slots row: owning database, holder, wal_status as typed constants including the server's own 'lost' verdict for ST-4, restart and confirmed-flush positions, and the WAL the slot retains measured from restart_lsn to the write position. dbconn.ConnectReplication dials the replication-mode connection with the pool's TLS posture, connect timeout, startup parameters, and BeforeConnect hook, as a bare pgconn: a walsender accepts only the simple protocol, and replication commands are not bounded by statement_timeout at all. The depguard core rule gains a decode variant admitting pglogrepl there and nowhere else; pglogrepl is pinned to a commit as the module publishes no tags. Integration tests on a dedicated wal_level=logical server, PG14-18: the exported snapshot predates a row written after the slot and dies with the connection; a new slot has confirmed exactly its consistent point; the publication is exactly the target; a second create reports the route's own slot and reuses the publication; a role without CREATE on the database is refused before any slot exists; a replication connection or a pool on another database is refused; a quiesced target and the zero target are refused; drop is idempotent and removes the publication; drop waits while a START_REPLICATION holder streams and completes when it closes; a drop cut off by its context leaves slot and publication in place; a foreign database's slot and a physical slot of the name are refused and left; retained WAL grows with writes from an unmoving restart_lsn; wal_status reads 'lost' after the server discards the slot's WAL; a wider or FOR ALL TABLES publication is refused with no slot.
…back OpenStream decodes the route's slot with pgoutput on a dedicated replication connection proven to be on the target's database, from the consistent point or a checkpointed applied position, into one ChangeEvent per committed row change. A text value or NULL is a present column and the unchanged-TOAST marker an absent one; an UPDATE that moved the key carries OldKey under either admitted replica identity. A relation other than the target, a tuple that does not line up with it, or a change outside a transaction fails closed; a changed table shape is ErrSourceShapeChanged and TRUNCATE is ErrUnsupportedChange. The stream keeps Delivered — moved by a commit or by a keepalive between transactions, never past an unyielded change — apart from Confirmed, the one position it reports to the server. Confirm refuses a regression or a position beyond Delivered, and a keepalive reply carries only what the caller confirmed, so the slot's resume point moves on the caller's word alone.
…fusals
CreateSlot and DropSlot took a replication config and a pool as separate
inputs and tied them together by database name alone. A name identifies a
database only within one cluster, so a config pointing at a same-named
database on another cluster passed the check: CreateSlot created the slot
there, and DropSlot counted the other cluster's "no such slot" as a drop
while the pool's slot stayed, holding WAL, with its publication gone. Both
now prove the connection and the pool are sessions of one database on one
cluster — IDENTIFY_SYSTEM's system identifier against pg_control_system()
— before any slot command, and DropSlot re-reads the catalog after an
accepted drop and refuses to report success while the slot is still there.
The publication check counted tables in pg_publication_tables, which says
which tables and not which changes: a publication of the derived name
publishing inserts only, no truncate, a row filter, a column list, or a
whole schema was adopted as the route's own, and the changes it left out
never reached the stream. DropSlot skipped the check entirely and dropped
whatever publication wore the name. One predicate — every operation, no
root publishing, no filters, exactly one plain table membership, and on the
create path that table being the target — now gates both paths; DropSlot
refuses a publication that fails it and leaves it as found.
Refusals are typed so an orchestrator can route on them: *ForeignStateError
{Object, Name, Detail} wrapping ErrForeignDecodingState, and
*PublicationPrivilegeError naming the privilege and the database and
wrapping the server's error. A slot the server created but described
unusably is dropped before the error is returned, so no CreateSlot error
leaves a slot. SlotStatus.RetainedBytes becomes Retained{Bytes, Known} so
an unset restart_lsn is not read as zero bytes retained.
The derived name hashes the pgx-quoted database.schema.table instead of the
dotted join, so two tables whose dotted spellings coincide derive different
names.
Tests cover each: same-named database on a second cluster for both create
and drop, each narrower publication shape, another table's publication,
the foreign publication left by the drop, the context-cut create leaving
no slot, the typed errors, the name derivation, and the holder's PID in
the fixture. Docs: ST-3 and D11 state the server proof and the
publication-scoped drop.
…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
…tream end Positions order transactions by their commit. pgoutput sends a transaction whole when it commits, so a change's own LSN can lie below a position the stream already delivered and the caller confirmed; the server replays by commit position, so nothing is lost, but the contract said otherwise and the applier's confirm bound was keyed on the change's LSN. Every ChangeEvent now carries the Delivered it arrived with, Buffer keys FirstLSN on it, and the Delivery, Confirm, package, ST-4, CO-5, SAFETY.md and design-doc text is restated in commit terms. Confirm sends first and records the position only when the server was told; a send that fails latches the stream, as it already did from a keepalive. sendStatus sets write, flush and apply explicitly. A server ending replication with the slot intact is ErrStreamEnded, not a violation. ServerWALEnd exposes the server's position for lag. Preflight now requires two free WAL senders for a copy-and-swap run: the slot's snapshot connection and the decoding stream hold one each. Tests: an interleaved transaction below the confirmed position composes with Buffer.OldestPending; Confirm after stop and after a failed send; the nine stream refusals (outside a transaction, another relation, quiesced database, renamed source, key leaving the identity, decoder error, the caller's deadline); the never-moves-backwards assertion waits for the walsender's flush position; preflight at max_wal_senders 1 and 2.
Add Catchup to pkg/applier: the loop that consumes the decode.Stream into the Buffer, drains it against the copier's Position, flushes, and confirms. It is the stream's only consumer and the first piece that drives the buffer and the flusher end to end against a live copy. Each cycle gathers changes until MaxChanges are buffered or Interval has passed, releases held images at the stream's delivered position, reads the copier's Position after the last Add — so a key the copier cuts between the read and the drain is judged in flight and waits for its chunk rather than being discarded (CO-4) — drains, and flushes the batch in the flusher's guarded transaction. A batch a constraint refuses twice comes back as ErrBatchDeferred and is requeued for the next pass; images the flush completed but did not write are held. The cycle then confirms the lesser of the delivered position and the buffer's oldest pending position, each keyed on the Delivered the change arrived with, so the slot never advances past a change still buffered, held, or deferred and never falls below what the stream already confirmed (ST-4). An empty cycle on a quiet table still confirms, so the slot releases WAL while nothing changes. Run runs once per catch-up, since its stream is single-use, and returns the lock session's loss ahead of whatever the cancelled flush reported, or the context's cause wrapped with the table, so a caller can tell a stop it asked for from a stream that failed. Status is one consistent snapshot of the positions and counters; Lag is how far the confirmed position trails the server's write position. While it runs, the catch-up is the tracker's progress.WorkSource for a catch-up step and stops being one before Run returns. Work reports changes_applied, changes_buffered, and lag_bytes, all read from memory so a poll never waits on the database. The progress JSON contract gains the catch-up operation and the three counters at format_version 6, and docs/progress-report.md pins the new shape. The convergence tests wire the stage the way the orchestrator will: slot under the table lock before the shadow is built, stream from the consistent point, copier and catch-up sharing the lock session, the catch-up started before the copier so every key is uncut at first. Each runs the workload generator against the source, quiesces by waiting for the confirmed position to reach the server's write position, stops the run, and compares the source with its shadow row for row: mixed insert/update/delete with TOAST rewrites across the whole copy (CO-4, plus the tracker's work before and after Run), primary-key moves (CO-4, D13), hot-row contention (CO-5), unchanged-TOAST updates (CO-8), and unique-key swaps through the whole-row fallback (CO-6). The workload generator gains a KeyMove mutation that moves a row to a fresh identity value. Key-move and unique-swap loads start after the copy has landed, because a copy chunk's own insert collides with a stale shadow row the flush has not yet deleted — the copier's race, not the catch-up's. Unit tests cover the option defaults and refusals, the constructor's refusals, the confirm bound on every side of the oldest pending position, Lag, and Work mirroring Status.
ServerWALEnd was overwritten by every change with that change's own WAL position: for logical decoding the walsender writes the record's position into both XLogData fields, not its own send position. Under commit ordering a change can lie below Delivered, so the lag the doc defined wrapped below zero, and read as zero after every commit. Only a keepalive moves it now, it never reads below Delivered, and the doc, SAFETY.md, and the design row name it for what it is: the walsender's send position, a lower bound on lag — lag against the server's WAL end is a pool-side measurement. OpenStream now takes the pool and runs the same cluster-and-database proof as CreateSlot and DropSlot. The derived slot name hashes database, schema, and table, so the same table on two clusters shares a name and a database-name check alone cannot tell them apart. proveSameServer returns the proven identity so the target's database is checked without a second IDENTIFY_SYSTEM. ErrStreamEnded's doc no longer promises the slot keeps the last confirm across the shutdown: a confirm that only moves confirmed_flush is not always written out, and a resume starts from its own checkpoint. Tests: a unit and an integration test for ServerWALEnd, replacing the one that pinned the per-change overwrite; a two-cluster refusal for OpenStream; and a buffer test that keys an UPDATE of an unheld key, a DELETE, and a key move on Delivered.
|
You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard. |
morgo
left a comment
There was a problem hiding this comment.
🤖 Automated adversarial review, posted on Morgan Tocker's behalf.
Approving. I went at the confirm bound hardest, since a bound one position too high silently loses a change and the row-for-row oracle would not catch it if the lost change were overwritten later. It holds, and the reason it holds is not the one the code makes obvious. The notes below are a coverage claim in docs/invariants.md that is broader than the tests behind it, a Run signature whose success path is always an error, and a counter that freezes exactly when an operator needs it moving.
What I tried to break and could not
Confirming a position a buffered change still sits at. confirmBound lowers to oldestPending only when it is strictly below delivered, so a pending entry recorded at exactly the current delivered position does not lower the bound. If an entry's position were its own commit LSN, that would confirm an unapplied transaction and lose it on resume.
It is not. decode sets ev.Delivered = s.delivered when it yields the change, and raiseDelivered runs on the commit message, which arrives after every change in its transaction. So a change's recorded position is the end of the previous committed transaction, strictly below its own. Equality can therefore only mean no commit has been processed since the change arrived — in which case everything at or below the bound genuinely is applied or discarded, and the bound is right. The >= is load-bearing and correct, and it is the single most dangerous line in the PR to change later. The unit tests pinning the bound below, equal to, and above the oldest pending are the right three cases.
A keepalive raising Delivered past a buffered entry. This is the other way the bound could outrun the buffer. raiseDelivered(s.serverWALEnd) on a keepalive does move Delivered independently of any transaction, but the buffered entry's own recorded position does not move with it, so oldestPending < delivered and the bound is pulled back down. The stall that follows is the intended ST-4 trade and lag_bytes is the signal for it.
Hold and Requeue landing after the drain but before the confirm. The ordering in cycle is what makes the bound honest: Drain removes entries, Hold(result.Held) and Requeue(batch) put the unapplied ones back, and only then does confirm read OldestPending(). Reading the oldest pending before those two calls would see an emptier buffer and confirm past entries that were about to return to it. Worth a comment saying the confirm must follow them, because the three calls read as independent steps.
Reading the copier position before the last Add. The // INV: CO-4 comment sits on the right line: c.positions.Position() is evaluated as the argument to Drain, after gather has returned and after Release. The before/after diagram in the description is the failure this ordering prevents and it is accurate.
The progress poll racing the cycle. Work goes through c.Status(), which takes c.mu, and every writer (recordFlush, recordDeferral, recordPositions) takes the same lock, so a poll never sees a half-updated Status. report's stop is called by defer before Run returns and is documented to fence an in-flight poll, so no poll completes against a finished catch-up. That is the part a -race run would not necessarily have caught on its own.
Lag underflowing. if s.ServerWALEnd <= s.Confirmed { return 0 } before the subtraction. ServerWALEnd can legitimately sit below Confirmed under commit ordering — decode's own comment says the server's write position can lie below Delivered — so this guard is not defensive, it is reachable, and without it lag_bytes would report a number near 2^64 to an operator deciding whether to cut over.
Option validation being bypassable. withDefaults runs before validate, and the comment says so, so a zero Interval becomes 250 ms rather than failing the sub-millisecond check. A negative interval still fails. The refusal message telling the caller to pass zero for the default is the right shape.
The convergence tests being weaker than they read. Two of the five are, but the test file says so plainly: the key-move and unique-swap loads start after the copy lands, and the comment explains the 23505 collision that makes mid-copy versions fail and attributes it to the copier rather than this stage. That is a level of disclosure most PRs this size do not reach, and it is why the note below is about docs/invariants.md and not about the tests.
Non-blocking
docs/invariants.md now cites coverage the tests do not provide, and names only one of the two open holes. The new CO-4 paragraph says the race is ordered rather than locked out, backed by "convergence tests under mixed load, key moves, hot-row contention, unchanged-TOAST updates, and unique-key swaps", and the test-plan table row for CO-4, CO-5, CO-6 and CO-8 now points at those tests. The *Open:* sentence that follows names the synchronous-commit visibility window — but not the mid-copy key move, which the test file documents as a reproduced failure (SQLSTATE 23505, for both primary-key moves and unique-value swaps) and works around by starting that load only after the copy has landed.
So a reader who goes to the invariants doc to ask "is a key that moves mid copy covered?" is told yes by the evidence list and finds nothing in the *Open:* text to correct it. The limitation lives in the test file and in this description, neither of which that reader is in. In a repo where the invariants doc is the safety contract and where you have already written one *Open:* paragraph for the other hole, the symmetrical fix is a second sentence: a key move or unique-value swap that straddles an in-flight chunk collides with the stale shadow row on the copier's own insert, the catch-up's half of the exchange is correct, the resolution is on the copier, and the convergence tests for those two races run after the copy lands. One sentence, and the doc stops over-claiming.
Run can never return nil. The loop is for { if err := c.cycle(...); err != nil { return c.finish(ctx, err) } }, so every exit is an error — including the expected one, where the orchestrator watches lag_bytes reach zero, cancels, and gets catch-up on s.t stopped: context.Canceled. A method whose only success path is a wrapped context.Canceled is easy to mis-handle, and there is no caller in the tree yet to establish the pattern. The doc comment describes the behaviour but does not say outright that nil is unreachable. Either say so in one line, or give the caller a sentinel to compare against, so the first orchestrator does not have to decide for itself which errors mean "this worked".
Related and smaller: finish attributes any error to the caller whenever ctx.Err() != nil. That is the documented precedence and it is the right default, but it means a genuine stream failure arriving in the same cycle as a cancellation is reported as a clean stop. Since the expected shutdown always takes this branch, a real failure during shutdown is the case most likely to be swallowed. Wrapping the original error alongside the cause — rather than replacing it — would cost nothing and keep that case legible.
changes_buffered and lag_bytes freeze while a cycle is in flight, which is when they matter most. Status only moves in recordFlush, recordDeferral, and recordPositions, all of which run after Flush returns. A flush stuck behind a lock wait, or a very large batch, leaves Status pinned at the previous cycle's numbers for as long as it takes. The field comments are honest — "after the last cycle", "as of the last cycle" — but the operator-facing reading of a lag_bytes that has stopped changing is "caught up and steady", which is the opposite of what a stalled flush means. The counters are the cutover signal, so the failure mode is an operator cutting over on a stale zero.
Recording positions at the top of each cycle as well as the end would keep lag_bytes and ServerWALEnd moving independently of the flush, since neither depends on the flush's outcome. Failing that, a timestamp on the snapshot would let a consumer tell a steady value from a stuck one.
gather drops delivery.Delivered and the position key lives elsewhere. c.buffer.Add(*delivery.Change) passes only the change, so the fact that the entry is keyed on the Delivered it arrived with — the fact the whole confirm bound rests on — is established in decode by ev.Delivered = s.delivered and is invisible at the one place a reader of catchup.go would look for it. confirmBound's comment explains why the buffer's remaining entries bound the confirmation but not what their positions are. One sentence pointing at Change.Delivered would connect the two halves of the argument.
Why
pkg/applierhas aBufferthat merges the stream into one entry per key and drains it against the copier'sPosition, and aFlusherthat applies one batch under the table lock — but nothing drives them.pkg/decodehas aStreamwith no consumer. This PR lands the catch-up loop CO-4 and ST-4 describe: the one consumer of the stream, which gathers, drains against a copier position read after the last add, flushes, requeues or holds, and confirms a position the slot can safely resume from. It is also the first time the buffer and the flusher run end to end against a live copy under concurrent writes, so the convergence tests the design promised for CO-4, CO-5, CO-6, and CO-8 land here, one per race, each comparing the source with its shadow row for row.Stacked on
kiran01bm/cs7-pgoutput(#151); this PR's diff is the loop, its tests, and the progress counters.Before / after
Source
app.workload, schema changeALTER TABLE … DROP COLUMN label, copier chunking 50 rows at a time, workload writing throughout.Why the position is read after the last add, with chunk
[1496, 1545]in flight:Read the position before the add and 1500 can be judged uncut while the copier's
SELECThas already run: the change is discarded, the chunk lands the stale row, and nothing replays it. Read after, and the key is in flight, so it waits and is flushed once the chunk lands.What
pkg/applier/catchup.go—PositionSource,CatchupOptions{Interval, MaxChanges, Clock, Tracker}with defaults and refusals (ErrInvalidCatchupOptions),NewCatchup(a nil stream, flusher, or position source is anErrInvariantViolation),Status{Delivered, Confirmed, ServerWALEnd, Buffered, Deferred, Held, Applied, Discarded, Fallbacks, BatchesDeferred}andStatus.Lag,Run(ErrCatchupAlreadyRunon a second call), the cycle (gather→Release→Drain→ flush →confirm),confirmBound, andfinish(lock loss first, then the context's cause wrapped with the table, else the stream's error).pkg/applier/catchup_work.go—Work, theprogress.WorkSourcethe catch-up registers for the lifetime ofRun; all three counters read from memory.pkg/progress—OperationCatchUp,Work.ChangesApplied / ChangesBuffered / LagBytes,FormatVersion6;docs/progress-report.mdgains the engine-measured catch-up section, the counter table, the operations row, and a fourth JSON example.internal/testutil/workload.go—Mix.KeyMoveandSummary.KeyMoves: anUPDATE … SET id = nextval(pg_get_serial_sequence(…))on a random row, so the primary key moves to a fresh identity value.catchup_test.go(defaults, refusals, confirm bound below / equal / above the oldest pending and with nothing pending,Lag,WorkmirrorsStatus);catchup_integration_test.go(fixture: slot under the table lock before the shadow is built, stream from the consistent point, copier and catch-up sharing the lock session, catch-up started before the copier; five convergence tests andTestCatchupRunsOnce);work_source_test.go(TestSnapshotJSONShapeForACatchUpStep).pkg/applier/doc.go;SAFETY.mdpkg/applierrow (scheduling no longer planned; ST-4 added); design package map;docs/invariants.mdCO-4 enforcement and the test-plan row.Decisions to veto
min(Delivered, OldestPending), both keyed on theDelivereda change arrived with. Nothing below the bound is still in the buffer, held, or deferred, so a stream reopened from it misses nothing; keying on arrival position rather than the change's own LSN keeps the bound from falling below what the stream already confirmed (the interleaved-transaction case decode: stream the slot into change events with caller-confirmed feedback #151 documents). The slot therefore stalls while any key waits on an in-flight chunk or a deferred batch; that is the intended ST-4 trade, andlag_bytesmakes the stall visible.Add, not mutual exclusion between copy and flush. The design's CO-4 text called the enforcement "scheduling that makes a chunk copy and a backlog flush mutually exclusive". Ordering the read is enough: an in-flight key is never discarded, a landed key's flush and a later chunk cannot overlap on the same key because the copier never re-reads a landed chunk, and the flusher already takes the copier's guard inside its transaction. No lock between the two stages, so a slow flush never stalls the copy and a slow chunk never stalls the flush. The CO-4 invariant text is updated to say so.INSERTof the row under its new key then collides on the unique index with the stale shadow row the flush has not yet deleted (SQLSTATE 23505, reproduced for both primary-key moves and unique-key swaps). That is the copier's race — the catch-up's part of the exchange is already correct — and its resolution is tracked as an internal follow-up on the copier. Mixed load, hot rows, and TOAST updates run across the whole copy.synchronous_commitwaiting on a standby, the walsender can send a transaction's commit while its backend is still in the proc array, so a chunk snapshot taken after the stream delivered the change may not see it: the drain judges the key uncut and discards the change, the chunk copies the pre-change row, and nothing replays it. Reading the position after the last add does not close this — the copier's snapshot comes later still. The CO-4 text names the window as open; closing it (discard an uncut entry only when its transaction's xid is below the copier's snapshotxminand not in itsxip, or hold uncut entries until the chunk that covers the key claims) is tracked as an internal follow-up. The convergence tests run without synchronous replication and do not exercise it.Runreturns the context's cause, wrapped, when the caller stopped it. The stream reports a cancelled receive as its own error; a caller that ended the context getscatch-up on s.t stopped: <cause>instead, and a lost table lock outranks both, so the orchestrator can tell a stop it asked for from a stream or lock failure.Runruns once because the stream is single-use; a resume opens a new stream and a new catch-up.MaxChangesbounds the buffer and the flush transaction under a burst. Both are options; an interval below one millisecond or a zeroMaxChangesis refused rather than defaulted silently. Both are callers' knobs later, not CLI flags now.Status, in the progress JSON.changes_applied(exact, monotone),changes_buffered(the backlog an operator watches drain before cutover), andlag_bytes(how far the confirmed position trails the server's write position). Deferred / held / fallback counts stay onStatusfor the library caller; the JSON contract grows when a consumer needs them. This is a contract change, henceformat_version6.Discardedcounts entries, not events. The buffer already merged events per key before the drain, so what it drops is one entry per uncut key; this is the figure that says the copy's own read is doing the work mid copy, and the mixed-load test waits on it before starting the copier.Verification
make lint0 issues;gofmt -lclean;go build ./...;make test-unitgreen (docs guards included).go test -race -count=1 ./pkg/applier/ ./pkg/progress/ ./internal/testutil/ ./pkg/decode/green on PG 16 (applier ~63 s, decode ~67 s).scripts/test-flaky.sh 'TestCatchupConverges.*' 10 ./pkg/applier/— all five convergence tests, 10/10 iterations with-race, ~18 s each.Discarded> 0 before the copy,Applied> 0 after,BufferedandBatchesDeferredzero at the end; the tracker's poll carrieschanges_appliedmid-run and no work afterRunreturns;WorkafterRunstill reports the applied count); primary-key moves after the copy (KeyMoves> 0); hot-row contention on the ten lowest keys (HotRowFraction1); unchanged-TOAST updates (ToastRewriteFraction0); unique-key swaps after the copy (Fallbacks> 0); a secondRun→ErrCatchupAlreadyRun; a stop → wrappedcontext.Cancelednaming the table.Lagwith and without a server position,WorkmirroringStatus;pkg/progressJSON pins for the catch-up step atformat_version6.