Skip to content

decode: stream the slot into change events with caller-confirmed feedback - #151

Merged
Kiran01bm merged 10 commits into
mainfrom
kiran01bm/cs7-pgoutput
Oct 8, 2026
Merged

Kiran01bm merged 10 commits into
mainfrom
kiran01bm/cs7-pgoutput

Conversation

@Kiran01bm

@Kiran01bm Kiran01bm commented Oct 8, 2026 •

Copy link
Copy Markdown
Collaborator

Why

pkg/decode can create, drop, and inspect the route's slot, but nothing reads it: ChangeEvent has no producer, and the applier, the catch-up loop, and the cutover gate all wait on one. This PR lands the pgoutput stream reader that ST-3, ST-4, CO-4, and CO-8 describe: one connection proven to be a session of the pool's cluster on the target's database, decoding the slot from the consistent point (first run) or the checkpointed applied position (resume) into one ChangeEvent per committed row change — with every column's presence carried so the server's unchanged-TOAST marker is never mistaken for a value, and the key a moved row had — and holding the slot's confirmed position where the caller says applied is, never where the server says WAL is. The reaper and the checkpoint-driven resume are the next PRs.

Stacked on kiran01bm/cs7-slot (#150); this PR's diff is the stream only, plus the preflight headroom it needs.

Before / after

Target app.ledger in database shop; slot and publication pgsprite_3f2a9c01.

Before                                        After
──────                                        ─────
decode.CreateSlot ──► Slot{ConsistentPoint}   decode.OpenStream(cfg, pool, target, from)
                      (nothing consumed it;     ├─ dbconn.ConnectReplication(cfg)      ◄── fresh conn; Slot's conn
ChangeEvent, Column, LSN — contract types       │                                          and snapshot untouched
with no producer)                               ├─ IDENTIFY_SYSTEM ≡ pool's cluster,     ◄── proof (ST-3): same cluster and
                                                │    dbname = 'shop'                           database as the pool
                                                └─ START_REPLICATION SLOT pgsprite_3f2a9c01
preflight: one free WAL sender                        LOGICAL from (proto_version '1',
                                                      publication_names 'pgsprite_3f2a9c01')
                                                   ──► Stream{Start=from, Delivered=from, Confirmed=0}

                                              preflight: two free WAL senders ◄── the slot's snapshot
                                                                                   conn and the stream

                                              stream.Next(ctx, wait) → Delivery{Change, Delivered}
                                                Begin ───────────────────────── consumed, inTransaction
                                                Relation(app.ledger) ─────────── cached; another table → ErrInvariantViolation
                                                                                 same id, new shape   → ErrSourceShapeChanged
                                                Insert/Update/Delete ─────────── Delivery{Change: &ChangeEvent{..., Delivered}}
                                                   't' value │ 'n' NULL ──────── Column{Present: true}
                                                   'u' unchanged TOAST ───────── Column{Present: false}        (D6 / CO-8)
                                                   key moved ('K' or 'O') ────── OldKey set
                                                Commit ───────────────────────── Delivery{Change: nil, Delivered: TransactionEnd}
                                                Truncate ─────────────────────── ErrUnsupportedChange
                                                Keepalive, between txns ──────── Delivery{nil, Delivered: ServerWALEnd}
                                                Keepalive, inside a txn ──────── Delivered unchanged; ServerWALEnd() still moves
                                                any change ───────────────────── ServerWALEnd() unmoved (message carries the
                                                                                 change's own position); never < Delivered
                                                Keepalive ReplyRequested ──────── standby status with Confirmed only
                                                nothing within wait ──────────── Delivery{nil, Delivered}  (conn stays usable)
                                                CopyDone / CommandComplete ───── ErrStreamEnded (slot intact; not a violation)
                                                WARNING from the server ──────── ErrInvariantViolation (ST-4): the walsender will
                                                                                 withhold changes (PG 18 skips a dropped publication)

                                              stream.Confirm(ctx, lsn)
                                                lsn < Confirmed ──────────────── ErrInvariantViolation (ST-4)
                                                lsn > Delivered ──────────────── ErrInvariantViolation (ST-4)
                                                otherwise ────────────────────── StandbyStatusUpdate{Write=Flush=Apply=lsn}
                                                                                 sent first; Confirmed=lsn only on success
                                                                                 → slot confirmed_flush_lsn moves

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:

session A: INSERT 201 @ 0/1927790 ─────────────────────────── COMMIT
session B:           INSERT 202, COMMIT → Delivered 0/19278E0
caller:                                  Confirm(0/19278E0)
stream:                                                        yields 201, LSN 0/1927790 < Confirmed

Nothing is lost: the server replays by commit position, so a stream reopened from 0/19278E0 yields 201 again. What the caller may confirm while 201 is unapplied is the Delivered it arrived with — carried on every ChangeEvent, and what applier.Buffer keys OldestPending on — never 201's own LSN.

What

  • pkg/decode/stream.go — OpenStream (takes the pool and runs the same cluster-and-database proof as CreateSlot and DropSlot), Stream{Start, Delivered, ServerWALEnd, Next, Close}, Delivery, ErrUnsupportedChange, ErrStreamEnded. Message dispatch (handleCopyData → keepalive / WAL data; handleNotice stops the stream on a warning), per-message decoders, fail (the first error is stored and returned from every later call).
  • pkg/decode/feedback.go — Confirm, Confirmed, sendStatus: the one place a position is reported to the server, as write, flush, and apply alike.
  • pkg/decode/relation.go — the single relation cache built from pgoutput's RelationMessage: refuses a relation that is not the target (ST-3), records column names, key flags, and replica identity; sameShape decides ErrSourceShapeChanged.
  • pkg/decode/tuple.go — decodeColumns ('t' present, 'n' present-NULL, 'u' absent, 'b' refused, count mismatch refused) and decodeKey (the PKColumn() must be a flagged key column; parsed as int64).
  • pkg/decode/types.go — ChangeEvent.Delivered, the stream's position when the change arrived.
  • pkg/decode/server_identity.go — proveSameServer returns the proven identity, so OpenStream checks the target's database against it without a second IDENTIFY_SYSTEM.
  • pkg/applier/buffer.go — FirstLSN keyed on the event's Delivered, so OldestPending never falls below what the stream confirmed.
  • pkg/preflight/copy_swap_environment.go — a copy-and-swap run needs two free WAL senders; the refusal's detail names both holders.
  • Tests — tuple_test.go, stream_test.go (hand-encoded Begin/Commit/Insert bytes: keepalive handling inside and between transactions, unpaired boundaries, a change outside a transaction or for another relation, the Delivered stamp, ErrStreamEnded, ServerWALEnd moved by keepalives only, Confirm refusals without a server), buffer_test.go (every event kind keys its entry on Delivered), stream_fixture_integration_test.go (shared wal_level=logical fixture, a wal_sender_timeout=2s variant, assertWALSenderFlushBecomes), stream_integration_test.go, stream_refusal_integration_test.go, feedback_integration_test.go.
  • Docs — pkg/decode/doc.go; SAFETY.md pkg/decode row; design package map (decode and applier rows); docs/invariants.md CO-5, CO-8, ST-3 (OpenStream joins the cluster proof), and ST-4 (a warning stops the stream); docs/engine-role.md, docs/refusal-classes.md, docs/tcb-model.md on sender headroom.

Decisions to veto

  • OpenStream dials its own connection; there is no Slot.Start. The slot's connection holds the exported snapshot the copier imports; streaming on it would end that transaction and the snapshot with it. A fresh connection means a resume (no Slot in hand) and a first run take the same path. Preflight therefore proves two free WAL senders, not one, so a server with exactly one never passes preflight, builds a shadow, and then fails to open the stream.
  • Next(ctx, wait) is a wall-clock wait, not a context deadline. The caller's context cancels the stream; wait bounds one call so a quiet table never blocks the applier's checkpoint cadence. Implemented with pgconn's timeout-tolerant receive, so the connection is reusable after an elapsed wait.
  • Begin and Commit are not surfaced as events. The applier does not need transaction boundaries to apply rows idempotently (CO-4); it needs a position it may confirm, which the commit's progress delivery gives it. Protocol version 1 is fixed, so every change yielded has already committed.
  • Delivered moves on a keepalive only between transactions. A keepalive's ServerWALEnd can lie inside a transaction the stream is mid-way through; advancing to it would let the caller confirm past a change not yet handed over. Inside a transaction the keepalive is answered but the position holds; ServerWALEnd() still reports it.
  • A keepalive reply carries Confirmed alone — zeros before the first Confirm. The server treats a zero flush position as "no information" and leaves the slot where it is, which is the ST-4 behaviour wanted: the slot moves only on the caller's word.
  • Confirm records the position only after the server was told, and a failed send ends the stream. A send that fails leaves Confirmed where it was, so a retry from a lower durable checkpoint is not refused; and since a connection that cannot carry a status report is in the same state whichever path tried to send one, the failure latches the stream from Confirm as it already did from a keepalive reply.
  • Confirm below Start is accepted. A resumed stream may be handed the checkpoint's applied position, which is below the start it was reopened from after a forwarded resume; the server never moves a slot backwards, so the report is harmless and the test proves the slot stays put.
  • Start() is the position asked for, not the one the server granted. The server forwards a start below the slot's confirmed position; the first delivery's position shows where decoding began. Reading confirmed_flush_lsn before START_REPLICATION belongs to the checkpoint-driven resume, which opens a stream without a Slot in hand, and is tracked as an internal follow-up there.
  • Shape change is a stop, not a reconcile. A second RelationMessage for the target with a different column list or replica identity is ErrSourceShapeChanged, persistent for the stream; TRUNCATE is ErrUnsupportedChange. Whether the route restarts the copy or refuses is the orchestrator's call.
  • An orderly end of replication is ErrStreamEnded, not a violation. The slot is intact with a position at or below the last confirm — a confirm that only moves confirmed_flush is not always written out before a shutdown — and a stream reopened from the caller's own checkpoint continues, since the server forwards a start below the slot's position; a resume can tell this from slot loss.
  • ServerWALEnd is the walsender's send position from its last keepalive, never below Delivered, and a lower bound on lag. A change's XLogData carries the change's own position in both fields, which under commit ordering can lie below Delivered, so a change does not move it; max(·, Delivered) keeps ServerWALEnd − Delivered from wrapping. Even so the keepalive reports how far the walsender has decoded, not how far the server has written, so the figure stays small while a backlog is undecoded; lag against the server's WAL end is a pool-side measurement against pg_current_wal_lsn(), which the catch-up progress source takes up. The decode doc, SAFETY.md, and the design row say so.
  • OpenStream takes the pool and runs proveSameServer. The derived slot name hashes database, schema, and table, so the same table on two clusters — shards, or a staging and a production copy — shares a slot name, and a config pointing at the other cluster would decode that cluster's slot as this target's changes. The database-name check from IDENTIFY_SYSTEM alone cannot tell the two apart; the proof CreateSlot and DropSlot already run can. The signature change reaches only test callers today and the stacked catch-up PR.
  • A warning from the server in copy-both mode stops the stream. Through PostgreSQL 17 a dropped publication fails the decoder (42704); from 18 the walsender skips loading it with a WARNING (55000), sends nothing for the change, and the next keepalive would move Delivered — and the caller's Confirm the slot — past a change the slot will never resend. The walsender sends a warning only to say it will not send something, so any warning is ErrInvariantViolation with the server's SQLSTATE reachable; a notice below warning is informational and ignored. Matching on 55000 alone would be narrower; failing on every warning is the fail-closed choice.
  • One relation cache, not a map. The publication is FOR TABLE <target>; a second relation id is already an invariant violation, so a map would only hide the check.
  • Event LSN is XLogData.WALStart, the change's own position, which can lie below Delivered and Confirmed (see above). Delivered after the commit is TransactionEndLSN. Adjacent transactions share a boundary (the next WALStart can equal the previous TransactionEndLSN), which the tests account for.
  • Column keeps Value and Present as fields. An accessor that makes an absent value unreadable is tracked as an internal follow-up to land with the applier's flush path, its first caller.
  • No new refusal class row. ErrSourceShapeChanged, ErrUnsupportedChange, and ErrStreamEnded are library errors; the orchestrator that maps them to verdict reasons is where docs/refusal-classes.md grows.

Verification

  • make lint 0 issues; gofmt -l clean; go build ./...; SKIP_INTEGRATION=1 go test ./... green (docs guards included).
  • go test -race -count=1 ./pkg/decode/ ./pkg/preflight/ ./pkg/applier/ green on PG 16 (decode ~70 s, preflight ~100 s, applier ~54 s); after the review round, go test -race -count=1 ./pkg/decode/ ./pkg/applier/ green again and scripts/test-flaky.sh 'TestServerWALEndNeverTrailsDelivered|TestOpenStreamRefusesAReplicationConnectionOnAnotherServer' 5 ./pkg/decode/ 5/5.
  • Review-round mutants: restoring the per-change serverWALEnd overwrite fails TestServerWALEndIsNotMovedByAChange and TestServerWALEndNeverTrailsDelivered (the latter with the 0/1927790 < 0/19278E0 underflow case); reverting OpenStream to the database-name check alone fails TestOpenStreamRefusesAReplicationConnectionOnAnotherServer, which then falls through to the other cluster's 42704; each of the four FirstLSN sites reverted to ev.LSN fails its subtest of TestBufferKeysEveryEventKindOnItsDeliveredPosition. Earlier rounds: PG_VERSION={14,15,17,18} go test -count=1 ./pkg/decode/ -run 'TestStream|TestOpenStream|TestConfirm|TestKeepaliveReplies|TestReopened' green on each.
  • Delta round on PostgreSQL 18: TestStreamReturnsTheDecodersError expects 42704 before 18 and ErrInvariantViolation carrying 55000 from 18; with the warning handling removed it times out on 18 (16.9 s) and passes with it (1.7 s). TestStreamStopsOnAWarningFromTheServer (unit) fails with the handling removed; TestServerWALEndKeepsAKeepaliveAboveALaterChange fails with the per-change overwrite restored (expected 0xe6, actual 0xd2); TestOpenStreamRefusesAPoolOnAnotherDatabase fails with the target-database check removed (only the server's 55000 is left in the chain). PG_VERSION=18 go test -race -count=1 ./pkg/decode/ and go test -race -count=1 ./pkg/decode/ ./pkg/applier/ on 16 green.
  • go test -count=5 -run 'TestNextReturnsTheCallersDeadline|TestStreamDeliversAnInterleavedTransactionBelowTheConfirmedPosition' ./pkg/decode/ — 5/5; scripts/test-flaky.sh 'TestKeepaliveRepliesDoNotMoveTheSlot|TestStreamDeliversATransactionsChangesBeforeMovingThePosition|TestReopenedStreamReplaysFromTheConfirmedPosition' 5 ./pkg/decode — 5/5 with -race.
  • Integration (stream_integration_test.go): insert yields every column present; a NULL is present with a nil value; an unchanged out-of-line value (SET STORAGE EXTERNAL, repeat('x', 4000), another column updated) is absent; a changed out-of-line value is present; an UPDATE that moves the key carries OldKey under default replica identity and under REPLICA IDENTITY FULL, and omits it when the key did not move; DELETE carries the key and nil columns; all of a transaction's changes are yielded before Delivered reaches its commit; a rolled-back transaction yields nothing; ADD COLUMN → ErrSourceShapeChanged on the next change and from every later call; TRUNCATE → ErrUnsupportedChange; zero target and zero start position refused before any connection.
  • Integration (stream_refusal_integration_test.go): a quiesced database and another database are refused (ST-3); a database of the target's name on a second cluster is refused before START_REPLICATION, by the stream's own proof rather than the other server's missing-slot error; a renamed source arrives as another relation and is refused; a key column leaving the replica identity is refused; a dropped publication surfaces as *pgconn.PgError 42704 before 18 and as an ST-4 violation carrying 55000 from 18; a pool on another database of the same cluster is refused by the stream's own proof; Next returns the caller's own context deadline.
  • Integration (feedback_integration_test.go): Confirm moves confirmed_flush_lsn to exactly the confirmed position; five seconds of keepalive replies on a wal_sender_timeout=2s server leave the slot where it was while Delivered rises past writes to an unpublished table; Confirm beyond Delivered refused and the slot unmoved; a stream reopened from the confirmed position replays the unconfirmed transaction and not the confirmed one; reopened from a forwarded position; Confirm below start accepted and — after the walsender's flush position shows the report was read — the slot unmoved; an interleaved transaction arrives below the confirmed position and Confirm(min(Delivered, Buffer.OldestPending)) succeeds; ServerWALEnd() ≥ Delivered() holds after a commit and after a change written below Delivered; Confirm after the stream stopped returns the stored error; a Confirm whose send fails leaves Confirmed where it was and ends the stream.
  • Integration (pkg/preflight): a free slot with exactly two free senders passes; max_wal_senders=1 is refused with the two-sender detail; a live sender counts against the headroom.
  • Unit: null vs unchanged-TOAST marker, binary tuple refused, column-count mismatch refused, key parsed at the int64 boundary with a non-integer, NULL, or missing key refused, sameShape on id / columns / replica identity; keepalive moves Delivered only outside a transaction and ServerWALEnd always, a change never moves ServerWALEnd (including a change that follows a mid-transaction keepalive), it reads Delivered until the server reports, a NOTICE leaves the stream running and a WARNING stops it with its SQLSTATE reachable; an UPDATE of an unheld key, a DELETE, and a key move each key FirstLSN on Delivered; Commit without Begin, Begin inside Begin, a change outside a transaction, and a change for another relation refused; a change carries the Delivered it arrived with; CopyDone ends the stream with ErrStreamEnded; Confirm regression and beyond-delivered refused; preflight refuses at max_wal_senders − 1 free.

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.
@Kiran01bm
Kiran01bm marked this pull request as ready for review October 8, 2026 02:10
@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.

@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. I reviewed this as the stack delta against kiran01bm/cs7-slot, not against main. I went looking specifically for a way to move the slot's confirmed position past something the caller never applied — the one failure here that destroys data irrecoverably, since the server cannot replay discarded WAL — and I could not construct one. All 16 checks green across PostgreSQL 14–18.

The position discipline is the part that had to be right, and it is. delivered rises in exactly two places, and both are safe for the reason the comments give:

  • on CommitMessage, to TransactionEndLSN, after every one of that transaction's changes has already been yielded;
  • on a keepalive, only when !inTransaction.

The keepalive case is the one worth spelling out, because it rests on a server property the comment asserts without naming: a walsender's keepalive reports sentPtr, the position it has streamed to, not the server's WAL insert position. So ServerWALEnd can never name a position with an undelivered change below it — including during a large transaction, where protocol version 1 buffers until commit and sentPtr therefore stays put, making raiseDelivered a no-op. Holding the position inside a transaction is belt-and-braces on top of that, and correctly so.

Verified rather than assumed:

  • Confirm's bounds are the right two and they close the hole. lsn > s.delivered is the one that matters, and delivered is initialised to from rather than zero, so a caller cannot confirm above the start before any data arrives either.
  • sendStatus really does move confirmed_flush_lsn despite setting only WALWritePosition. pglogrepl fills WALFlushPosition and WALApplyPosition from the write position when they are zero, so the wire message carries all three as the description says. Worth having checked, because logical slots advance on the flush position — had pglogrepl not defaulted, Confirm would have been a silent no-op, and the integration test asserting the slot moves to exactly the confirmed position is what pins it.
  • decodeColumns cannot turn an unchanged-TOAST marker into a value. 't' sets value and present, 'n' sets present with a nil value, 'u' falls through leaving Present false, and every other data type including 'b' is refused rather than defaulted. The column-count mismatch check in front of it means a stale relation cache cannot shift values by one position either.
  • newRelation cannot leave keyIndex dangling. It starts at -1, a missing PK column is refused before use, and the key column must carry the replica-identity flag — so decodeKey's unchecked tuple.Columns[rel.keyIndex] is safe for every relation that exists.
  • Next(ctx, 0) degrades safely rather than killing the stream. An already-expired waitCtx still satisfies all three arms of waitElapsed (pgconn.Timeout, waitCtx.Err() == DeadlineExceeded, ctx.Err() == nil), so a caller whose wait is an unset config field gets an immediate empty progress delivery instead of a permanently failed stream. That is the right direction and it is not obvious from reading Next alone.
  • Every OpenStream failure path closes the connection via errors.Join(err, conn.Close(ctx)), including the database-identity refusal.

Non-blocking

Confirm advances s.confirmed before the send, so the field can outrun what the server was actually told.

s.confirmed = lsn
return s.sendStatus(ctx)

If sendStatus fails, confirmed is already lsn and the stream is not failed, so the field now disagrees with its own doc comment — "Confirmed is the position last reported to the server as applied" — which is exactly what it is not.

Nothing is lost by this: confirmed is only read as the floor for later confirmations and as the value sent, any later successful send carries a position at or above the failed one, and the server never moves a slot backwards. The concrete cost is a legitimate caller getting refused. An applier that fails to confirm and then falls back to its last durable checkpoint — below the position whose send just failed — gets ErrInvariantViolation: confirm X below the Y already confirmed, naming a confirmation that never reached the server. That reads as a bug in the caller, and the remedy it suggests (don't go backwards) is the opposite of what the caller should do.

Setting s.confirmed only after sendStatus returns nil makes the field mean what it says and makes that retry work, at the cost of nothing — a failed send leaves the floor where it was, which is the conservative direction.

Related, smaller: the same sendStatus error is treated two ways. From handleKeepalive it propagates to Next, which latches it in fail and ends the stream permanently; from Confirm it is returned with the stream still live. Both are defensible in isolation, but a connection that cannot send a status update is in the same state either way, and one of the two callers decides that is terminal.

This PR is the first thing to produce Present: false, which makes the Value/Present ambiguity live. An unchanged-TOAST column and a SQL NULL both arrive with Value == nil; only Present separates them, and getting it wrong writes NULL over a value that was never touched — the precise corruption CO-8 exists to prevent. The producer side is right and the comment says so. The types predate this PR so the shape is not yours to change here, but since the applier is the next consumer and the compiler will not catch a switch on Value == nil, it may be worth Column growing an accessor that makes reading an absent value impossible rather than merely documented, before there is a second consumer to migrate.

@aparajon

aparajon commented Oct 8, 2026

Copy link
Copy Markdown
Collaborator

🤖 Review of acf57cf (1/2): 1 blocking, 4 non-blocking.

This is the adversarial correctness pass over the whole diff at acf57cf66a0c4bbfdef3dc650909b25ec2268a67, against its base fee00a1 (#150). I read OpenStream, Next, the message dispatch, Confirm, sendStatus, relation.go and tuple.go, and the applier code that will consume them. Then I ran the stream against PostgreSQL 14, 16 and 18 with two concurrent sessions, confirmed below the slot, terminated the walsender, dropped the publication, and stopped the server under an open stream.

The feedback path is sound where it matters most: nothing the stream sends can move the slot past a transaction the caller has not been handed. A keepalive reply carries only Confirmed, Confirm refuses a position beyond Delivered, and a reopened stream replays every transaction that committed above the confirmed position. The gap is in how the PR describes positions. It says a change's own LSN orders against Delivered, but PostgreSQL orders the slot by commit position. Two sessions are enough to show the difference, and the applier's existing confirm rule is written in the same terms.

On invariants: ST-4 is extended, and with B1 open its new Enforced text states a property the stream does not have. CO-8 is extended to the decoder (presence carried; the TOAST and NULL mutants are killed). CO-4 is upheld (OldKey only on a move, under both identities). ST-3 is upheld for the database and target checks, though N3 shows four of those guards have no test. CO-5 is affected through B1: its OldestPending bound does not compose with Confirm.

Blocking

B1. A change can arrive with an LSN below a position the stream already delivered and the caller confirmed, so the Delivered contract and ST-4's new Enforced text are false. stream.go:35-38, stream.go:208-217, stream.go:365-366, invariants.md:621-626, buffer.go:42-46

Delivery.Delivered promises that "every change at or below it has been yielded". The ST-4 Enforced line says it "never names a position with an unyielded change below it". Confirm says it reports "every change at or below lsn" as applied. pgoutput sends a transaction only when it commits, but each change keeps the WAL position where it was written. Two sessions break all three statements:

  1. Session A inserts key 201 at 0/1927790 and stays open.
  2. Session B inserts key 202 and commits. The stream yields it, and its commit delivers 0/19278E0. The caller confirms that position, and the slot moves to it.
  3. A commits. The stream yields key 201 with LSN = 0/1927790, below both Delivered and Confirmed.

Nothing is lost. The server replays by commit position, so a stream reopened from 0/19278E0 yields key 201 again; I measured that on 14, 16 and 18. But the docs describe the other ordering, and the next PR will build on them:

  • Buffer.OldestPending (CO-5, invariants.md:209-211) bounds the confirm by each entry's FirstLSN, which is ChangeEvent.LSN, on the reasoning that a stream confirmed at or past it "could not replay those entries". Step 3 shows the slot does replay them. After step 3, that bound is 0/1927790. Confirm(min(Delivered, OldestPending)) returns invariant violation: ST-4: confirm 0/1927790 below the 0/19278E0 already confirmed, so the composition the two docs describe fails on an ordinary interleaving.
  • A resume that trusts the stated contract, and skips replayed events at or below its checkpoint, drops key 201. Nothing in the docs warns against it.
  • The unit test comment at stream_test.go:34-36 describes this case: a keepalive "can lie past changes of an open transaction". The inTransaction guard does not cover it, because the stream sees no BEGIN until the open transaction commits. The guard is still right for a transaction being sent. It just is not what keeps step 3 safe; commit ordering is.

Suggested fix:

  • State the contract in commit terms: every transaction that committed at or below Delivered has been yielded in full, and a change's LSN can lie below a position already delivered or confirmed. That wording goes in Delivery, Confirm, doc.go, ST-4's rule and Enforced line, the SAFETY.md row, and the design row.
  • Give the applier a bound that cannot fall below Confirmed. That bound is the Delivered a change arrived with: it is below that change's commit, and the stream never lowers it. Either carry it on ChangeEvent, or have Buffer.Add take the Delivery and key FirstLSN on it. Then OldestPending stays below every unapplied commit, and Confirm never refuses it.
Test case: the documented contract fails on acf57cf, and a replacement pins the real one

The documented contract, as a test:

func TestDeliveredCoversEveryChangeBelowIt(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())

	open, err := f.pool.Begin(t.Context())
	require.NoError(t, err)
	t.Cleanup(func() { _ = open.Rollback(context.WithoutCancel(t.Context())) })
	_, err = open.Exec(t.Context(), fmt.Sprintf(`INSERT INTO %s.ledger (id, note) VALUES (201, 'written first')`, f.schema))
	require.NoError(t, err)
	f.exec(t, `INSERT INTO %s.ledger (id, note) VALUES (202, 'committed first')`)
	first := nextChange(t, stream)
	awaitDelivered(t, stream, first.LSN+1)
	confirmed := stream.Delivered()
	require.NoError(t, stream.Confirm(t.Context(), confirmed))

	require.NoError(t, open.Commit(t.Context()))
	late := nextChange(t, stream)
	assert.Greater(t, late.LSN, confirmed, "every change at or below Delivered has been yielded")

	buf := applier.NewBuffer()
	require.NoError(t, buf.Add(late))
	oldest, _ := buf.OldestPending()
	assert.NoError(t, stream.Confirm(t.Context(), min(stream.Delivered(), oldest)))
}

--- FAIL: TestDeliveredCoversEveryChangeBelowIt (1.43s) on acf57cf (PG 16): "0/1927790" is not greater than "0/19278E0", then the Confirm refusal above.

The contract the stream actually keeps, which passes on acf57cf on 14, 16 and 18 and belongs in the suite once the docs say it:

// A transaction that writes the target before another transaction commits,
// and commits after it, arrives after that commit with a change position
// below what the stream already delivered and the caller confirmed. The
// slot replays by commit position, so a stream reopened from the
// confirmation sees the change again; the position the caller may confirm
// while that change is unapplied is the Delivered it arrived with, never
// the change's own LSN.
func TestStreamDeliversAnInterleavedTransactionBelowTheConfirmedPosition(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())

	open, err := f.pool.Begin(t.Context())
	require.NoError(t, err)
	t.Cleanup(func() { _ = open.Rollback(context.WithoutCancel(t.Context())) })
	_, err = open.Exec(t.Context(), fmt.Sprintf(`INSERT INTO %s.ledger (id, note) VALUES (201, 'written first')`, f.schema))
	require.NoError(t, err)
	f.exec(t, `INSERT INTO %s.ledger (id, note) VALUES (202, 'committed first')`)
	first := nextChange(t, stream)
	require.EqualValues(t, 202, first.Key)
	awaitDelivered(t, stream, first.LSN+1)
	confirmed := stream.Delivered()
	require.NoError(t, stream.Confirm(t.Context(), confirmed))
	f.assertConfirmedFlushBecomes(t, slot, confirmed)

	require.NoError(t, open.Commit(t.Context()))
	late := nextChangeDelivery(t, stream) // nextChange, returning the whole Delivery
	require.EqualValues(t, 201, late.Change.Key)
	assert.Less(t, late.Change.LSN, confirmed, "the change was written below the confirmed position")
	assert.GreaterOrEqual(t, late.Delivered, confirmed)
	require.NoError(t, stream.Confirm(t.Context(), late.Delivered), "the delivery's own position is confirmable")
	require.NoError(t, stream.Close(t.Context()))

	replayed := nextChange(t, f.openStream(t, confirmed))
	assert.Equal(t, *late.Change, replayed, "the slot replays the transaction that committed above it")
}

--- PASS on acf57cf with PG 16 (1.28s), 14 (1.35s) and 18 (7.65s).

Non-blocking

N1. An open stream that has not confirmed the server's send position holds a fast shutdown of the server open, until the stream closes. stream.go:41-47, stream.go:392-403

On shutdown, a logical walsender exits only once the client's reported flush position (its write position, while flush is zero) equals what the walsender has sent. Until then it keeps asking for replies. This stream answers each request with Confirmed, as ST-4 wants, so the walsender never gets there. On PG 16 I opened a stream, let it yield one change, and ran pg_ctl stop -m fast while the caller kept calling Next:

  • Confirming nothing, pg_ctl gave up after 30.3s with server does not shut down. The shutdown checkpoint ran only once the test closed the stream.
  • Confirming each Delivered as it arrived, the server stopped in about 0.05s.

This is latent until the orchestrator exists. It fires when a caller holds a stream open with Confirmed behind Delivered during a server restart: a deferred buffer, a slow flush, or a long cutover wait. The cost is the target server's availability, so the contract belongs in the Stream doc now:

  • confirm Delivered whenever nothing the caller holds is unapplied;
  • Close the stream when it cannot confirm and is not making progress.

One trap for whoever fixes this in the stream rather than in the caller: do not report Delivered as the write position. pglogrepl.SendStandbyStatusUpdate copies a zero flush position from the write position, so before the first Confirm that would move the slot. Today's sendStatus is safe only because write, flush and apply are all Confirmed; setting all three explicitly would make that independent of the library default.

Probe (not a suite test): needs a server the test can stop

Run against a dedicated postgres:16 -c wal_level=logical container, using newSlotFixtureOn with its URL. The probe opens a stream, yields one change, then runs docker exec -u postgres <ctr> pg_ctl stop -m fast -t 30 while looping on stream.Next(ctx, streamWait). In one variant it calls stream.Confirm(d.Delivered) after every delivery.

Without confirms: pg_ctl returned after 30.337654375s: exit status 1: waiting for server to shut down... failed. The server log shows shutting down at :54.894 and the shutdown checkpoint at :25.204, right after the stream closed.

With confirms: the stream ended after 52.98ms, and the server stopped.

N2. A server that ends replication in an orderly way is reported as ErrInvariantViolation. stream.go:163-171

In the confirming variant of N1's probe, the walsender finished with CommandComplete. The stream returned invariant violation: ST-4: stream from slot pgsprite_c4847cdb received *pgproto3.CommandComplete outside of copy-both mode. A server restart with the slot intact is the clean resume ST-4 separates from slot loss. ErrInvariantViolation means state this package cannot have been handed, and an orchestrator that routes on it will treat a routine restart as an engine bug. A distinct sentinel for the server ending the stream (CommandComplete, CopyDone) would let the resume path tell the two apart.

N3. Nine refusals have no test that fails without them. Nine mutants survive the PR's tests. The proposed tests below pass on acf57cf and each kills its mutant:

Refusal removed Where Killed by
change outside a transaction stream.go:290-293 TestStreamRefusesAChangeOutsideATransaction
change to another relation id stream.go:298-301 TestStreamRefusesAChangeToAnotherRelation
quiesced target (DecodesWAL) stream.go:77-79 TestOpenStreamRefusesAQuiescedTargetAndAnotherDatabase
connection on another database (ST-3) stream.go:98-101 same test
the caller's deadline read as a quiet wait stream.go:178-180 TestNextReturnsTheCallersDeadline
Confirm after the stream stopped stream.go:375-377 TestConfirmRefusesAfterTheStreamStopped
a server ERROR mid-stream stream.go:163-165 TestStreamReturnsTheDecodersError
a relation that is not the target relation.go:37-40 TestStreamStopsWhenTheSourceIsRenamed
key not in the replica identity relation.go:52-55 TestStreamStopsWhenTheKeyLeavesTheReplicaIdentity

The deadline mutant matters most. Without ctx.Err() == nil, a caller whose own deadline passes inside the wait gets a progress delivery with a nil error from every later call, so its loop spins instead of stopping.

Test case: nine tests that pass on acf57cf and each kill a surviving mutant

Unit tests (package decode, next to stream_test.go):

// insertMessage encodes a pgoutput INSERT of a one-column row whose only
// column is the text key.
func insertMessage(relationID uint32, key string) []byte {
	b := []byte{byte(pglogrepl.MessageTypeInsert)}
	b = binary.BigEndian.AppendUint32(b, relationID)
	b = append(b, 'N')
	b = binary.BigEndian.AppendUint16(b, 1)
	b = append(b, pglogrepl.TupleDataTypeText)
	b = binary.BigEndian.AppendUint32(b, uint32(len(key)))
	return append(b, key...)
}

// keyOnlyRelation is the stream's record of a one-column table keyed on id.
func keyOnlyRelation(id uint32) *relation {
	return &relation{id: id, columns: []pglogrepl.RelationMessageColumn{{Flags: keyFlag, Name: "id"}}, keyIndex: 0}
}

// A row change the server sent outside BEGIN … COMMIT is refused rather
// than yielded: nothing would ever deliver a position covering it.
func TestStreamRefusesAChangeOutsideATransaction(t *testing.T) {
	s := &Stream{delivered: 100, relation: keyOnlyRelation(7)}

	d, yielded, err := s.handleWALData(walData(110, insertMessage(7, "101")))

	require.ErrorIs(t, err, ErrInvariantViolation)
	assert.False(t, yielded)
	assert.Nil(t, d.Change)
}

// A row change naming a relation other than the one the stream described
// is refused rather than decoded against the target's columns (ST-3).
func TestStreamRefusesAChangeToAnotherRelation(t *testing.T) {
	s := &Stream{delivered: 100, relation: keyOnlyRelation(7)}
	_, _, err := s.handleWALData(walData(105, beginMessage(140)))
	require.NoError(t, err)

	d, yielded, err := s.handleWALData(walData(110, insertMessage(8, "101")))

	require.ErrorIs(t, err, ErrInvariantViolation)
	assert.False(t, yielded)
	assert.Nil(t, d.Change)
}

Integration tests (package decode_test):

// A stream is refused for a quiesced target, whose environment was never
// checked for decoding, and on a replication connection to another
// database, which would read another database's slot of the same name (ST-3).
func TestOpenStreamRefusesAQuiescedTargetAndAnotherDatabase(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)

	_, err := decode.OpenStream(t.Context(), f.cfg, f.mintTarget(t, f.pool, false), slot.ConsistentPoint())
	require.ErrorIs(t, err, decode.ErrInvariantViolation)

	other := dbconn.Config{URL: testutil.NewDatabase(t, f.serverURL)}
	_, err = decode.OpenStream(t.Context(), other, f.target, slot.ConsistentPoint())
	require.ErrorIs(t, err, decode.ErrInvariantViolation)
}

// A source renamed under the stream is no longer the target: its next
// change stops the stream rather than arriving under the old name.
func TestStreamStopsWhenTheSourceIsRenamed(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())

	f.exec(t, `ALTER TABLE %s.ledger RENAME TO ledger_moved`)
	f.exec(t, `INSERT INTO %s.ledger_moved (id, note) VALUES (101, 'renamed')`)

	require.ErrorIs(t, nextError(t, stream), decode.ErrInvariantViolation)
}

// A source whose key is no longer part of the replica identity cannot
// carry its key in an old tuple, so the stream stops at its description.
func TestStreamStopsWhenTheKeyLeavesTheReplicaIdentity(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())

	f.exec(t, `ALTER TABLE %s.ledger REPLICA IDENTITY NOTHING`)
	f.exec(t, `INSERT INTO %s.ledger (id, note) VALUES (101, 'no identity')`)

	require.ErrorIs(t, nextError(t, stream), decode.ErrInvariantViolation)
}

