From 7ae426a01dc92b9a89a6e4affc21154d149019bd Mon Sep 17 00:00:00 2001 From: "glm-5.3-flash" Date: Fri, 9 Oct 2026 22:31:05 +0000 Subject: [PATCH] =?UTF-8?q?pg=20forwarder:=20the=20reconnect=20arm=20retri?= =?UTF-8?q?es=20failed=20connects=20until=20success=20or=20shutdown=20?= =?UTF-8?q?=E2=80=94=20one=20failed=20connect=20no=20longer=20kills=20the?= =?UTF-8?q?=20loop=20permanently;=20every=20loop=20exit=20path=20releases?= =?UTF-8?q?=20the=20fanout=20sender=20via=20the=20shared-slot=20drop=20gua?= =?UTF-8?q?rd=20(receivers-terminal-iff-loop-gone);=20reconnect-config=20t?= =?UTF-8?q?est=20seam=20+=20failed-connect=20pins:=20six-cycle=20exhausted?= =?UTF-8?q?-backoff=20server-less=20unit,=20pre-flip=20abort,=20guard=20dr?= =?UTF-8?q?op,=20and=20the=20harness=20end-to-end=20unreachable-outage=20?= =?UTF-8?q?=E2=86=92=20full-recovery=20test=20(bug=20replay-proven=20again?= =?UTF-8?q?st=20the=20old=20shape)=20(task=20pg-fix-forwarder-reconnect,?= =?UTF-8?q?=20review=20002=20Finding=201)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- alkstore-postgres/Cargo.toml | 3 +- alkstore-postgres/src/forwarder.rs | 229 +++++++++++++++++----- alkstore-postgres/src/store/open_tests.rs | 218 ++++++++++++++++++++ tasks/pg-fix-forwarder-reconnect.md | 90 ++++++++- 4 files changed, 482 insertions(+), 58 deletions(-) diff --git a/alkstore-postgres/Cargo.toml b/alkstore-postgres/Cargo.toml index 2496dae..1369ed5 100644 --- a/alkstore-postgres/Cargo.toml +++ b/alkstore-postgres/Cargo.toml @@ -14,4 +14,5 @@ tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "time" tokio-postgres = "0.7" [dev-dependencies] -alkstore-contract-suite = { path = "../alkstore-contract-suite" } \ No newline at end of file +alkstore-contract-suite = { path = "../alkstore-contract-suite" } +tokio = { version = "1", features = ["test-util"] } \ No newline at end of file diff --git a/alkstore-postgres/src/forwarder.rs b/alkstore-postgres/src/forwarder.rs index f438a2d..4d6861c 100644 --- a/alkstore-postgres/src/forwarder.rs +++ b/alkstore-postgres/src/forwarder.rs @@ -41,6 +41,7 @@ //! per-receiver bridging are the notify-listen task's; the skeleton //! broadcasts raw notification channel names. +use std::future::Future; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; @@ -50,6 +51,20 @@ use tokio_postgres::NoTls; use crate::seam::database_error; +/// The raw fanout sender's shared slot — the forwarder handle and the +/// loop both take it away on their exit paths, so receivers see the +/// terminal arm from *either* side. +pub(crate) type FanoutSlot = Arc>>>; + +/// The reconnect config's shared slot — the loop's connect step reads +/// it fresh per attempt; the `#[cfg(test)]` seam swaps it (the +/// unreachable-endpoint injection the failed-connect tests ride). +pub(crate) type ReconnectConfigSlot = Arc>; + +/// A fresh connection generation the reconnect produces: the poll +/// loop's next `Connection` and the `Client` driving its queries. +pub(crate) type GenerationPair = (tokio_postgres::Client, ListenerConnection); + /// The listener connection's concrete type for the `NoTls` posture /// (the engine's TLS story rides the consumer's `Config` sslmode for /// the pooled path; the dedicated listener connects with `NoTls` — @@ -154,11 +169,13 @@ impl ChannelSet { #[derive(Debug)] pub(crate) struct Forwarder { /// The wake fanout — every receiver bridge holds a subscription; - /// its release (this handle's take at shutdown, the loop's drop at - /// its exit) is the receivers' terminal `Closed`/`None` arm. - /// `Option` because shutdown empties it while receiver-held - /// `Arc`s keep the struct alive. - fanout: std::sync::Mutex>>, + /// its release (the shutdown take, or the loop's exit-path release + /// guard) is the receivers' terminal `Closed`/`None` arm. The + /// shared [`FanoutSlot`] exactly because receiver-held + /// `Arc`s keep this struct alive — the release fires + /// from both the handle's shutdown and the loop's exit + /// (receivers-terminal-iff-loop-gone). + fanout: FanoutSlot, /// The channel registry the reconnect re-issues from; read by the /// register/unregister surfaces below. channels: Arc, @@ -172,6 +189,12 @@ pub(crate) struct Forwarder { /// gap fails with `Database`; the registry entry's recovery rides /// the reconnect re-issue). Read by `register` below. connected: Arc, + /// The reconnect config's shared slot — test-observation surface + /// for the failed-connect seam (the loop's own [`LoopCtx`] clone + /// is the live read path; the field itself exists for the + /// `#[cfg(test)]` swap accessors). + #[cfg(test)] + reconnect_config: ReconnectConfigSlot, shutdown: watch::Sender, } @@ -201,9 +224,12 @@ impl Forwarder { reconnect_config: tokio_postgres::Config, ) -> Forwarder { let (fanout, _) = broadcast::channel(FANOUT_CAPACITY); + let fanout_slot: FanoutSlot = Arc::new(std::sync::Mutex::new(Some(fanout.clone()))); let channels = Arc::new(ChannelSet::default()); let (commands_tx, commands_rx) = mpsc::unbounded_channel(); let (shutdown_tx, shutdown_rx) = watch::channel(false); + let reconnect_config_slot: ReconnectConfigSlot = + Arc::new(std::sync::Mutex::new(reconnect_config)); // The connection handed in was just established — the liveness // flag starts true and the loop clears it at each generation's // death (the register path's honest transient state). @@ -214,9 +240,10 @@ impl Forwarder { // lifetime). tokio::spawn(forwarder_loop(LoopCtx { client, - connection: Some(connection), - reconnect_config, + connection, + reconnect_config: reconnect_config_slot.clone(), fanout: fanout.clone(), + fanout_slot: fanout_slot.clone(), channels: channels.clone(), commands: commands_rx, connected: connected.clone(), @@ -224,18 +251,22 @@ impl Forwarder { })); Forwarder { - fanout: std::sync::Mutex::new(Some(fanout)), + fanout: fanout_slot.clone(), channels, commands_tx, connected, + #[cfg(test)] + reconnect_config: reconnect_config_slot, shutdown: shutdown_tx, } } /// Subscribe to the raw fanout (the notify task bridges broadcast - /// receivers into `WakeReceiver`s). `None` after shutdown (the - /// senders are being released) — a subscribe racing `close()` or - /// arriving past it fails closed rather than panicking. + /// receivers into `WakeReceiver`s). `None` after shutdown — or + /// after *any* forwarder-loop exit (the loop's release guard takes + /// the shared sender too, the defense-in-depth arm: a subscribe + /// racing a dead loop fails closed rather than parking on a + /// forever-silent broadcast). pub(crate) fn subscribe(&self) -> Option> { self.fanout .lock() @@ -302,9 +333,10 @@ impl Forwarder { /// Shutdown: flip the switch — the loop exits, dropping the /// listener `Client` and its session (the listener connection is /// gone) — and release this handle's fanout sender. The broadcast - /// closes once every sender is dropped (the loop's clone dies with - /// its exit; this handle's dies here — receivers unaffected by - /// `Arc` lifetimes see the terminal arm). Idempotent. + /// closes once every sender is dropped (the loop's exit-path + /// release guard empties the shared slot too; this handle's take + /// fires the terminal arm at the flip — receivers unaffected by + /// `Arc` lifetimes see it). Idempotent. pub(crate) fn shutdown(&self) { let _ = self.shutdown.send(true); // Take the sender even if a receiver holds this Arc: the @@ -312,32 +344,114 @@ impl Forwarder { // whatever lifetime stands on the handle. *self.fanout.lock().expect("fanout mutex poisoned") = None; } + + /// Aim the reconnect arm's connect step at `config` (the failed- + /// connect test seam — the unreachable-endpoint injection). The + /// loop reads the slot fresh per attempt, so a swap lands on the + /// next attempt. + #[cfg(test)] + pub(crate) fn swap_reconnect_config(&self, config: tokio_postgres::Config) { + *self + .reconnect_config + .lock() + .expect("reconnect config mutex poisoned") = config; + } + + /// The reconnect config's current content (the failed-connect + /// test's restore point — the swap-back carries the original). + #[cfg(test)] + pub(crate) fn reconnect_config_snapshot(&self) -> tokio_postgres::Config { + self.reconnect_config + .lock() + .expect("reconnect config mutex poisoned") + .clone() + } +} + +/// The reconnect arm's retry core (review 002 Finding 1's fix): sleep +/// → attempt the connect — repeat until an attempt succeeds or the +/// shutdown flip, **never returning with nothing connected** (a failed +/// connect retries inside this arm; the skeleton's fall-back-into-the- +/// loop-with-an-empty-slot shape — one failed connect killing the +/// loop forever — is structurally gone). +/// +/// `backoff_ms` carries the live exponential backoff (doubled per +/// attempt, capped at [`RECONNECT_MAX_MS`]; the caller resets it on +/// the next *verified-alive* generation). The shutdown flip aborts +/// with `None` — the only exit besides `Some`'s fresh generation. +pub(crate) async fn reconnect_with_retry( + mut attempt: F, + backoff_ms: &mut u64, + shutdown: &mut watch::Receiver, +) -> Option +where + F: FnMut() -> Fut, + Fut: Future>, +{ + loop { + tokio::select! { + biased; + + // Shutdown (or the sender's drop) exits mid-backoff — the + // subscribers' terminal arm is not held hostage by a + // backoff cycle. + _ = shutdown.changed() => return None, + + _ = tokio::time::sleep(Duration::from_millis(*backoff_ms)) => {} + } + *backoff_ms = backoff_ms.saturating_mul(2).min(RECONNECT_MAX_MS); + if let Some(generation) = attempt().await { + return Some(generation); + } + } +} + +/// Releases the forwarder handle's fanout sender at the loop's exit — +/// *every* exit path: Drop runs on return, unwind, and task abort +/// alike. This is the receivers-terminal-iff-loop-gone invariant's +/// enforcement side (the shutdown take covers the flip; this guard +/// covers every other exit the defense-in-depth posture guards +/// against) — a dead loop can never leave subscribers parked on a +/// forever-silent broadcast. +pub(crate) struct FanoutRelease { + pub(crate) slot: FanoutSlot, +} + +impl Drop for FanoutRelease { + fn drop(&mut self) { + *self.slot.lock().expect("fanout mutex poisoned") = None; + } } /// The forwarder loop (the POC's ~90-line shape, restructured for the -/// dynamic channel set + command issuer + shutdown). +/// dynamic channel set + command issuer + shutdown + the reconnect +/// retry arm). /// /// Per connection generation: spawn the poll task (it owns the /// `Connection` and fans out `poll_message`'s notifications) → issue /// LISTENs from the live channel-set snapshot → serve commands and -/// wait for the connection's death → reconnect with exponential -/// backoff (50 ms → 2 s cap). The loop consumes the handed-in -/// connection on its first generation and reconnects into the slot -/// afterwards. +/// wait for the connection's death → reconnect via +/// [`reconnect_with_retry`] (exponential backoff 50 ms → 2 s cap; +/// failed connects never leave the arm). The loop-carried connection +/// is *always present* at the generation's top by construction: the +/// first iteration holds the handed-in connection, and only a +/// successful connect re-binds it — there is no empty-slot arm left to +/// fall back into. /// /// The receiver-facing close shape (ADR-021 §5's pg arm): the fanout -/// sender clones live here (loop) and on the handle (taken at -/// shutdown); a connection death and its reconnect are invisible to -/// subscribers (the fanout never drops on a generation boundary — -/// receivers stay open across reconnects, they receive the reserved -/// reconnect-wake instead); the senders drop only at shutdown — the -/// loop's clone at its exit, the handle's at the shutdown take — the -/// receivers' terminal arm. +/// sender clones live here (loop) and on the shared slot (taken by the +/// shutdown flip or the exit-path release guard — see +/// [`FanoutRelease`]); a connection death and its reconnect are +/// invisible to subscribers (the fanout never drops on a generation +/// boundary — receivers stay open across reconnects, they receive the +/// reserved reconnect-wake instead); the senders drop only when the +/// loop is gone — the invariant's shape. struct LoopCtx { client: tokio_postgres::Client, - connection: Option, - reconnect_config: tokio_postgres::Config, + connection: ListenerConnection, + reconnect_config: ReconnectConfigSlot, fanout: broadcast::Sender, + fanout_slot: FanoutSlot, channels: Arc, commands: mpsc::UnboundedReceiver, connected: Arc, @@ -346,18 +460,22 @@ struct LoopCtx { async fn forwarder_loop( LoopCtx { - mut client, - mut connection, + client, + connection, reconnect_config, fanout, + fanout_slot, channels, mut commands, connected, mut shutdown, }: LoopCtx, ) { + let _fanout_release = FanoutRelease { slot: fanout_slot }; + let mut client = client; let mut backoff_ms = RECONNECT_BASE_MS; let mut first_listen = true; + let mut connection = connection; loop { // The dedicated poll loop, spawned BEFORE any client query can @@ -366,10 +484,7 @@ async fn forwarder_loop( // the poll task co-resident with client queries forever). The // poll task owns the Connection; the loop drives queries // through the Client. - let Some(poll_connection) = connection.take() else { - return; - }; - let mut poll_task = tokio::spawn(poll_loop(poll_connection, fanout.clone())); + let mut poll_task = tokio::spawn(poll_loop(connection, fanout.clone())); // Issue LISTENs from the dynamic channel set's snapshot (on // first connect it carries pre-registered channels; on @@ -461,7 +576,8 @@ async fn forwarder_loop( // The command sender is gone (the store — // and its forwarder handle — went away // without a shutdown flip): serve no more, - // exit. + // exit (the release guard fires; the + // receivers' terminal arm). None => { poll_task.abort(); return; @@ -486,27 +602,32 @@ async fn forwarder_loop( // Reconnect: the Client is dead with the session (dropping it // here closes nothing server-side — pitfall 2 was satisfied - // for the connection's whole lifetime). Back off exponentially - // (50 ms → 2 s cap) and reconnect. + // for the connection's whole lifetime). The retry core sleeps + // and attempts the connect — with the *fresh* slot content per + // attempt (the test seam swaps it) — until a connect succeeds + // or the shutdown flip; a failed connect never falls back into + // the loop above without a connection (the loop's `None` arm + // is restructured away entirely). connected.store(false, Ordering::Release); - tokio::time::sleep(Duration::from_millis(backoff_ms)).await; - backoff_ms = backoff_ms.saturating_mul(2).min(RECONNECT_MAX_MS); - - if *shutdown.borrow() { + let attempt = { + let slot = reconnect_config.clone(); + move || { + let guard = slot.lock().expect("reconnect config mutex poisoned"); + let config = guard.clone(); + drop(guard); + async move { config.connect(NoTls).await.ok() } + } + }; + let Some((next_client, next_connection)) = + reconnect_with_retry(attempt, &mut backoff_ms, &mut shutdown).await + else { + // Shutdown mid-backoff/mid-retry: the poll task is already + // done (the generation's end preceded this arm); plain + // return — the guard and the client drop carry the rest. return; - } - match reconnect_config.connect(NoTls).await { - Ok((new_client, new_connection)) => { - client = new_client; - connection = Some(new_connection); - } - Err(_e) => { - // Still unreachable — back off again and retry. (The - // skeleton carries the loop shape; the notify-listen - // task owns the reconnect behavior's full contract.) - continue; - } - } + }; + 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 9d2a5bf..c1ba633 100644 --- a/alkstore-postgres/src/store/open_tests.rs +++ b/alkstore-postgres/src/store/open_tests.rs @@ -734,3 +734,221 @@ fn seam_mappings_preserve_the_source_chain() { let mapped = pool_error(deadpool_postgres::PoolError::Closed); assert!(matches!(mapped, alkstore::Error::Database(_))); } + +/// Unit pin (server-less, paused time) — the disconnect-fix's retry +/// core (review 002 Finding 1): a *failing* connect is retried once +/// per backoff cycle across multiple doublings of the 50 ms base, and +/// the backoff reaches the 2 s cap — the loop's shutdown flip is the +/// only way out (`None`), never a spurious success. +#[tokio::test(start_paused = true)] +async fn reconnect_retry_attempts_failed_connects_across_backoff_cycles() { + use tokio::sync::watch; + + let (shutdown_tx, mut shutdown_rx) = watch::channel(false); + let attempts = std::sync::Arc::new(AtomicU64::new(0)); + let flip_after = attempts.clone(); + let let_flip_tx = shutdown_tx.clone(); + let always_failing = move || { + let attempts = flip_after.clone(); + let flip_tx = let_flip_tx.clone(); + async move { + let n = attempts.fetch_add(1, Ordering::SeqCst) + 1; + if n == 6 { + let _ = flip_tx.send(true); + } + None:: + } + }; + + let mut backoff_ms = crate::forwarder::RECONNECT_BASE_MS; + let result = + crate::forwarder::reconnect_with_retry(always_failing, &mut backoff_ms, &mut shutdown_rx) + .await; + + assert!( + result.is_none(), + "with every connect failing, the retry's only exit is the shutdown flip" + ); + assert_eq!( + attempts.load(Ordering::SeqCst), + 6, + "one connect attempt per backoff cycle: 50 → 100 → 200 → 400 → 800 → 1600 ms \ + (multiple failed connects retried, the dead-loop shape would give one and stop)" + ); + assert_eq!( + backoff_ms, + crate::forwarder::RECONNECT_MAX_MS, + "the doubling caps at 2 s" + ); +} + +/// Unit pin (server-less): the shutdown flip *before* the retry core +/// starts aborts immediately — no connect attempt rides a shut-down +/// forwarder. +#[tokio::test(start_paused = true)] +async fn reconnect_retry_aborts_before_any_attempt_on_preflipped_shutdown() { + use tokio::sync::watch; + + let (shutdown_tx, mut shutdown_rx) = watch::channel(false); + shutdown_tx.send(true).expect("watch alive"); + let attempts = std::sync::Arc::new(AtomicU64::new(0)); + let attempts2 = attempts.clone(); + let always_failing = move || { + let attempts = attempts2.clone(); + async move { + attempts.fetch_add(1, Ordering::SeqCst); + None:: + } + }; + + let mut backoff_ms = crate::forwarder::RECONNECT_BASE_MS; + let result = + crate::forwarder::reconnect_with_retry(always_failing, &mut backoff_ms, &mut shutdown_rx) + .await; + + assert!(result.is_none(), "the flip aborts the retry"); + assert_eq!( + attempts.load(Ordering::SeqCst), + 0, + "no connect attempt on a shut-down forwarder" + ); +} + +/// Unit pin (server-less) — the fanout release guard: its drop empties +/// the shared sender slot, which is what surfaces the receiver's +/// terminal arm on *any* loop exit (the receivers-terminal-iff-loop- +/// gone invariant's mechanism). +#[test] +fn fanout_release_guard_empties_the_shared_slot_on_drop() { + let (sender, _) = + tokio::sync::broadcast::channel::(FANOUT_TEST_CAPACITY); + let slot = std::sync::Arc::new(std::sync::Mutex::new(Some(sender))); + let guard = crate::forwarder::FanoutRelease { slot: slot.clone() }; + assert!( + slot.lock().unwrap().is_some(), + "the guard holds the release until its own drop" + ); + drop(guard); + assert!( + slot.lock().unwrap().is_none(), + "the guard's drop released the fanout sender (any loop exit surfaces the terminal arm)" + ); +} + +const FANOUT_TEST_CAPACITY: usize = 1; + +/// 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 +/// backend kill forces the reconnect, and the loop survives *repeated* +/// failed connects across several backoff cycles (outage wider than +/// four doubling steps) — the skeleton's shape had the loop dead +/// after the first one, every later `listen()` failing with the +/// spurious "mid-reconnect" `Database` forever. Then the real config +/// comes back and full recovery is pinned: the synthetic +/// reconnect-wake broadcasts, a registered channel's post-reconnect +/// NOTIFY delivers, and `listen()` succeeds. +#[tokio::test(flavor = "multi_thread")] +async fn forwarder_survives_failed_reconnects_and_recovers_fully() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("failed-conn")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let forwarder = store.forwarder().clone(); + let admin = harness_client().await.unwrap(); + + let channel = format!("fconn_probe_{}", instance_namer("ch")()); + let mut rx = forwarder.subscribe().unwrap(); + forwarder.register(&channel).await.unwrap(); + + // Pre-outage delivery verified (the normal path is not what broke). + let notify_sql = |payload: &str| { + format!( + "NOTIFY {}, '{payload}'", + crate::schema::quote_identifier(&channel) + ) + }; + admin + .batch_execute(¬ify_sql("pre-outage")) + .await + .unwrap(); + assert!( + wait_for_channel(&mut rx, &channel, Duration::from_secs(3)).await, + "pre-outage delivery must work first" + ); + + // The seam: reconnects now point at an unreachable endpoint (the + // open-time config was fine — the outage is post-open). Then the + // backend kill forces the reconnect the skeleton never survived. + let real_config = forwarder.reconnect_config_snapshot(); + let outage_config: tokio_postgres::Config = unreachable_dsn().parse().unwrap(); + forwarder.swap_reconnect_config(outage_config); + let killed = admin + .query_one( + "SELECT pg_terminate_backend(pid) FROM pg_stat_activity + WHERE application_name = $1 AND pid <> pg_backend_pid()", + &[&store.listener_application_name()], + ) + .await + .unwrap() + .get::<_, bool>(0); + assert!(killed, "the listener backend was found and killed"); + + // Mid-outage: the loop is alive but disconnected — `register` + // fails with the honest transient `Database` (the mid-reconnect + // posture, registry entry kept for the reconnect re-issue). + tokio::time::sleep(Duration::from_millis(250)).await; + let outage_channel = instance_namer("outage-ch")(); + let err = forwarder.register(&outage_channel).await.unwrap_err(); + assert!( + matches!(err, alkstore::Error::Database(_)), + "mid-outage registration fails with the honest transient, got: {err:?}" + ); + + // The outage window: wider than four backoff doubling steps + // (attempts at 50/150/350/750 ms with the 50 ms base) — failed + // connects repeated across several cycles; the dead-loop shape + // would have stopped after the first attempt. + tokio::time::sleep(Duration::from_millis(1450)).await; + + // Restore: the next backoff attempt connects, and recovery is + // observable end-to-end. + forwarder.swap_reconnect_config(real_config); + let got_wake = wait_for_channel( + &mut rx, + alkstore::RESERVED_LISTENER_RECONNECTED, + Duration::from_secs(10), + ) + .await; + let mut delivered = false; + for _ in 0..20 { + let _ = admin.batch_execute(¬ify_sql("post-recovery")).await; + if wait_for_channel(&mut rx, &channel, Duration::from_millis(150)).await { + delivered = true; + break; + } + } + assert!( + got_wake && delivered, + "the forwarder must survive repeated failed connects and recover \ + (reconnect-wake {got_wake}, post-recovery delivery {delivered})" + ); + + // The trait surface is back: listen() succeeds post-recovery + // (never the spurious forever-"mid-reconnect" `Database`). + let recovery_channel = instance_namer("recovery-ch")(); + let mut receiver = store.listen(&recovery_channel).await.unwrap(); + assert!( + receiver + .recv_timeout(Duration::from_millis(100)) + .await + .unwrap() + .is_none(), + "the recovered listener idles healthily" + ); + + store.close(); + drop_schema(&admin, &schema).await; +} diff --git a/tasks/pg-fix-forwarder-reconnect.md b/tasks/pg-fix-forwarder-reconnect.md index 2904585..3acfc29 100644 --- a/tasks/pg-fix-forwarder-reconnect.md +++ b/tasks/pg-fix-forwarder-reconnect.md @@ -1,7 +1,7 @@ --- id: pg-fix-forwarder-reconnect name: Fix forwarder permanent death after one failed reconnect (review 002 Finding 1) -status: pending +status: completed depends_on: [] scope: moderate risk: medium @@ -103,8 +103,92 @@ test-exercised, not reasoned-about. ## Notes -> To be filled by implementation agent +> Decisions of record the implementation made that the description +> didn't pin: + +- **The retry core is an extracted helper (`reconnect_with_retry`)** + rather than an inline restructure: sleep → shutdown-select → attempt + → repeat, generic over an `FnMut() -> Future>` attempt closure. It is only ever exited + with a fresh generation pair (`Some`) or at the shutdown flip + (`None`) — the loop top's empty-slot arm is restructured away + entirely: `LoopCtx.connection` is no longer `Option`; the loop + re-binds a plain loop-carried connection from the successful connect + at the arm's bottom, so there is no fall-through path left to die + on. +- **The attempt closure is the connector injection** (the + implementer's-choice seam, one mechanism for both shapes): the real + closure clones `tokio_postgres::Config` out of a shared + `ReconnectConfigSlot` (`Arc>`) and calls + `Config::connect(NoTls)`, mapping errors to `None`; the server-less + unit tests substitute their own always-failing closures. The slot is + also the config-swap seam: `Forwarder::swap_reconnect_config` / + `reconnect_config_snapshot`(`#[cfg(test)]`) let a test point + reconnects at an unreachable endpoint post-open and restore the + real one. The field is `#[cfg(test)]` on `Forwarder` (dead in prod + otherwise — the loop reads its own slot clone); the shutdown + watch's `changed()` is selected *with* the backoff sleep so a close + surfaces the terminal arm mid-outage instead of waiting out a cycle. +- **The fanout hardening rides a shared slot + drop guard**: the + handle's sender became the shared `FanoutSlot` (`Arc>>`) the loop also receives, and the loop holds a + `FanoutRelease` guard whose `Drop` empties the slot — Drop runs on + return, unwind, and task abort alike, so *any* loop exit (plus + `Forwarder::shutdown`'s take, unchanged) releases the fanout. + `subscribe` reads the shared slot — a subscribe arriving past any + loop exit fails closed (`None`) rather than parking on a + forever-silent broadcast. receivers-terminal-iff-loop-gone is now + enforced by construction, not just the shutdown path. +- **Backoff reset semantics unchanged**: the cap/growth live in the + helper (doubled per attempt, `saturating_mul(2)` → `min(2000)`); + the reset to the 50 ms base still fires on the *verified-alive* + generation (successful LISTEN/probe), so a connect that succeeds + but dies before its LISTEN reconnects at the carried backoff. +- **The success path of the retry helper has no server-less unit pin** + (a real `(Client, Connection)` pair cannot be fabricated); it is + pinned end-to-end by the harness recovery test. The unit pins carry: + repeated failed connects across six backoff cycles with the cap, + pre-flipped shutdown aborts with zero attempts, and the guard's + drop-releases-slot mechanism. +- `tokio` gained a `test-util` feature on the crate's + `[dev-dependencies]` (feature-unified into test builds only) for the + paused-time retry unit test — the gate battery runs deterministic + and instant. +- The bug was replay-proofed: with the old one-shot `Err → return` + shape temporarily reintroduced, the new harness test + (`forwarder_survives_failed_reconnects_and_recovers_fully`) fails as + expected; with the fix it passes (verified both ways live). ## Summary -> To be filled on completion \ No newline at end of file +> What landed, verified how: + +- **`alkstore-postgres/src/forwarder.rs`**: the forwarder loop + restructured — the reconnect arm now carries + `reconnect_with_retry` (sleep with live backoff → shutdown-arms → + connect attempt, repeating until success or shutdown; failed + connects never leave the arm), the loop-carried connection is a + plain non-`Option` rebinding (the empty-slot `take() else return` + that killed the loop permanently after one failed connect is gone by + construction), and the shared `FanoutSlot` + `FanoutRelease` guard + make every loop exit path release the fanout sender + (receivers-terminal-iff-loop-gone). `LoopCtx`/`Forwarder::spawn` + adjusted; the `cfg(test)` reconnect-config seam (swap/snapshot) + added; subscribe/shutdown doc comments updated for the invariant. +- **`alkstore-postgres/src/store/open_tests.rs`** (4 new tests): + `forwarder_survives_failed_reconnects_and_recovers_fully` (harness: + config swapped to unreachable → backend kill → repeated failed + connects across ~4 backoff cycles mid-outage with the honest + transient `register` error pinned → config restored → reconnect-wake + + post-recovery NOTIFY delivery + working `listen()`), two + server-less paused-time unit pins of the retry core (six failed + attempts across six backoff cycles capped at 2 s, shutdown-abort + before any attempt), and the fanout release guard's drop mechanism. +- **`alkstore-postgres/Cargo.toml`**: `tokio` `test-util` added to + `[dev-dependencies]` only. +- **Gates**: `cargo test -p alkstore-postgres` with the harness server + green (115 lib + 10 contract-suite + 9 schema, run twice full); + every pre-existing forwarder/reconnect/no-replay/shutdown pin green + unchanged; workspace `cargo build`, `cargo test` (server-less — new + tests skip or run server-less cleanly), `cargo clippy --all-targets + -- -D warnings`, `cargo fmt --check` all green. \ No newline at end of file