From ace94f0e13b78bd2f50b6bdebb36a8fc14a43bd7 Mon Sep 17 00:00:00 2001 From: "glm-5.3-flash" Date: Sat, 10 Oct 2026 04:49:35 +0000 Subject: [PATCH] =?UTF-8?q?pg=20forwarder:=20generation-tagged=20commands?= =?UTF-8?q?=20=E2=80=94=20a=20reconnect=20drops=20its=20predecessor's=20qu?= =?UTF-8?q?eued=20commands=20before=20replaying,=20closing=20the=20stale-U?= =?UTF-8?q?NLISTEN=20window=20(a=20stale=20UNLISTEN=20replayed=20after=20t?= =?UTF-8?q?he=20snapshot's=20re-issued=20LISTENs=20silently=20cancelled=20?= =?UTF-8?q?a=20re-registered=20channel,=20its=20wakes=20lost=20until=20the?= =?UTF-8?q?=20next=20reconnect);=20the=20reconcile=20decision=20is=20the?= =?UTF-8?q?=20pure=20helper=20(stale=20LISTEN=20and=20UNLISTEN=20dropped,?= =?UTF-8?q?=20current-generation=20commands=20replay=20in=20queue=20order,?= =?UTF-8?q?=20dropped=20stale=20LISTENs=20ack=20the=20transient=20mid-LIST?= =?UTF-8?q?EN=20error)=20pinned=20by=20a=20four-combination=20server-less?= =?UTF-8?q?=20unit=20test=20(replay-proven=20against=20the=20neutered=20sh?= =?UTF-8?q?ape);=20'queued=20commands=20replay=20harmlessly'=20doc=20corre?= =?UTF-8?q?cted=20and=20the=20post-reconnect=20LISTEN-set=20invariant=20(s?= =?UTF-8?q?erver=20LISTEN=20set=20=3D=20registry=20snapshot,=20dead-genera?= =?UTF-8?q?tion=20commands=20can't=20undo=20it)=20stated=20in=20the=20modu?= =?UTF-8?q?le=20+=20loop=20docs=20(task=20pg-fix-stale-unlisten,=20review?= =?UTF-8?q?=20002=20Finding=204)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- alkstore-postgres/src/forwarder.rs | 244 ++++++++++++++++++---- alkstore-postgres/src/store/open_tests.rs | 93 +++++++++ tasks/pg-fix-stale-unlisten.md | 108 +++++++++- 3 files changed, 399 insertions(+), 46 deletions(-) diff --git a/alkstore-postgres/src/forwarder.rs b/alkstore-postgres/src/forwarder.rs index 4d6861c..c47c107 100644 --- a/alkstore-postgres/src/forwarder.rs +++ b/alkstore-postgres/src/forwarder.rs @@ -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, /// 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, + /// 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>, }, - 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; + +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, + generation: u64, +) -> (Vec, Vec) { + 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, commands: mpsc::UnboundedReceiver, + generation: GenerationSlot, connected: Arc, shutdown: watch::Receiver, } @@ -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; } diff --git a/alkstore-postgres/src/store/open_tests.rs b/alkstore-postgres/src/store/open_tests.rs index 39b6142..bb52739 100644 --- a/alkstore-postgres/src/store/open_tests.rs +++ b/alkstore-postgres/src/store/open_tests.rs @@ -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 diff --git a/tasks/pg-fix-stale-unlisten.md b/tasks/pg-fix-stale-unlisten.md index 8153fc3..f11550b 100644 --- a/tasks/pg-fix-stale-unlisten.md +++ b/tasks/pg-fix-stale-unlisten.md @@ -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`) — 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 \ No newline at end of file +> 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. \ No newline at end of file