// The caller's own deadline ends the stream even when it falls inside the
// wait: only the wait running out is a quiet table.
func TestNextReturnsTheCallersDeadline(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())

	ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond)
	defer cancel()
	until := time.Now().Add(streamDeadline)
	for time.Now().Before(until) {
		if _, err := stream.Next(ctx, 5*time.Second); err != nil {
			require.ErrorIs(t, err, context.DeadlineExceeded)
			return
		}
	}
	require.FailNow(t, "Next kept yielding progress after the caller's deadline")
}

// A stopped stream refuses to confirm: the caller's position is no longer
// one the stream can vouch for.
func TestConfirmRefusesAfterTheStreamStopped(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())
	f.exec(t, `TRUNCATE %s.ledger`)
	require.ErrorIs(t, nextError(t, stream), decode.ErrUnsupportedChange)

	err := stream.Confirm(t.Context(), stream.Delivered())

	require.ErrorIs(t, err, decode.ErrUnsupportedChange)
	assert.Equal(t, decode.LSN(0), stream.Confirmed())
}

// A publication dropped under the stream fails the decoder on the server;
// the stream returns the server's error rather than a protocol violation.
func TestStreamReturnsTheDecodersError(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())
	_, err := f.pool.Exec(t.Context(), `DROP PUBLICATION `+slot.Name())
	require.NoError(t, err)
	f.exec(t, `INSERT INTO %s.ledger (id, note) VALUES (101, 'unpublished')`)

	var pgErr *pgconn.PgError
	require.True(t, errors.As(nextError(t, stream), &pgErr))
	assert.Equal(t, "42704", pgErr.Code)
}

