pg forwarder: generation-tagged commands — a reconnect drops its predecessor's queued commands before replaying, closing the stale-UNLISTEN window (a stale UNLISTEN replayed after the snapshot's re-issued LISTENs silently cancelled a re-registered channel, its wakes lost until the next reconnect); the reconcile decision is the pure helper (stale LISTEN and UNLISTEN dropped, current-generation commands replay in queue order, dropped stale LISTENs ack the transient mid-LISTEN error) pinned by a four-combination server-less unit test (replay-proven against the neutered shape); 'queued commands replay harmlessly' doc corrected and the post-reconnect LISTEN-set invariant (server LISTEN set = registry snapshot, dead-generation commands can't undo it) stated in the module + loop docs (task pg-fix-stale-unlisten, review 002 Finding 4)

This commit is contained in:
glm-5.3-flash committed 2026-10-10 04:49:35 +00:00
1 parent 86719a39cc
commit ace94f0e13
3 files changed
+399 -46

No files matched your search

+206 -38
View File
@@ -40,10 +40,25 @@
//! bridge posture). The wake content (`Wake { channel }` only) and
//! per-receiver bridging are the notify-listen task's; the skeleton
//! broadcasts raw notification channel names.
//!
//! # The post-reconnect LISTEN-set invariant
//!
//! After every reconnect, the server-side LISTEN set matches the
//! registry snapshot and nothing queued from a dead generation can undo
//! it: the fresh generation's loop top re-issues LISTENs from the
//! snapshot (the re-issue's authority), and every queued command
//! predating the reconnect's generation is dropped rather than
//! replayed ([`reconcile_queued_commands`] — the generation tags, the
//! review of 002 Finding 4's fix; a stale UNLISTEN replayed after the
//! re-issued LISTEN would otherwise silently cancel a channel the
//! registry believes is listening, its wakes lost until the next
//! reconnect or a fresh registration). Queued commands carry the
//! generation they were queued under; the loop bumps the shared
//! counter at each reconnect.
use std::future::Future;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::time::Duration;
use tokio::sync::{broadcast, mpsc, watch};
@@ -180,10 +195,16 @@ pub(crate) struct Forwarder {
/// register/unregister surfaces below.
channels: Arc<ChannelSet>,
/// LISTEN/UNLISTEN commands to the live connection (the loop
/// issues them; queued commands replay harmlessly across
/// reconnects alongside the registry re-issue). Read by the
/// register/unregister surfaces below.
/// issues them; each carries the generation it was queued under,
/// and a reconnect drops its predecessor's queued commands before
/// replaying — a stale UNLISTEN must never undo the re-issue,
/// review 002 Finding 4). Read by the register/unregister
/// surfaces below.
commands_tx: mpsc::UnboundedSender<ListenCommand>,
/// The connection-generation counter the queued commands are
/// tagged against (the loop bumps it at each reconnect; the
/// register/unregister surfaces load it at send time).
generation: GenerationSlot,
/// Whether a listener connection is currently live — the
/// register-path's honest transient state (registration during a
/// gap fails with `Database`; the registry entry's recovery rides
@@ -199,19 +220,125 @@ pub(crate) struct Forwarder {
}
/// Commands to the live listener connection (built by
/// `register`/`unregister` below). `Listen` carries an ack the
/// forwarder resolves after the server processed the LISTEN — the
/// registration is synchronous end-to-end, so a `listen()` that
/// returned sees every later commit (`listen` **starts from now** —
/// a notify racing the LISTEN cannot be lost; delivery only starts
/// at server-side registration).
/// `register`/`unregister` below). Every command carries the
/// generation it was queued under (the shared counter the loop bumps
/// at each reconnect): a reconnect drops queued commands predating
/// its generation before replaying — a stale UNLISTEN replayed after
/// the fresh generation's re-issued LISTENs would silently cancel a
/// channel the registry believes is listening (review 002 Finding 4).
///
/// `Listen` carries an ack the forwarder resolves after the server
/// processed the LISTEN — the registration is synchronous end-to-end,
/// so a `listen()` that returned sees every later commit (`listen`
/// **starts from now** — a notify racing the LISTEN cannot be lost;
/// delivery only starts at server-side registration).
#[derive(Debug)]
enum ListenCommand {
pub(crate) enum ListenCommand {
Listen {
generation: u64,
channel: String,
ack: tokio::sync::oneshot::Sender<alkstore::Result<()>>,
},
Unlisten(String),
Unlisten {
generation: u64,
channel: String,
},
}
/// The shared generation counter (the `Forwarder` handle tags queued
/// commands with its value at send time; the loop is the sole writer,
/// bumping it at each reconnect — a reader can never see a value the
/// loop has not yet re-issued from).
pub(crate) type GenerationSlot = Arc<AtomicU64>;
impl ListenCommand {
/// The generation the command was queued under (the reconnect's
/// drop decision reads this).
fn generation(&self) -> u64 {
match self {
ListenCommand::Listen { generation, .. } => *generation,
ListenCommand::Unlisten { generation, .. } => *generation,
}
}
/// Resolve a dropped stale command: a queued-unprocessed LISTEN's
/// caller awaits its ack — it receives the transient mid-LISTEN
/// error (the connection it targeted is gone; the registration is
/// recovered by the re-issue, `register`'s documented posture). A
/// dropped UNLISTEN is best-effort — nothing to resolve.
pub(crate) fn abandon_stale(self) {
if let ListenCommand::Listen { ack, .. } = self {
let _ = ack.send(Err(mid_listen_error()));
}
}
}
/// The stale-command reconcile (review 002 Finding 4's decision, the
/// pure helper the reconnect and the replay loop share): from the
/// commands queued when the previous connection generation ended and
/// the generation just re-connected into, `(to replay, dropped)`.
///
/// A command replays iff it was queued under the **current**
/// generation; everything older is dropped and never touches the
/// fresh connection. The invariant this enforces: after a reconnect,
/// the server-side LISTEN set matches the registry snapshot (the
/// re-issue) and nothing queued from a dead generation can undo it —
/// the snapshot re-issue covers every stale LISTEN (the registry is
/// the source of truth; registry-write-first ordering), and the stale
/// UNLISTENs are exactly what must not replay (the finding's silent
/// wake-loss). Note the decision rides the generation tags alone:
/// the snapshot's role is the re-issue's, not this filter's (the
/// generation shape of the review's two candidates).
pub(crate) fn reconcile_queued_commands(
queued: Vec<ListenCommand>,
generation: u64,
) -> (Vec<ListenCommand>, Vec<ListenCommand>) {
let mut replay = Vec::new();
let mut dropped = Vec::new();
for command in queued {
if command.generation() < generation {
dropped.push(command);
} else {
replay.push(command);
}
}
(replay, dropped)
}
/// The mid-LISTEN transient error — a `register` whose LISTEN did not
/// land (died mid-command, or dropped stale at a reconnect) acks/
/// returns this; the retry recovers it (the registry entry stays
/// either way and the re-issue covers the channel on reconnect).
fn mid_listen_error() -> alkstore::Error {
database_error(
"the listener connection died mid-LISTEN; \
the channel is registered and re-issued on reconnect",
)
}
/// Execute one command against the live connection (the replay's and
/// the serve loop's shared arm). A `Listen` acks its caller
/// synchronously; an `Unlisten` is best-effort.
async fn run_command(client: &tokio_postgres::Client, command: ListenCommand) {
match command {
ListenCommand::Listen { channel, ack, .. } => {
let result = if client
.batch_execute(&format!("LISTEN {};", quote_ident(&channel)))
.await
.is_ok()
{
Ok(())
} else {
Err(mid_listen_error())
};
let _ = ack.send(result);
}
ListenCommand::Unlisten { channel, .. } => {
let _ = client
.batch_execute(&format!("UNLISTEN {};", quote_ident(&channel)))
.await;
}
}
}
impl Forwarder {
@@ -228,6 +355,7 @@ impl Forwarder {
let channels = Arc::new(ChannelSet::default());
let (commands_tx, commands_rx) = mpsc::unbounded_channel();
let (shutdown_tx, shutdown_rx) = watch::channel(false);
let generation: GenerationSlot = Arc::new(AtomicU64::new(0));
let reconnect_config_slot: ReconnectConfigSlot =
Arc::new(std::sync::Mutex::new(reconnect_config));
// The connection handed in was just established — the liveness
@@ -246,6 +374,7 @@ impl Forwarder {
fanout_slot: fanout_slot.clone(),
channels: channels.clone(),
commands: commands_rx,
generation: generation.clone(),
connected: connected.clone(),
shutdown: shutdown_rx,
}));
@@ -254,6 +383,7 @@ impl Forwarder {
fanout: fanout_slot.clone(),
channels,
commands_tx,
generation,
connected,
#[cfg(test)]
reconnect_config: reconnect_config_slot,
@@ -303,6 +433,7 @@ impl Forwarder {
let sent = self
.commands_tx
.send(ListenCommand::Listen {
generation: self.generation.load(Ordering::Acquire),
channel: channel.to_string(),
ack: ack_tx,
})
@@ -321,12 +452,18 @@ impl Forwarder {
/// write first (ordering symmetric with [`Forwarder::register`]'s
/// recovery logic); the UNLISTEN is best-effort — a lost one
/// re-issues nothing (UNLISTEN is not recovery-relevant: the
/// snapshot won't carry the removed channel).
/// snapshot won't carry the removed channel). The queued UNLISTEN
/// rides its send-time generation: if the connection died before
/// processing it, the reconnect drops it rather than replaying —
/// a stale UNLISTEN after the re-issued LISTENs would otherwise
/// silently cancel a re-registered channel (review 002 Finding 4;
/// the fresh connection's LISTEN set is the snapshot's).
pub(crate) fn unregister(&self, channel: &str) {
if self.channels.remove_subscriber(channel) {
let _ = self
.commands_tx
.send(ListenCommand::Unlisten(channel.to_string()));
let _ = self.commands_tx.send(ListenCommand::Unlisten {
generation: self.generation.load(Ordering::Acquire),
channel: channel.to_string(),
});
}
}
@@ -438,6 +575,17 @@ impl Drop for FanoutRelease {
/// successful connect re-binds it — there is no empty-slot arm left to
/// fall back into.
///
/// The reconnect bumps the shared generation counter (before the
/// snapshot re-issue, so commands queued past the bump — a healthy
/// `register`'s among them — tag as current-generation and replay);
/// at each fresh generation's top the predecessor's queued commands
/// are drained through [`reconcile_queued_commands`] — everything
/// tagged under a dead generation is dropped, never replayed onto the
/// new connection (the stale-UNLISTEN hazard, review 002 Finding 4).
/// The invariant: **after a reconnect, the server-side LISTEN set
/// matches the registry snapshot and nothing queued from a dead
/// generation can undo it.**
///
/// The receiver-facing close shape (ADR-021 §5's pg arm): the fanout
/// sender clones live here (loop) and on the shared slot (taken by the
/// shutdown flip or the exit-path release guard — see
@@ -454,6 +602,7 @@ struct LoopCtx {
fanout_slot: FanoutSlot,
channels: Arc<ChannelSet>,
commands: mpsc::UnboundedReceiver<ListenCommand>,
generation: GenerationSlot,
connected: Arc<AtomicBool>,
shutdown: watch::Receiver<bool>,
}
@@ -467,6 +616,7 @@ async fn forwarder_loop(
fanout_slot,
channels,
mut commands,
generation: generation_slot,
connected,
mut shutdown,
}: LoopCtx,
@@ -476,6 +626,7 @@ async fn forwarder_loop(
let mut backoff_ms = RECONNECT_BASE_MS;
let mut first_listen = true;
let mut connection = connection;
let mut generation = generation_slot.load(Ordering::Acquire);
loop {
// The dedicated poll loop, spawned BEFORE any client query can
@@ -506,6 +657,25 @@ async fn forwarder_loop(
if listen_ok {
backoff_ms = RECONNECT_BASE_MS;
connected.store(true, Ordering::Release);
// The stale command drain (review 002 Finding 4): the
// re-issue above is the server-side LISTEN set's authority
// — the predecessor generation's queued commands are
// reconciled before the serve loop starts, so a stale
// UNLISTEN can never replay after it and silently cancel a
// live channel. Commands queued past the bump (a
// register/unregister that raced the reconnect) replay
// normally.
let mut queued = Vec::new();
while let Ok(cmd) = commands.try_recv() {
queued.push(cmd);
}
let (replay, dropped) = reconcile_queued_commands(queued, generation);
for cmd in dropped {
cmd.abandon_stale();
}
for cmd in replay {
run_command(&client, cmd).await;
}
if !first_listen {
// The synthetic reconnect-wake (ADR-008 §4's reserved
// channel; ADR-006's recovery semantics) broadcast to
@@ -543,7 +713,17 @@ async fn forwarder_loop(
cmd = commands.recv() => {
match cmd {
Some(ListenCommand::Listen { channel, ack }) => {
Some(command) if command.generation() < generation => {
// A stale command queued under a dead
// generation (a send racing the
// reconnect's drain): dropped here for
// the same reason the generation-top
// drain drops its batch — it must not
// touch this connection (review 002
// Finding 4).
command.abandon_stale();
}
Some(command) => {
// Acked after the server processed the
// LISTEN: a `listen()` return stamps
// "registered from now" end-to-end. A
@@ -551,27 +731,7 @@ async fn forwarder_loop(
// mid-command) acks the failure — the
// caller sees it; the registry entry
// rides the reconnect re-issue.
let result = if client
.batch_execute(&format!(
"LISTEN {};",
quote_ident(&channel)
))
.await
.is_ok()
{
Ok(())
} else {
Err(database_error(
"the listener connection died mid-LISTEN; \
the channel is registered and re-issued on reconnect",
))
};
let _ = ack.send(result);
}
Some(ListenCommand::Unlisten(ch)) => {
let _ = client
.batch_execute(&format!("UNLISTEN {};", quote_ident(&ch)))
.await;
run_command(&client, command).await;
}
// The command sender is gone (the store —
// and its forwarder handle — went away
@@ -626,6 +786,14 @@ async fn forwarder_loop(
// return — the guard and the client drop carry the rest.
return;
};
// The generation bump happens here — *before* the fresh
// generation's loop top re-issues the snapshot — so every
// command queued under the dead connection is stale-tagged
// against the new one and the generation-top drain drops it
// (review 002 Finding 4); commands queued from here on
// (register's honest error path, the re-issue) tag current.
generation += 1;
generation_slot.store(generation, Ordering::Release);
client = next_client;
connection = next_connection;
}
+93
View File
@@ -926,6 +926,99 @@ fn fanout_release_guard_empties_the_shared_slot_on_drop() {
const FANOUT_TEST_CAPACITY: usize = 1;
/// Unit pin (server-less) — the stale-command reconcile (review 002
/// Finding 4's decision), **all four command/staleness combinations**:
/// a stale-generation UNLISTEN is dropped (the hazard — replayed after
/// the fresh generation's re-issued LISTEN it would silently cancel a
/// channel the registry believes is listening), a stale-generation
/// LISTEN is dropped (redundant — the snapshot re-issue covers it),
/// and both current-generation commands replay normally. Within the
/// replay the queue order is preserved; each dropped LISTEN resolves
/// its awaiting caller with the transient mid-LISTEN `Database` error
/// (the registration is recovered by the re-issue — `register`'s
/// documented posture).
#[tokio::test]
async fn stale_generation_commands_drop_and_current_generation_replay() {
use crate::forwarder::{ListenCommand, reconcile_queued_commands};
let listen_with = |generation: u64, channel: &str| {
let (ack_tx, _ack_rx) = tokio::sync::oneshot::channel();
ListenCommand::Listen {
generation,
channel: channel.to_string(),
ack: ack_tx,
}
};
let stale_channel = "stale-channel";
let current_channel = "current-channel";
// Queue order preserved: the finding's sequence has the stale
// commands queued before anything fresh.
let queued = vec![
listen_with(1, stale_channel),
ListenCommand::Unlisten {
generation: 1,
channel: stale_channel.to_string(),
},
listen_with(2, current_channel),
ListenCommand::Unlisten {
generation: 2,
channel: current_channel.to_string(),
},
];
let (replay, dropped) = reconcile_queued_commands(queued, 2);
fn describe(command: &ListenCommand) -> (u64, String, bool) {
match command {
ListenCommand::Listen {
generation,
channel,
..
} => (*generation, channel.clone(), true),
ListenCommand::Unlisten {
generation,
channel,
} => (*generation, channel.clone(), false),
}
}
let replayed: Vec<_> = replay.iter().map(describe).collect();
assert_eq!(
replayed,
vec![
(2, current_channel.to_string(), true),
(2, current_channel.to_string(), false),
],
"current-generation commands replay, queue order preserved"
);
let dropped: Vec<_> = dropped.iter().map(describe).collect();
assert_eq!(
dropped,
vec![
(1, stale_channel.to_string(), true),
(1, stale_channel.to_string(), false),
],
"stale-generation commands drop — the snapshot re-issue covers the \
stale LISTEN, and the stale UNLISTEN never undoes it"
);
// A dropped stale LISTEN resolves its awaiting caller with the
// transient mid-LISTEN `Database` error (never a spurious
// "loop is gone" ack-drop, never an Ok — the LISTEN it acked did
// not land; the re-issue covers the registration).
let (ack_tx, ack_rx) = tokio::sync::oneshot::channel();
crate::forwarder::ListenCommand::Listen {
generation: 1,
channel: stale_channel.to_string(),
ack: ack_tx,
}
.abandon_stale();
let acked = ack_rx.await.expect("the ack resolves");
assert!(
matches!(acked, Err(alkstore::Error::Database(_))),
"the dropped stale LISTEN acks the transient mid-LISTEN error, got: {acked:?}"
);
}
/// Acceptance (review 002 Finding 1, end-to-end): a failed reconnect
/// connect no longer kills the forwarder. The reconnect-config seam
/// points the loop's connect attempts at an unreachable endpoint, a
+100 -8
View File
@@ -1,7 +1,7 @@
---
id: pg-fix-stale-unlisten
name: Fix stale queued UNLISTEN cancelling a re-issued LISTEN after reconnect (review 002 Finding 4)
status: pending
status: completed
depends_on: [pg-fix-forwarder-reconnect]
scope: narrow
risk: medium
@@ -66,18 +66,18 @@ shape.
## Acceptance Criteria
- [ ] A stale queued UNLISTEN can no longer unlisten a channel the
- [x] A stale queued UNLISTEN can no longer unlisten a channel the
registry believes is listening — the generation-drop (or
reconcile) decision pinned by server-less unit tests covering
all four command/staleness combinations
- [ ] Current-generation commands (acked registrations included)
- [x] Current-generation commands (acked registrations included)
still process normally after a reconnect — existing
register/ack tests stay green
- [ ] The "replay harmlessly" doc comment corrected to the actual
- [x] The "replay harmlessly" doc comment corrected to the actual
posture
- [ ] The invariant stated in the forwarder docs (post-reconnect
- [x] The invariant stated in the forwarder docs (post-reconnect
server LISTEN set = registry snapshot)
- [ ] `cargo test -p alkstore-postgres` (harness server), clippy
- [x] `cargo test -p alkstore-postgres` (harness server), clippy
`-D warnings`, fmt clean; gates green server-less
## References
@@ -93,8 +93,100 @@ shape.
## Notes
> To be filled by implementation agent
> Decisions of record the implementation made that the description
> didn't pin:
- **Shape chosen: candidate 1 (generation-tagging), pure.** The
reconcile-against-registry-snapshot shape was not taken; the
decision therefore rides the generation tags *alone* — the helper's
registry-snapshot input from the task's helper sketch is not a
parameter (its role is the re-issue's, not the filter's; recorded in
the helper's doc comment). The register/unregister *same-generation*
race (an UNLISTEN racing a fresh register on the live connection,
queue order LISTEN-then-UNLISTEN) is adjacent but out of this
finding's scope (the finding is about dead-generation replay only)
— unchanged, pre-existing; the reconcile shape would have been the
fix for it, but it was not among the finding's cited effects.
- **Tagging mechanics**: a shared `GenerationSlot`
(`Arc<AtomicU64>`) — the loop is the sole writer, bumping the
counter right after `reconnect_with_retry` returns a fresh pair and
*before* the fresh generation's snapshot re-issue, so every command
queued under the dead connection is stale against the new one; the
`register`/`unregister` surfaces load the counter at send time.
During a gap `register` sends nothing (its `connected==false` honest
error stands; registry-write-first covers the recover), so the only
commands queued under a dead generation are pre-death ones and
unregister's best-effort UNLISTENs.
- **The reconcile is applied at two sites** with one shared rule:
(a) the generation-top drain — each fresh, re-issue-verified
generation drains the queue via `try_recv`, feeds it through the
pure helper `reconcile_queued_commands(queued, generation) →
(replay, dropped)`, and runs the survivors through the *extracted*
`run_command` arm (the serve loop's old LISTEN/UNLISTEN body,
deduped); (b) the serve loop's recv arm — a staleness guard
(`command.generation() < generation`) dropping the one racy case the
drain misses (a send whose tag loaded before the bump but that
landed after the drain reached empty). Site (b) reuses the same
decision (`< current`) and the same `abandon_stale` drop path.
- **A dropped stale LISTEN acks its awaiting caller with the
transient mid-LISTEN `Database` error** (`abandon_stale`, the
pre-existing "died mid-LISTEN; re-issued on reconnect" message,
now a shared `mid_listen_error()` helper) — never a spurious
"loop is gone" ack-drop, never an `Ok` (the LISTEN it acks did not
land; the registration is recovered by the re-issue). A dropped
stale UNLISTEN resolves nothing (best-effort, `unregister`'s
posture). The drain runs *after* the re-issue, so even the
load-race-dropped LISTEN's channel was already covered by the
re-issue at ack time.
- **The behavioral pin through the live forwarder was not added** —
the task calls the unit pin the acceptance bar; a live arrangement
would need the UNLISTEN to deterministically lose the race against
the backend kill before the loop's next `recv` (timing-narrow, and
the drain site is unreachable without forced mid-flight deaths), so
the harness battery keeps the existing pins (register/ack,
backend-kill reconnect, failed-connect survival/recovery — all
green unchanged) and the replay-proof below carries the
regression evidence.
- **Replay-proofed**: with `reconcile_queued_commands` temporarily
neutered (everything replayed, the old verbatim-replay shape) the
new unit pin fails exactly on the stale-generation commands
replaying; with the fix restored it passes (verified both ways).
## Summary
> To be filled on completion
> What landed, verified how:
- **`alkstore-postgres/src/forwarder.rs`**: `ListenCommand` is
generation-tagged (`Listen { generation, channel, ack }`,
`Unlisten { generation, channel }`), the shared
`GenerationSlot` counter threads the loop (sole writer, bumped at
each reconnect before the snapshot re-issue) and the
register/unregister send sites; the pure
`reconcile_queued_commands(queued, generation) → (replay, dropped)`
is applied at the fresh generation's post-re-issue drain and mirrored
by the serve loop's recv-arm staleness guard; dropped stale LISTENs
resolve their callers' acks with the transient error via
`abandon_stale`/`mid_listen_error`; the old inline serve-arm
LISTEN/UNLISTEN body is now the shared `run_command`.
- **Docs**: the "queued commands replay harmlessly across
reconnects" field comment corrected to the actual posture
(generation-tagged drop-before-replay), the invariant ("after a
reconnect, the server-side LISTEN set matches the registry snapshot
and nothing queued from a dead generation can undo it") stated in
both the module doc (new "post-reconnect LISTEN-set invariant"
section) and the `forwarder_loop` doc, and the
register/unregister doc comments extended with their side of the
posture.
- **`alkstore-postgres/src/store/open_tests.rs`** (1 new test):
`stale_generation_commands_drop_and_current_generation_replay` —
the four command/staleness combinations (stale/current ×
LISTEN/UNLISTEN) in queue order through the pure helper, the
in-replay order preservation, and the dropped stale LISTEN's
transient `Database` ack.
- **Gates**: `cargo test -p alkstore-postgres` with the harness
server green (119 lib + 10 contract-suite + 9 schema, run twice
full — before and after the fmt pass); every pre-existing
forwarder/reconnect/register/ack pin green unchanged; workspace
`cargo build`, server-less `cargo test` (188 lib tests across the
crates), `cargo clippy --all-targets -- -D warnings`,
`cargo fmt --check` all green.