All nine pass on acf57cf (PG 16). Run against the nine mutants, each one fails exactly the test named in the table. A walsender terminated with pg_terminate_backend does not reach the ErrorResponse case: pgconn returns a FATAL as an error. That is why the last test uses a dropped publication, which the server reports at ERROR.

N4. The "server never moves the slot backwards" assertion reads the slot before the server can have read the status update. feedback_integration_test.go:129-131

The walsender applies a status update when it next reads the socket, which is why the file's own assertConfirmedFlushBecomes polls. Line 131 reads pg_replication_slots straight after the send, so it would pass even if the server did move the slot back. The claim does hold. I confirmed below the slot, called Next ten times (2s), then closed the stream, and the slot stayed at the confirmed position on 14 and 16. The test just needs a round trip before it reads the slot, for example a few Next calls, so that it pins the behavior Confirm's doc relies on.

The verified list and the mutation table are in 2/2.

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

@aparajon

aparajon commented Oct 8, 2026

Copy link
Copy Markdown
Collaborator

🤖 Review of acf57cf (2/2): 0 blocking, 3 non-blocking.

This comment covers the two lenses, OSS adoption and integration ease for SchemaBot and other importers, and then what I verified, including the mutation table. These are not correctness findings.

For adoption, the stream is easy to pick up. Next(ctx, wait) returns either a change or a position, so one loop serves a busy table and a quiet one alike. The first error sticks, so a caller cannot read past a refusal by accident. ErrSourceShapeChanged and ErrUnsupportedChange are sentinels an importer can route on, and server errors arrive wrapped, so errors.As reaches the *pgconn.PgError. The "Decisions to veto" section in the description is the part an outside adopter will want most, and most of it belongs in doc.go, where it will outlive the PR.

The notes below are for the orchestrator that will drive this, and for the progress and resume surfaces an importer such as SchemaBot renders from it. Each is cheaper to shape now, before that caller exists.

1. The stream needs a second WAL sender, but preflight proves one. copy_swap_environment.go:170, stream.go:65-73

The description raises this as a decision to veto, and I would take the preflight change. CreateSlot's connection must stay open while the copier imports its snapshot, and OpenStream dials a walsender of its own. On a server with exactly one free sender, a run would pass preflight, create the slot and start the copy, and only then fail to open the stream, with a shadow built and a slot held. An importer reports that as a mid-run failure, when the whole point of the headroom check is to refuse before anything is created. Requiring two free senders when the route decodes WAL keeps the refusal in preflight, where the operator can act on it.

2. The server's position is read and then dropped, so a lag figure has to be projected. stream.go:213-224

Every keepalive carries ServerWALEnd. Between transactions it becomes Delivered. Inside a transaction, which is the whole of a long catch-up burst, it is discarded. A catch-up progress source wants "how far behind the server is the applier". Today the stream can answer only with its own positions, so an importer would have to query pg_current_wal_lsn() on another connection and subtract, which mixes two clocks. An accessor such as ServerWALEnd(), holding the last value the server reported, would let the copy-and-swap WorkSource report a lag it measured. SchemaBot renders the engine's numbers rather than estimating its own, so a measured figure here is what it would show.

3. A forwarded resume reports the position it asked for, not the one it got. stream.go:122-126, feedback_integration_test.go:134-136

When from is below the slot's confirmed position, the server forwards the start. Start() and Delivered() still report from until the first keepalive or commit. The PR's own test opens such a stream from the consistent point, and the server starts it at the confirmed position. A resume that logs or checkpoints Start() records a position the server never decoded from. Either read the slot's confirmed_flush_lsn in OpenStream and start delivered at the higher of the two, or document Start() as the requested position. With B1's fix, the resume path's rule for replayed events is the other half of this: replay is decided by commit position, so a resume cannot filter by ChangeEvent.LSN.

Verified

  • At acf57cf, the PR's decode tests pass on PG 16: all 23 targeted Stream, OpenStream, Confirm, keepalive, reopen, tuple and relation tests. CI is green across 14 to 18.
  • Measured, and not findings:
    • The interleaved transaction in B1 is replayed from the confirmed position on 14, 16 and 18. The stream loses nothing; the defect is the stated contract.
    • Confirming below the slot's position does not move it backwards on 14 or 16, even after the walsender has read the update (N4 is about the test, not the behavior).
    • A walsender terminated with pg_terminate_backend ends the stream with the server's 57P01 as a *pgconn.PgError, because pgconn returns a FATAL as an error, not as a message.
    • A dropped publication ends the stream with 42704, wrapped.
  • Checked by reading, and not findings:
    • An elapsed wait leaves the connection usable. pgconn's default context watcher sets a read deadline, not a cancel request, and a read cut off mid-message resumes on the next call. A context already done before the receive is still reported as a timeout, so waitElapsed sees it.
    • A keepalive reply before the first Confirm sends write, flush and apply as zero, and the walsender confirms nothing for a zero flush position.
    • For logical decoding the walsender writes the same LSN into XLogData's start and end fields. That is why the ServerWALEnd mutant below is equivalent.
    • IdentifySystem proves the database before START_REPLICATION, and every error path after the dial closes the connection.

Mutants ran against ./pkg/decode/ with the PR's targeted tests on PG 16. 25 of 36 were killed. Of the 11 survivors, 2 are equivalent and the other 9 are killed by N3's tests in 1/2.

Mutant Caught by
keepalive raises Delivered inside a transaction TestKeepalivePositionIsDeliveredOnlyBetweenTransactions
keepalive never raises Delivered the same test, and TestKeepaliveRepliesDoNotMoveTheSlot
status update reports Delivered instead of Confirmed TestKeepaliveRepliesDoNotMoveTheSlot
Confirm regression check removed TestConfirmRefusesARegressionAndAnOverreach
Confirm beyond-delivered check removed the same test
commit does not raise Delivered TestStreamDeliversATransactionsChangesBeforeMovingThePosition and 1 more
commit raises to CommitLSN, not the end LSN TestReopenedStreamReplaysFromTheConfirmedPosition and 1 more
BEGIN inside a transaction / COMMIT outside one accepted TestUnpairedTransactionBoundariesStopTheStream
OldKey set when the key did not move TestStreamUnderReplicaIdentityFullSetsOldKeyOnlyOnAMove
OldKey never set TestStreamDecodesAKeyMoveWithOldKey and 1 more
TRUNCATE ignored TestStreamStopsOnATruncate
zero start accepted TestOpenStreamRefusesTheZeroTargetAndAMissingPosition
an elapsed wait fails the stream TestKeepaliveRepliesDoNotMoveTheSlot
failure not sticky in Next TestStreamStopsWhenTheSourceGainsAColumn
keepalive reply not sent TestKeepaliveRepliesDoNotMoveTheSlot
delivered starts at zero TestReopenedStreamReplaysFromTheConfirmedPosition and 1 more
shape ignores replica identity / compares names only TestRelationSameShapeComparesEveryColumnAttribute
shape change accepted TestStreamStopsWhenTheSourceGainsAColumn
TOAST marker carried as present TestDecodeColumnsDistinguishesNullFromAnUnchangedToastMarker and 1 more
NULL carried as absent the same unit test, and TestStreamDecodesANullAsAPresentColumn
column count not checked TestDecodeColumnsRefusesAColumnCountMismatch
binary data accepted as text TestDecodeColumnsRefusesBinaryData
another publication name passed 14 tests
change outside a transaction, foreign relation id, DecodesWAL, database name, caller deadline, Confirm after failure, server ERROR, target name, key flag survive; N3's nine tests in 1/2 kill them
event LSN taken from ServerWALEnd survives; equivalent (the walsender writes the same LSN into both fields)
key data type not checked survives; equivalent (a NULL or marker key fails ParseInt with the same ErrInvariantViolation)

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: 1 blocking, 7 non-blocking (see the review comments above).

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

Kiran01bm and others added 3 commits October 8, 2026 16:34
* 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
…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.
Base automatically changed from kiran01bm/cs7-slot to main October 8, 2026 05:43
…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.
@Kiran01bm

Copy link
Copy Markdown
Collaborator Author

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

Reviewed head was acf57cf. Every finding below is addressed in the follow-up commit on this branch (the branch is merged forward onto the current kiran01bm/cs7-slot first; the reviewed commits are untouched). Verification is in the refreshed PR body.

# Finding Status Explanation
B1 A change can arrive with an LSN below a position already delivered and confirmed; the Delivered contract, ST-4 Enforced, and CO-5's OldestPending bound are stated in the wrong terms Fixed Both halves of the suggested fix taken. The contract is restated in commit terms everywhere it is written — Delivery, Confirm, doc.go, ST-4's rule and Enforced line, CO-5's OldestPending text, the SAFETY.md row, and the design doc's decode and applier rows — and says outright that a change's own LSN can lie below a delivered or confirmed position and that the server replays by commit position, so a resume must never skip a replayed change for lying below its checkpoint. ChangeEvent now carries Delivered, the stream's position when the change arrived, stamped at the one delivery site; applier.Buffer keys FirstLSN on it instead of the change's LSN, so OldestPending can never fall below Confirmed. The proposed TestStreamDeliversAnInterleavedTransactionBelowTheConfirmedPosition is in the suite, composing the stream with applier.NewBuffer() exactly as the review wrote it: the late change arrives below the confirmed position, and Confirm(min(Delivered, OldestPending)) succeeds. The stream_test.go keepalive comment is reworded so it no longer attributes the interleaving case to the inTransaction guard.
N1 An open stream that has not confirmed the server's send position holds a fast shutdown open until the stream closes Fixed Stream's doc now tells the caller the two obligations: confirm Delivered whenever nothing is unapplied so the walsender can drain, and Close a stream it has stopped reading, since a stuck stream holds the server's shutdown. sendStatus sets write, flush, and apply explicitly rather than relying on the library filling the zeros. Confirm, Confirmed, and sendStatus move to pkg/decode/feedback.go (one file per concern).
N2 A server that ends replication in an orderly way is reported as ErrInvariantViolation Fixed CopyDone / CommandComplete now end the stream with the new ErrStreamEnded sentinel, distinct from a fail-closed violation, so a resume can tell a clean restart from slot loss. Unit-tested; ST-4's Enforced line and the SAFETY.md row name it.
N3 Nine refusals have no test that fails without them Fixed All nine tests added, in the form the review proposed. Unit (stream_test.go): change outside a transaction, change for another relation, and the Delivered stamp on a change. Integration (stream_refusal_integration_test.go): quiesced database and another database refused by IDENTIFY_SYSTEM, renamed source refused as another relation, key column leaving the replica identity, a decoder error surfaced as *pgconn.PgError 42704, and Next returning the caller's own deadline. Each was run against the mutant it targets before landing.
N4 The "server never moves the slot backwards" assertion reads the slot before the server can have read the status update Fixed The fixture gained assertWALSenderFlushBecomes, which polls pg_stat_replication.flush_lsn for the slot's walsender (joined through pg_replication_slots.active_pid) until it shows the position just sent, and TestReopenedStreamReplaysFromTheConfirmedPosition waits on it before reading confirmed_flush_lsn. The assertion now holds on PostgreSQL 16 with the race closed: the server does keep the slot where it was.
L1 The stream needs a second WAL sender, but preflight proves one Fixed Taken. pkg/preflight has copySwapWALSenders = 2 and refuses when fewer than two senders are free, with the refusal's detail naming both holders (the slot's snapshot connection and the decoding stream). The cause's doc comment, docs/engine-role.md, docs/refusal-classes.md, docs/tcb-model.md, and the design doc say two. Unit test extended with the max_wal_senders − 1 case; integration tests run at max_wal_senders=2 (slot headroom, live-sender counting) and max_wal_senders=1 (refused). The "Decisions to veto" bullet that raised this is retired from the body.
L2 The server's position is read and then dropped, so a lag figure has to be projected Fixed Stream.ServerWALEnd() returns the server's position from the last keepalive or WAL message, inside or outside a transaction; Delivered subtracted from it is the stream's lag. Unit-tested.
L3 A forwarded resume reports the position it asked for, not the one it got Fixed (doc); resume read tracked Start()'s doc now says it is the position the stream asked for, that the server forwards a start below the slot's confirmed position, and that Delivered reflects the server's position from the first keepalive or commit. Reading confirmed_flush_lsn before START_REPLICATION belongs to the checkpoint-driven resume that opens a stream without a Slot in hand, which is the next PR in this stack; tracked as an internal follow-up there rather than adding a catalog read this PR's callers cannot yet use.
NB1 Confirm advances s.confirmed before the send, so the field can outrun what the server was told; the same sendStatus error is treated two ways Fixed Confirm sends first and sets confirmed only when the send returns nil, so the field means what its comment says and a retry from a lower durable checkpoint is no longer refused. The second point is resolved toward the keepalive path: a status report the connection cannot carry latches the stream through fail, from Confirm as it already did from a keepalive, since the connection is in the same state either way. Tests: TestConfirmRefusesAfterTheStreamStopped, TestConfirmLeavesTheRecordWhenTheSendFails. Listed under "Decisions to veto" in the body.
NB2 This PR is the first to produce Present: false, making the Value/Present ambiguity live; a Column accessor would make the absent value unreadable Tracked Agreed on the direction. The accessor is a change to the applier's consumption of Column as much as to the type, and the flush path that reads columns lands with the applier's flush PR; the accessor is tracked as an internal follow-up to land there, where its first caller is, rather than as an unused method here. The producer-side guarantee the review verified is unchanged.

Source: adversarial review comments on this PR; fixes in the follow-up commit.

@aparajon

aparajon commented Oct 8, 2026

Copy link
Copy Markdown
Collaborator

🤖 Delta review of 30f16d7 (1/2): 1 blocking, 3 non-blocking.

This is the adversarial correctness pass over what changed since my review at acf57cf, at head 30f16d74ef23d367a3926738bced9173bd4c9189. The branch merged #150's final head in 847d3f9. That merge's second parent has exactly the tree of main at e0e0bc2, so the PR's own diff is now e0e0bc2..30f16d7. Its own content since acf57cf is 30f16d7 plus the merge's hand resolutions: the Buffer, Entry and Batch changes that key FirstLSN on Delivered landed in the merge, not in the commit. I read both. Then I ran pkg/decode, pkg/applier and pkg/preflight at the head on PostgreSQL 16 (all pass), probes on 14, 16 and 18, a live fast shutdown on 16.15, and 26 mutants of the delta.

Every earlier finding is fixed. The commit-ordering rework is right: ChangeEvent.Delivered is stamped at the one delivery site, the buffer keys on it, and Confirm(min(Delivered, OldestPending)) now composes. The new accessor ServerWALEnd does not hold its contract: it is overwritten by every change with that change's own WAL position.

On invariants: ST-4 is extended, and its Enforced text now matches what the stream does. N5 qualifies one sentence of the ErrStreamEnded doc it cites. CO-5 is upheld: the OldestPending bound can no longer fall below Confirmed, but four of the five sites that set it have no test (N6). ST-3 is upheld for the database and target checks. N7 notes that OpenStream does not run the cluster proof #150 added for create and drop.

Earlier finding At 30f16d7 Evidence
B1 Fixed ChangeEvent.Delivered (types.go:69-85) is stamped in deliverChange (stream.go:317-323). Buffer keys FirstLSN on it (buffer.go:102-171). CO-5 and ST-4 are restated in commit terms (invariants.md:211-217, invariants.md:668-692). The interleaving test composes with OldestPending (feedback_integration_test.go:152-192). Mutants D01 and D11 are killed.
N1 Fixed The Stream doc states both obligations (stream.go:60-65). sendStatus sets write, flush and apply explicitly (feedback.go:54-63). Measured on 16.15: a stream that confirms each Delivered let pg_ctl stop -m fast finish, and the stream ended after 60ms.
N2 Fixed CopyDone and CommandComplete return ErrStreamEnded (stream.go:201-204). The same live shutdown returned errors.Is(err, ErrStreamEnded) == true and ErrInvariantViolation == false.
N3 Fixed All nine tests are in stream_test.go and stream_refusal_integration_test.go. I re-ran the nine original mutants at 30f16d7. Each fails exactly its named test.
N4 Fixed The test waits on the walsender's flush_lsn before reading the slot (feedback_integration_test.go:134-137).

Blocking

B2. ServerWALEnd takes every change's own WAL position, so the lag its doc defines underflows or reads zero. stream.go:153-158, stream.go:271, stream_test.go:86-105, SAFETY.md:27, copy-and-swap-design.md:465

The doc says ServerWALEnd is "the server's position as it last reported it", that it "runs ahead of Delivered", and that "Delivered subtracted from it is the stream's lag". In logical replication, though, the walsender does not put its own position in an XLogData message. It writes the position of the record being sent into both the start and the end field. So line 271 sets serverWALEnd to where the change was written. Under B1's ordering, that lies below Delivered whenever the change's transaction was still open when the stream last delivered a position. I measured it on 14, 16 and 18 with the B1 interleaving:

after the first commit:   ServerWALEnd = 0/19278E0   Delivered = 0/19278E0
after the late change:    ServerWALEnd = 0/1927790   Delivered = 0/19278E0
ServerWALEnd() - Delivered() = 18446744073709551280

That is the doc's formula on LSN, a uint64. The same overwrite makes a commit read as no lag. A commit's XLogData carries the commit's end, which is exactly the new Delivered, however far the server is ahead. The PR body says lag "is readable during a long catch-up burst". In such a burst nearly every message is a change, so the keepalive's value lasts only until the next change.

This is blocking because it breaks the documented contract of a new exported accessor, and SAFETY.md and the design doc repeat that contract. The unit test pins the wrong model: line 105 feeds a change whose end field (210) differs from its start (170), which a server does not send. D07, the mutant that deletes line 271, is killed only by that assertion. Nothing in the tree reads ServerWALEnd yet. Its first readers will be a catch-up progress source and a throttle, and the value is wrong in both directions for them: the underflow reads as an enormous lag, and the zero after each commit reads as caught up.

Suggested fix: take the position from keepalives only, and never report it below Delivered:

// ServerWALEnd is the walsender's position from its last keepalive, and
// never less than Delivered.
func (s *Stream) ServerWALEnd() LSN { return max(s.serverWALEnd, s.delivered) }

Then delete line 271, and say in the doc, SAFETY.md and the design row what the value is. 2/2 has the caveat an importer needs about what that figure measures. TestServerWALEndFollowsEveryServerReport fails with the fix at lines 88 and 105, as it should. The tests below replace it.

Tests that fail on 30f16d7 and pass with the fix

Unit (package decode, next to stream_test.go):

// For logical decoding the walsender writes the record's own position into
// XLogData's end field, so a change reports where it was written, not where
// the server is: only a keepalive moves ServerWALEnd, and it never reads
// below Delivered.
func TestServerWALEndIsNotMovedByAChange(t *testing.T) {
	s := &Stream{delivered: 100, relation: keyOnlyRelation(7)}
	_, _, err := s.handleKeepalive(t.Context(), pglogrepl.PrimaryKeepaliveMessage{ServerWALEnd: 200})
	require.NoError(t, err)

	_, _, err = s.handleWALData(walData(160, beginMessage(240)))
	require.NoError(t, err)
	_, _, err = s.handleWALData(walData(170, insertMessage(7, "101")))
	require.NoError(t, err)
	assert.Equal(t, LSN(200), s.ServerWALEnd())

	_, _, err = s.handleWALData(walData(250, commitMessage(240, 250)))
	require.NoError(t, err)
	assert.Equal(t, LSN(250), s.ServerWALEnd(), "a commit past the last keepalive is the furthest position known")
}

Integration (package decode_test):

// A transaction that wrote before another committed arrives after that
// commit. Its change carries the WAL position it was written at, which lies
// below what the stream already delivered, so the server's position the
// stream reports for lag must not follow it there.
func TestServerWALEndNeverTrailsDelivered(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())

	open, err := f.pool.Begin(t.Context())
	require.NoError(t, err)
	t.Cleanup(func() {
		if err := open.Rollback(context.WithoutCancel(t.Context())); err != nil && !errors.Is(err, pgx.ErrTxClosed) {
			t.Logf("roll back the open transaction: %v", err)
		}
	})
	_, err = open.Exec(t.Context(), fmt.Sprintf(`INSERT INTO %s.ledger (id, note) VALUES (201, 'written first')`, f.schema))
	require.NoError(t, err)
	f.exec(t, `INSERT INTO %s.ledger (id, note) VALUES (202, 'committed first')`)
	first := nextChange(t, stream)
	awaitDelivered(t, stream, first.LSN+1)
	assert.GreaterOrEqual(t, stream.ServerWALEnd(), stream.Delivered(), "after a commit")

	require.NoError(t, open.Commit(t.Context()))
	late := nextChangeDelivery(t, stream)
	require.EqualValues(t, 201, late.Change.Key)
	assert.GreaterOrEqual(t, stream.ServerWALEnd(), stream.Delivered(), "after a change written below Delivered")
}

On 30f16d7: --- FAIL: TestServerWALEndIsNotMovedByAChange (0.00s) (expected: 0xc8, actual: 0xaa), and --- FAIL: TestServerWALEndNeverTrailsDelivered (1.61s) with "0/1927790" is not greater than or equal to "0/19278E0" on PG 16. The integration test fails the same way on 14 (0/17351A0 vs 0/17352F0) and on 18 (0/1BAA7B0 vs 0/1BAA900). With the fix applied, both tests pass, and the whole decode package passes except the two assertions of TestServerWALEndFollowsEveryServerReport named above.

Non-blocking

N5. On PostgreSQL 16, the slot does not keep the last confirmed position across the clean shutdown that ErrStreamEnded reports. stream.go:24-29

The doc says "The slot keeps the position last confirmed". In the live shutdown above, the stream confirmed each Delivered and ended with ErrStreamEnded at 0/2193720. The walsender finishes a fast shutdown only once the client's reported flush equals what it sent, so the server had received that position. After a restart, pg_replication_slots showed confirmed_flush_lsn = 0/2193668 for the slot, which is below the last confirm. Before 17, a confirm that only moves confirmed_flush does not dirty the slot, so the shutdown checkpoint does not write it. PostgreSQL 17 writes logical slots at the shutdown checkpoint. I measured 16 only.

Nothing is lost. A stream reopened from the caller's own checkpoint is honored, since the server forwards only a start that lies below the slot. A resume can still misread it, though. L3's follow-up reads confirmed_flush_lsn before START_REPLICATION, and after a clean restart a slot behind the checkpoint is normal; it is not evidence that the slot was replaced. Suggested wording: the slot keeps a position at or below the last confirm, and a resume starts from its own checkpoint.

N6. Four of the five places that key FirstLSN on Delivered have no test that fails without them. buffer.go:116, buffer.go:151, buffer.go:164, buffer.go:171

TestBufferOldestPendingIsTheDeliveredPositionNotTheChangesLSN covers the INSERT. Reverting any of the other four sites to ev.LSN survives the whole decode and applier suites: an UPDATE of a key the buffer does not hold, a DELETE, a key move's image, and its old-key marker (D12, D13, D14, D17). Each of those mutants brings back B1 for that event kind.

Test case: passes on 30f16d7, and each of the four mutants fails it
// Every kind of event keys its entry on the delivered position it arrived
// with, not on its own LSN: an UPDATE of a key the buffer does not hold, a
// DELETE, and a key move, whose moved image and old-key marker both owe the
// stream the same position.
func TestBufferKeysEveryEventKindOnItsDeliveredPosition(t *testing.T) {
	oldKey := int64(1)
	for name, ev := range map[string]decode.ChangeEvent{
		"update":   {Kind: decode.Update, LSN: 50, Delivered: 80, Key: 2, Columns: []decode.Column{col("label", "a")}},
		"delete":   {Kind: decode.Delete, LSN: 50, Delivered: 80, Key: 2},
		"key move": {Kind: decode.Update, LSN: 50, Delivered: 80, Key: 2, OldKey: &oldKey, Columns: []decode.Column{col("label", "a")}},
	} {
		t.Run(name, func(t *testing.T) {
			b := NewBuffer()
			require.NoError(t, b.Add(ev))

			oldest, ok := b.OldestPending()
			require.True(t, ok)
			assert.Equal(t, decode.LSN(80), oldest)
			for key, e := range b.entries {
				assert.Equal(t, decode.LSN(80), e.FirstLSN, "entry for key %d", key)
			}
		})
	}
}

--- PASS on 30f16d7. Against the mutants: line 116 fails .../update, line 171 fails .../delete, and lines 151 and 164 each fail .../key_move.

N7. OpenStream is the one replication entry point that does not run the cluster proof. stream.go:112-122

#150 made CreateSlot and DropSlot prove that the replication connection and the pool are on one cluster (server_identity.go:28). OpenStream checks only IDENTIFY_SYSTEM's database name, and it takes no pool to prove against. The derived name hashes the database, schema and table, so the same table on two clusters gets the same slot name. Examples are the shards of one application, or a staging copy and a production copy. A cfg that points at the other cluster opens that cluster's slot, and its rows arrive as this target's changes. newRelation compares names, not the relation OID, so nothing downstream notices.

This is latent: today only tests open streams. It fires when the resume path opens a stream from a config that CreateSlot never proved in that process. Suggested fix: have OpenStream take the pool and call proveSameServer before START_REPLICATION. Whatever closes #150's standby case then covers this path too. Comparing RelationID with the target's OID would be a cheap second check, but it needs an accessor on CopySwapTarget. I did not run this one, since it needs two clusters with the same database name.

The adoption lens, the verified list and the mutation table are in 2/2.

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

@aparajon

aparajon commented Oct 8, 2026

Copy link
Copy Markdown
Collaborator

🤖 Delta review of 30f16d7 (2/2): 0 blocking, 1 non-blocking.

This comment covers the two lenses, OSS adoption and integration ease for SchemaBot and other importers, for what changed since acf57cf. Then it lists what I verified, including the mutation table. These are not correctness findings.

For adoption, the delta makes the stream easier to build on. ChangeEvent.Delivered gives an importer the confirmable position on the event itself, so a buffer never has to reason about commit order. ErrStreamEnded is a sentinel a resume loop can route on. The Stream doc now spells out the two caller obligations that keep a server's shutdown from hanging. The earlier lens notes:

Earlier note At 30f16d7
1. two WAL senders Taken. copySwapWALSenders = 2 (copy_swap_environment.go:31-35, :175-181). Both mutants (D15, D16) are killed by the unit test and by the max_wal_senders 1 and 2 integration tests.
2. a measured lag Taken, but not yet right. See B2 in 1/2 and A1 below.
3. forwarded resume Doc fixed (stream.go:143-147). The catalog read moves to the resume PR, where its caller is. That is reasonable.
Confirm records before the send Fixed (feedback.go:28-47). D02 and D03 are killed by TestConfirmLeavesTheRecordWhenTheSendFails.
Column accessor Deferred to the flush PR, where its first caller lands. That is reasonable.

Non-blocking

A1. Even with B2 fixed, ServerWALEnd is how far the walsender has decoded, not the server's WAL end, so a lag from it is a lower bound. stream.go:153-158

My earlier note 2 called the keepalive value the server's position. That was too strong. A logical walsender's keepalive carries its send position (sentPtr in walsender.c), which is the end of the last record it decoded. While it works through a backlog, most of the backlog is WAL it has not read yet, so ServerWALEnd - Delivered stays small exactly when the stream is furthest behind. I read this in the server source and did not measure it.

A progress surface that renders "lag" from this figure will show near zero during the catch-up an operator is waiting on. Two changes would let an importer show a number it can stand behind:

  • In the doc, name the value as the walsender's send position, a lower bound on lag.
  • Leave true lag to a measurement against pg_current_wal_lsn() on the pool, or have the stream's progress source take that measurement next to it.

Verified

  • The PR's own change is isolated. 847d3f9's second parent f13dc12 has the same tree as main at e0e0bc2. So everything decode: create, drop, and inspect the route's replication slot #150 brought in is upstream, and the PR's diff is e0e0bc2..30f16d7. The merge's hand resolutions carry the Buffer, Entry and Batch doc changes and the hold-test fixtures. 30f16d7 carries the rest.
  • Tests at the head. decode (60.6s), applier (41.6s) and preflight (82.4s) pass at 30f16d7 on PG 16.
  • The applier's other bound already uses commit terms. Buffer.Release(passed) compares a read's insert position against a delivered position, which already orders by commit. No other applier code reads ChangeEvent.LSN.
  • The boundary is honored. A Delivered raised by a keepalive between transactions is a record boundary that no unsent commit record starts below. So confirming exactly a change's Delivered still replays its transaction. My first B2 probe confirmed the late change's Delivered, reopened a stream from it, and saw key 201 again.
  • Live fast shutdown on 16.15. With per-delivery confirms, the stream ended after 60ms with ErrStreamEnded. N5's slot read was taken on the same server after the restart.
  • A test failure I could not reproduce. One run of the D12 and D13 mutants overlapped my probe runs. In that run, TestStreamDecodesADelete, TestInspectSlotMeasuresTheWALTheSlotRetains and a flusher hold test also failed. The re-runs of both mutants, four isolated runs of the two decode tests, and two concurrent full runs of decode and applier all passed. I did not keep that output, so I am not reporting it as a finding.
Mutant (delta) Caught by
D01 change not stamped with Delivered TestStreamDeliversAnInterleavedTransactionBelowTheConfirmedPosition, TestStreamStampsAChangeWithTheDeliveredPositionItArrivedWith
D02 Confirm records before the send TestConfirmLeavesTheRecordWhenTheSendFails
D03 a failed Confirm send does not latch TestConfirmLeavesTheRecordWhenTheSendFails
D04 status carries apply only TestConfirmMovesTheSlotToTheConfirmedPosition and three others
D05 orderly end reported as a violation TestStreamEndsWhenTheServerLeavesCopyBothMode
D06 keepalive does not record its position TestKeepalivePositionIsDeliveredOnlyBetweenTransactions and two others
D07 change does not overwrite ServerWALEnd only TestServerWALEndFollowsEveryServerReport:105, which pins the wrong model (B2)
D08 keepalive reply carries Delivered TestKeepaliveRepliesDoNotMoveTheSlot
D09, D10 Confirm overreach / regression check removed TestConfirmRefusesARegressionAndAnOverreach
D11 every FirstLSN keyed on the change LSN TestBufferOldestPendingIsTheDeliveredPositionNotTheChangesLSN, the interleaving test
D12, D13, D14, D17 one FirstLSN site keyed on the change LSN (update, delete, move image, old-key marker) survive; N6's test kills each
D15, D16 one free WAL sender accepted three preflight tests
The nine N3 refusals removed each fails exactly the test N3 named

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: 1 blocking, 4 non-blocking (see the review comments above).

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

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.
@Kiran01bm

Copy link
Copy Markdown
Collaborator Author

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

Delta review of 30f16d7 (1/2 and 2/2): B2 fixed; N5, N6, N7, and A1 taken.

# Finding Status Explanation
B2 ServerWALEnd takes every change's own WAL position, so the lag its doc defines underflows or reads zero Fixed handleWALData no longer touches serverWALEnd; only handleKeepalive moves it, and ServerWALEnd() returns max(serverWALEnd, delivered). The accessor doc, the field comment, the handleWALData doc, the SAFETY.md decode row, and the design-doc decode row now say what the value is: the walsender's send position from its last keepalive, never below Delivered, a lower bound on lag. TestServerWALEndFollowsEveryServerReport is replaced by TestServerWALEndIsNotMovedByAChange (unit: a change below the keepalive leaves it, a keepalive inside a transaction moves it while Delivered holds, a commit past the keepalive is the furthest position known, and it reads Delivered before the server reports) and TestServerWALEndNeverTrailsDelivered (integration: the interleaved-transaction case, with require.Less(late.Change.LSN, Delivered) so the test proves the change really arrived below Delivered). Restoring the overwrite fails both, the integration one with the same 0/1927790 < 0/19278E0 pair the review measured.
N5 ErrStreamEnded doc says the slot keeps the last confirmed position; on 16 a confirm-only move is not written at the shutdown checkpoint Fixed Reworded: the slot survives with a position at or below the last confirm (a confirm that only moves confirmed_flush is not always written out before a shutdown), and a new stream from the caller's own checkpoint continues, since the server forwards a start below the slot's position. The same caveat is in the PR body's ErrStreamEnded decision, and the resume PR that reads confirmed_flush_lsn before START_REPLICATION is on notice not to read a slot behind the checkpoint as a replaced slot.
N6 Four of the five FirstLSN sites keyed on Delivered have no test that fails without them Fixed TestBufferKeysEveryEventKindOnItsDeliveredPosition added as suggested, with a require.NotEmpty(b.entries) guard so the per-entry loop cannot pass vacuously. Reverting each of the four sites to ev.LSN fails its subtest: buffer.go update → update, delete marker → delete, move image and old-key marker → key move.
N7 OpenStream is the one replication entry point that does not run the cluster proof Fixed OpenStream takes the pool (in the position CreateSlot and DropSlot use) and calls proveSameServer before START_REPLICATION; proveSameServer now returns the proven identity, so the target-database check reuses the same IDENTIFY_SYSTEM round trip instead of adding one. TestOpenStreamRefusesAReplicationConnectionOnAnotherServer uses the existing second-cluster fixture: a replication config pointing at a database of the target's name on another cluster is refused with ErrInvariantViolation; without the proof it falls through to that cluster's 42704 replication slot does not exist, which the test does not accept. The RelationID-versus-target-OID check is tracked as an internal follow-up; it needs an accessor on CopySwapTarget. The signature change reaches only test callers and the stacked catch-up PR, which will follow.
A1 Even with B2 fixed, ServerWALEnd is how far the walsender has decoded, not the server's WAL end, so a lag from it is a lower bound Taken Both suggested changes are in: the accessor doc and the design row name the value as the walsender's send position, a lower bound on lag that stays small while a backlog is undecoded, and say that lag against the server's WAL end is measured on the pool against pg_current_wal_lsn(). The stacked catch-up PR's lag_bytes currently derives from ServerWALEnd − Confirmed; it will take the pool-side measurement, and its doc will say which figure it reports, as part of following this change.

Source: delta review comments 1/2 and 2/2 on head 30f16d7; fixes in the follow-up commit.

@Kiran01bm
Kiran01bm enabled auto-merge (squash) October 8, 2026 08:06
@aparajon

aparajon commented Oct 8, 2026

Copy link
Copy Markdown
Collaborator

🤖 Delta review of 8bcae74: 1 blocking, 2 non-blocking.

This covers what changed since my review at 30f16d7. The head is now 8bcae74fc9aaf73e813b6c8fac3a9edcca101e06. c667bb2 merges main at e0e0bc2, and its tree is identical to 30f16d7's, so the PR's own change is exactly 8bcae74. I read it whole. I mutation-tested its production lines and re-ran the four N6 buffer mutants. I ran pkg/decode with -race on PostgreSQL 16 and 18. I also read the red test (PostgreSQL 18) job and reproduced it.

8bcae74 fixes everything my last review asked for. The red PostgreSQL 18 job is a separate problem. It comes from a test that 30f16d7 added, and that test is right to fail: on 18, a dropped publication makes the stream silently skip a change and then deliver a position past it (B3). My last delta ran the package on 16 only, which is how I missed it.

Invariants: ST-4 is weakened on PostgreSQL 18 until B3 is fixed: Delivered passes a change the stream never yielded. ST-3 is extended: OpenStream now runs the cluster proof. Its Enforced line does not say so yet (N9).

Earlier finding At 8bcae74 Evidence
B2 Fixed handleWALData no longer writes serverWALEnd (stream.go:281-287). ServerWALEnd returns max(serverWALEnd, delivered) (stream.go:162-172). The interleaving test is in (feedback_integration_test.go:194-225) and passes on 16 and 18. Removing the floor (E2) fails both new tests. Putting the overwrite back (E1) fails neither; see N8.
N5 Fixed The doc now says the slot keeps a position "at or below the last confirm" and that a resume starts from the caller's own checkpoint (stream.go:25-33).
N6 Fixed buffer_test.go:366-390. I reverted each of the four sites to ev.LSN. Line 116 fails /update, 151 and 164 fail /key_move, and 171 fails /delete.
N7 Fixed OpenStream takes the pool and calls proveSameServer before START_REPLICATION (stream.go:121-131). Replacing the proof with the old database-name check (E3) fails TestOpenStreamRefusesAReplicationConnectionOnAnotherServer (stream_refusal_integration_test.go:33-48). Dropping the system-id comparison (E5) fails that test and #150's create-slot test. OpenStream has no caller outside tests.
A1 Taken The accessor doc, SAFETY.md:27 and copy-and-swap-design.md:465 call the figure the walsender's send position and a lower bound on lag. They point true lag at pg_current_wal_lsn() on the pool.

Blocking

B3. On PostgreSQL 18, a dropped publication makes the stream skip a change and deliver a position past it, and TestStreamReturnsTheDecodersError times out. stream.go:219-221, stream_refusal_integration_test.go:76-87

The red job (log) fails TestStreamReturnsTheDecodersError with the stream did not fail before the stream deadline. The failure is deterministic and has nothing to do with 8bcae74:

  • It fails on PostgreSQL 18.6 locally, both at 8bcae74 and at 30f16d7.
  • It passes on 16.
  • The stacked branch's run at 30f16d7 code failed it the same way on 18 (log).
  • 14 through 17 pass in CI.
  • 30f16d7 added the test, and no CI run on this PR covered that commit.

Rerunning the job will not clear it.

The failing test exposes a real fail-open. Through 17, pgoutput fails with 42704 when its publication is gone. Version 18 skips a publication it cannot find instead. The server sends a WARNING, and handleMessage drops every notice. I probed 18.6 by dropping the publication and inserting one row. The stream logged this:

NOTICE WARNING 55000: skipped loading publication "pgsprite_b1d552b3"
insert committed below 0/1BAAB88; Delivered=0/1BAAB88 changes=0 err=<nil>
Confirm(Delivered) err=<nil>; slot confirmed_flush=0/1BAAB88

So the change never arrives, a keepalive moves Delivered past its commit, and Confirm moves the slot past it as well. The slot will never resend that change. A copy-and-swap would then swap in a table that lacks the write. Any DROP PUBLICATION on the route's publication leads here: an operator cleaning up, DDL tooling, or another process's DropSlot. This is blocking for two reasons: it breaks ST-4 on a supported major, and the all-green check is red.

The job's other failure was TestReopenedStreamReplaysFromTheConfirmedPosition, at feedback_integration_test.go:119: the third change did not arrive within 15s. I could not reproduce it on 18 with -race: one full package run plus six looped runs of that test all passed. The stacked branch's run passed it too, and 8bcae74 does not touch how changes are delivered. I am not counting it as a finding.

Suggested fix: fail closed on a warning in copy-both mode. The walsender sends one only to say it will not send something. In my probe it said so once, when it reloaded publications. That makes the warning the stream's only chance to stop before a keepalive moves Delivered past the skipped change. Then make the test expect the SQLSTATE the server's major reports. Full pkg/decode passes with -race on 16 and 18 with this change, so no other test receives a warning:

	case *pgproto3.NoticeResponse:
		// INV: ST-4 — a warning in copy-both mode is the walsender saying it
		// will not send something, such as changes under a publication it
		// skipped loading; stopping here keeps a later keepalive from
		// delivering a position past them.
		if msg.SeverityUnlocalized == "WARNING" {
			pgErr := pgconn.ErrorResponseToPgError((*pgproto3.ErrorResponse)(msg))
			return Delivery{}, false, fmt.Errorf("%w: ST-4: stream from slot %s: %w", ErrInvariantViolation, s.slotName, pgErr)
		}
		return Delivery{}, false, nil
	case *pgproto3.ParameterStatus:
		// An asynchronous message the connection has already recorded.
		return Delivery{}, false, nil

If failing on every warning is too broad, matching SQLSTATE 55000 on the notice covers the case I measured.

Test case: fails on 8bcae74 on PostgreSQL 18, passes with the fix on 16 and 18
// A publication dropped under the stream ends it with the server's
// SQLSTATE: before 18 the decoder fails on the missing publication
// (42704); from 18 it skips loading it with a warning (55000) and sends
// nothing for the change, which the stream must refuse rather than step past.
func TestStreamReturnsTheDecodersError(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	stream := f.openStream(t, slot.ConsistentPoint())
	_, err := f.pool.Exec(t.Context(), `DROP PUBLICATION `+slot.Name())
	require.NoError(t, err)
	f.exec(t, `INSERT INTO %s.ledger (id, note) VALUES (101, 'unpublished')`)

	var version int
	require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT current_setting('server_version_num')::int`).Scan(&version))
	want := "42704"
	if version >= 180000 {
		want = "55000"
	}
	requireSQLState(t, nextError(t, stream), want)
}

On 8bcae74, PostgreSQL 18.6: --- FAIL: TestStreamReturnsTheDecodersError (16.22s) (the stream did not fail before the stream deadline). With the fix it passes on 18 (1.57s) and on 16 (1.38s).

Non-blocking

N8. Two of 8bcae74's production lines survive the suite: the Delivered floor hides the overwrite B2 removed, and the target-database check is never reached. stream.go:172, stream.go:128-131

The reply says putting the overwrite back fails both new tests. At this head it fails neither (E1). Every change in those tests carries a position at or below Delivered, so max hides the regression. It shows only when a change follows a keepalive sent mid-transaction, which is the case A1 describes. A keepalive sent during a long replay can carry a send position above the changes still to come. With the overwrite back, ServerWALEnd drops from the keepalive's value to the change's.

Separately, proveSameServer already makes the connection agree with the pool, so the identity.database != target.Database() check runs only when the pool is on a different database from the target. No test covers that (E4). Without the check, the server still refuses with 55000 replication slot ... was not created in this database. The stream's own ST-3 error is lost, though.

Test cases: pass on 8bcae74, and each fails its mutant

Unit (package decode):

// A keepalive sent while a transaction is being replayed carries a position
// above the changes the walsender has yet to send from it, so a change that
// follows the keepalive leaves ServerWALEnd where the keepalive put it,
// even while that position is still above Delivered.
func TestServerWALEndKeepsAKeepaliveAboveALaterChange(t *testing.T) {
	s := &Stream{delivered: 100, relation: keyOnlyRelation(7)}
	_, _, err := s.handleWALData(walData(160, beginMessage(240)))
	require.NoError(t, err)
	_, _, err = s.handleKeepalive(t.Context(), pglogrepl.PrimaryKeepaliveMessage{ServerWALEnd: 230})
	require.NoError(t, err)
	require.Equal(t, LSN(100), s.Delivered(), "a keepalive inside a transaction is not delivered")

	_, _, err = s.handleWALData(walData(210, insertMessage(7, "101")))
	require.NoError(t, err)
	assert.Equal(t, LSN(230), s.ServerWALEnd())
}

Integration (package decode_test):

// The cluster proof ties the replication connection to the pool, not to the
// target: a connection and a pool that agree on another database of the
// same cluster are refused as the stream's own ST-3 violation, before the
// server is asked for a slot that database does not hold.
func TestOpenStreamRefusesAPoolOnAnotherDatabase(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)

	other := dbconn.Config{URL: testutil.NewDatabase(t, f.serverURL)}
	otherPool, err := dbconn.NewPool(t.Context(), other)
	require.NoError(t, err)
	t.Cleanup(otherPool.Close)

	stream, err := decode.OpenStream(t.Context(), other, otherPool, f.target, slot.ConsistentPoint())
	require.ErrorIs(t, err, decode.ErrInvariantViolation)
	assert.Nil(t, stream)
}

Both pass on 8bcae74. With the overwrite restored: --- FAIL: TestServerWALEndKeepsAKeepaliveAboveALaterChange (0.00s) (expected: 0xe6, actual: 0xd2). With the target-database check removed: --- FAIL: TestOpenStreamRefusesAPoolOnAnotherDatabase (1.45s), whose error chain holds only the server's SQLSTATE 55000.

N9. ST-3's Enforced line still names only CreateSlot and DropSlot as running the cluster proof. invariants.md:657-665

OpenStream cites ST-3 at the call (stream.go:121-123). SAFETY.md and the design row already describe the new proof. Adding OpenStream to the Enforced line keeps the registry in step with them.

Mutant (8bcae74) Caught by
E1 a change overwrites serverWALEnd again survives; N8's unit test kills it
E2 ServerWALEnd without the Delivered floor TestServerWALEndIsNotMovedByAChange, TestServerWALEndNeverTrailsDelivered
E3 OpenStream checks the database name only TestOpenStreamRefusesAReplicationConnectionOnAnotherServer
E4 OpenStream drops the target-database check survives; N8's integration test kills it
E5 the proof skips the system-id comparison TestOpenStreamRefusesAReplicationConnectionOnAnotherServer, TestCreateSlotRefusesAReplicationConnectionOnAnotherServer
E6 the proof returns the replication side's identity equivalent: both identities are equal once the checks pass

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: 1 blocking, 2 non-blocking (see the review comment above).

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

From PostgreSQL 18 a walsender that cannot load the slot's publication
skips it with a WARNING instead of failing, and sends nothing for the
changes under it. The stream dropped every notice, so a later keepalive
moved Delivered -- and the caller's Confirm the slot -- past a change the
slot will never resend. A warning in copy-both mode is the walsender
saying it will withhold something, so handleNotice now ends the stream
with ErrInvariantViolation and the server's SQLSTATE reachable; a notice
below warning is still informational. The decoder-error test expects the
server's SQLSTATE for its major: 42704 before 18, 55000 from 18.

Two surviving mutants get tests: a change that follows a keepalive sent
mid-transaction leaves ServerWALEnd at the keepalive's position, and a
pool on another database of the same cluster is refused by OpenStream's
own check rather than the server's missing-slot error. ST-3's enforced
line names OpenStream among the entry points that run the cluster proof;
the SAFETY.md and design decode rows and ST-4 record the warning rule.
@Kiran01bm

Copy link
Copy Markdown
Collaborator Author

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

Delta review of 8bcae74: B3 fixed; N8 and N9 taken.

# Finding Status Explanation
B3 On PostgreSQL 18 a dropped publication makes the walsender skip it with a WARNING and send nothing; the stream dropped every notice, so a keepalive moved Delivered — and Confirm the slot — past a change the slot will never resend, and TestStreamReturnsTheDecodersError timed out on 18 Fixed handleMessage now routes NoticeResponse to handleNotice: a notice below warning stays informational; a WARNING ends the stream with ErrInvariantViolation (ST-4) and the server's SQLSTATE reachable through pgconn.ErrorResponseToPgError. Took the broad form — any warning, not 55000 alone — since the walsender sends one only to say it will not send something, and that is the stream's one chance to stop before a keepalive moves Delivered. TestStreamReturnsTheDecodersError expects 42704 before 18 and ErrInvariantViolation carrying 55000 from 18: on 18.6 it passes in 1.7 s with the handling and times out at 16.9 s without it. TestStreamStopsOnAWarningFromTheServer (unit) covers both branches on every major — NOTICE leaves the stream running and Delivered where it was, WARNING stops it with 55000 reachable — and fails with the handling removed. SAFETY.md, the design-doc decode row, and ST-4's Enforced line record the rule. Full pkg/decode with -race is green on 18 and on 16, so no other test receives a warning.
N8 Two production lines survive the suite: the Delivered floor hides the overwrite B2 removed, and the target-database check in OpenStream is never reached Taken Both suggested tests are in as written. TestServerWALEndKeepsAKeepaliveAboveALaterChange: a change that follows a keepalive sent mid-transaction leaves ServerWALEnd at the keepalive's 230; with the overwrite restored it reads the change's 210 (expected 0xe6, actual 0xd2). TestOpenStreamRefusesAPoolOnAnotherDatabase: a config and pool that agree on another database of the same cluster are refused as the stream's own ST-3 violation; with the check removed the chain holds only the server's 55000. The claim in the last response that restoring the overwrite failed the two B2 tests was wrong at that head; it fails this one.
N9 ST-3's Enforced line names only CreateSlot and DropSlot as running the cluster proof Taken The line now reads CreateSlot, DropSlot, and OpenStream prove the connection and the pool are sessions of one database on one cluster before any slot command or START_REPLICATION.

Not a finding, noted: TestReopenedStreamReplaysFromTheConfirmedPosition was red once on the 18 job and did not reproduce here either; it is unchanged by this commit.

Source: delta review comment on head 8bcae74; fixes in the follow-up commit.

@aparajon

aparajon commented Oct 8, 2026

Copy link
Copy Markdown
Collaborator

🤖 Delta review of e1fb311: 1 blocking, 2 non-blocking.

This covers e1fb311, the only commit since my review at 8bcae74. The head is now e1fb311d096e20c9b657f9ac4b9afc9087a3a90c. I read the commit whole and mutation-tested its production change. I re-ran the two N8 mutants and ran pkg/decode with -race on PostgreSQL 14, 16 and 18, which passed on all three. I also read the red test (PostgreSQL 14) job and traced its cause.

N8 and N9 are fixed. B3 is partly fixed. The stream now stops on the warning, but only if the warning reaches it. If client_min_messages = error is set on the role, the database or the server, the walsender never sends the warning. On 18 the change is then skipped and confirmed past, exactly as in B3 (B4). The red PostgreSQL 14 job has nothing to do with this PR (N10).

Invariants: ST-4 is extended by the warning rule, but it stays weakened on 18 for any session that filters warnings, until B4 is fixed. ST-3 is upheld, and its Enforced line now names OpenStream.

Earlier finding At e1fb311 Evidence
B3 Partly fixed handleNotice ends the stream on a WARNING with ErrInvariantViolation, and the server's *pgconn.PgError stays in the chain (stream.go:219-220, stream.go:232-248). After the failure, Confirm refuses (feedback.go:29-31). The test expects 42704 before 18 and 55000 from 18 (stream_refusal_integration_test.go:93-117). It passes on 14, 16 and 18 here, and on 14 through 18 in CI. One gap remains: the fix depends on a session setting the stream does not control (B4).
N8 Fixed Both tests are in (stream_test.go:199-215, stream_refusal_integration_test.go:49-65). E1 now fails TestServerWALEndKeepsAKeepaliveAboveALaterChange, and E4 now fails TestOpenStreamRefusesAPoolOnAnotherDatabase.
N9 Fixed invariants.md:660-662 names CreateSlot, DropSlot and OpenStream.

Does failing on every WARNING over-trigger? I don't see it. No healthy stream in the full suite received one on 14, 16 or 18. A NOTICE or anything lower still passes through, and the mutant that fails on every notice (F2) is caught. So taking the broad form looks right to me.

Blocking

B4. A role, database or server that sets client_min_messages = error keeps the warning from the stream, and on 18 B3's data loss comes back. replication.go:36, stream.go:239-241

The walsender filters its messages through client_min_messages like any other backend, and ConnectReplication does not set it. A role or a database set to error to quiet noisy clients passes that setting on to the replication connection, so the WARNING 55000 is never sent and handleNotice never runs. I probed this on 18.6 at e1fb311. I ran ALTER DATABASE … SET client_min_messages = error, then the same drop-and-insert as before:

insert committed below 0/1BAAD58; Delivered=0/1BAAD58 changes=0 err=<nil>; Confirm(Delivered) err=<nil>
slot confirmed_flush=0/1BAAD58

So the change never arrives, Delivered moves past its commit, and Confirm moves the slot past it too. That is the B3 outcome. Before 18 the server sends an ERROR, and client_min_messages cannot filter that out, so 14 through 17 still fail closed.

Suggested fix: have the replication connection request warnings itself. A startup parameter overrides role, database and server settings, and it also overrides a value in the URL, because it is assigned after parsing:

	connConfig.RuntimeParams["replication"] = "database"
	// A walsender reports what it will withhold from the stream as a
	// warning, so no role, database, or server setting may keep warnings
	// from reaching the connection.
	connConfig.RuntimeParams["client_min_messages"] = "warning"

With this change, the probe stops with invariant violation: ST-4: … WARNING: skipped loading publication "pgsprite_8b26b557" (SQLSTATE 55000). pkg/dbconn and pkg/decode pass with -race on 18.

Test case: fails on e1fb311 on PostgreSQL 18, passes with the fix on 14, 16 and 18
// A database that sends its sessions only errors cannot hide a withheld
// change from the stream: the replication connection asks for warnings
// itself, so a publication dropped under the stream still ends it with the
// server's SQLSTATE on every major.
func TestStreamStopsOnTheWarningWhenTheDatabaseSendsOnlyErrors(t *testing.T) {
	f := newSlotFixture(t)
	slot := f.createSlot(t)
	var database string
	require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT current_database()`).Scan(&database))
	_, err := f.pool.Exec(t.Context(), `ALTER DATABASE `+pgx.Identifier{database}.Sanitize()+` SET client_min_messages = error`)
	require.NoError(t, err)
	stream := f.openStream(t, slot.ConsistentPoint())
	_, err = f.pool.Exec(t.Context(), `DROP PUBLICATION `+slot.Name())
	require.NoError(t, err)
	f.exec(t, `INSERT INTO %s.ledger (id, note) VALUES (101, 'unpublished')`)

	var version int
	require.NoError(t, f.pool.QueryRow(t.Context(), `SELECT current_setting('server_version_num')::int`).Scan(&version))
	err = nextError(t, stream)
	if version >= 180000 {
		require.ErrorIs(t, err, decode.ErrInvariantViolation)
		requireSQLState(t, err, "55000")
		return
	}
	requireSQLState(t, err, "42704")
}

On e1fb311, PostgreSQL 18.6: --- FAIL: TestStreamStopsOnTheWarningWhenTheDatabaseSendsOnlyErrors (16.31s) (the stream did not fail before the stream deadline). With the fix it passes on 14 (1.47s), 16 (1.15s) and 18 (1.13s).

Non-blocking

N10. The red test (PostgreSQL 14) job comes from a contended executor test this PR does not touch.

The one failure in the job is TestQuarantineAbandonedIndexFailsClosedOnStaleObservation/index_replaced_by_a_concurrent_reindex. Its REINDEX INDEX CONCURRENTLY hit 55P03 canceling statement due to lock timeout. pkg/decode passed in that same job. Here is why the job's failure is not this PR's:

  • The PR changes nothing in pkg/executor.
  • pkg/decode tests start their own containers through StartPostgresWithSettings, which never uses PG_DSN. So they never connect to the job's shared server.
  • The test passes 3 out of 3 times on 14 in isolation.
  • It passed on 15 through 18 at this head.

The likely cause is how the matrix job runs: it puts every package on one long-lived server and database (ci.yml:140-146, Makefile:37). REINDEX CONCURRENTLY waits for every older snapshot in its own database. On 14.24, a repeatable-read snapshot held in the same database ran it into a 3s lock_timeout with this exact error. The same snapshot held in another database did not. Running that fixture on a database of its own (testutil.NewDatabase) would remove the cross-package wait. That belongs in a separate PR. I did not rerun the job.

N11. Matching the localized severity survives the suite. stream.go:240

SeverityUnlocalized is the right field: under a non-English lc_messages the server localizes Severity, for example to WARNUNG. But the unit test sets both fields to the same value, so F5 (msg.Severity != warningSeverity) passes. If the test's warning carries a localized Severity alongside SeverityUnlocalized: "WARNING", the test pins the field.

Test case: passes on e1fb311, fails F5
// The stream reads the unlocalized severity: a server running with
// non-English lc_messages localizes Severity, and its warning must still
// stop the stream.
func TestStreamStopsOnALocalizedWarning(t *testing.T) {
	s := &Stream{delivered: 100, confirmed: 90}
	_, yielded, err := s.handleMessage(t.Context(), &pgproto3.NoticeResponse{
		Severity: "WARNUNG", SeverityUnlocalized: "WARNING", Code: "55000", Message: "skipped loading publication",
	})
	require.ErrorIs(t, err, ErrInvariantViolation)
	assert.False(t, yielded)
}
Mutant (e1fb311) Caught by
F1 a warning is ignored like any notice TestStreamStopsOnAWarningFromTheServer, TestStreamReturnsTheDecodersError (18)
F2 every notice stops the stream TestStreamStopsOnAWarningFromTheServer
F3 the server's error leaves the chain (%v) TestStreamStopsOnAWarningFromTheServer, TestStreamReturnsTheDecodersError (18)
F4 the stop is not ErrInvariantViolation TestStreamStopsOnAWarningFromTheServer, TestStreamReturnsTheDecodersError (18)
F5 match the localized Severity survives; N11's test kills it
F6 a notice is dropped with ParameterStatus again TestStreamStopsOnAWarningFromTheServer, TestStreamReturnsTheDecodersError (18)
E1 (N8) a change overwrites serverWALEnd TestServerWALEndKeepsAKeepaliveAboveALaterChange
E4 (N8) OpenStream drops the target-database check TestOpenStreamRefusesAPoolOnAnotherDatabase

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: 1 blocking, 2 non-blocking (see the review comment above).

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

@Kiran01bm
Kiran01bm merged commit e45c2de into main Oct 8, 2026
27 of 29 checks passed
@Kiran01bm
Kiran01bm deleted the kiran01bm/cs7-pgoutput branch October 8, 2026 09:36
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