pg forwarder: the reconnect arm retries failed connects until success or shutdown — one failed connect no longer kills the loop permanently; every loop exit path releases the fanout sender via the shared-slot drop guard (receivers-terminal-iff-loop-gone); reconnect-config test seam + failed-connect pins: six-cycle exhausted-backoff server-less unit, pre-flip abort, guard drop, and the harness end-to-end unreachable-outage → full-recovery test (bug replay-proven against the old shape) (task pg-fix-forwarder-reconnect, review 002 Finding 1)
This commit is contained in:
1 parent
e067357de9
commit
7ae426a01d
4 files changed
+481
-57
No files matched your search
@@ -15,3 +15,4 @@ tokio-postgres = "0.7"
|
|||||||
|
|
||||||
[dev-dependencies]
|
[dev-dependencies]
|
||||||
alkstore-contract-suite = { path = "../alkstore-contract-suite" }
|
alkstore-contract-suite = { path = "../alkstore-contract-suite" }
|
||||||
|
tokio = { version = "1", features = ["test-util"] }
|
||||||
@@ -41,6 +41,7 @@
|
|||||||
//! per-receiver bridging are the notify-listen task's; the skeleton
|
//! per-receiver bridging are the notify-listen task's; the skeleton
|
||||||
//! broadcasts raw notification channel names.
|
//! broadcasts raw notification channel names.
|
||||||
|
|
||||||
|
use std::future::Future;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::sync::atomic::{AtomicBool, Ordering};
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
@@ -50,6 +51,20 @@ use tokio_postgres::NoTls;
|
|||||||
|
|
||||||
use crate::seam::database_error;
|
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<std::sync::Mutex<Option<broadcast::Sender<RawNotification>>>>;
|
||||||
|
|
||||||
|
/// 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<std::sync::Mutex<tokio_postgres::Config>>;
|
||||||
|
|
||||||
|
/// 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 listener connection's concrete type for the `NoTls` posture
|
||||||
/// (the engine's TLS story rides the consumer's `Config` sslmode for
|
/// (the engine's TLS story rides the consumer's `Config` sslmode for
|
||||||
/// the pooled path; the dedicated listener connects with `NoTls` —
|
/// the pooled path; the dedicated listener connects with `NoTls` —
|
||||||
@@ -154,11 +169,13 @@ impl ChannelSet {
|
|||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub(crate) struct Forwarder {
|
pub(crate) struct Forwarder {
|
||||||
/// The wake fanout — every receiver bridge holds a subscription;
|
/// The wake fanout — every receiver bridge holds a subscription;
|
||||||
/// its release (this handle's take at shutdown, the loop's drop at
|
/// its release (the shutdown take, or the loop's exit-path release
|
||||||
/// its exit) is the receivers' terminal `Closed`/`None` arm.
|
/// guard) is the receivers' terminal `Closed`/`None` arm. The
|
||||||
/// `Option` because shutdown empties it while receiver-held
|
/// shared [`FanoutSlot`] exactly because receiver-held
|
||||||
/// `Arc<Forwarder>`s keep the struct alive.
|
/// `Arc<Forwarder>`s keep this struct alive — the release fires
|
||||||
fanout: std::sync::Mutex<Option<broadcast::Sender<RawNotification>>>,
|
/// 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
|
/// The channel registry the reconnect re-issues from; read by the
|
||||||
/// register/unregister surfaces below.
|
/// register/unregister surfaces below.
|
||||||
channels: Arc<ChannelSet>,
|
channels: Arc<ChannelSet>,
|
||||||
@@ -172,6 +189,12 @@ pub(crate) struct Forwarder {
|
|||||||
/// gap fails with `Database`; the registry entry's recovery rides
|
/// gap fails with `Database`; the registry entry's recovery rides
|
||||||
/// the reconnect re-issue). Read by `register` below.
|
/// the reconnect re-issue). Read by `register` below.
|
||||||
connected: Arc<AtomicBool>,
|
connected: Arc<AtomicBool>,
|
||||||
|
/// 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<bool>,
|
shutdown: watch::Sender<bool>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -201,9 +224,12 @@ impl Forwarder {
|
|||||||
reconnect_config: tokio_postgres::Config,
|
reconnect_config: tokio_postgres::Config,
|
||||||
) -> Forwarder {
|
) -> Forwarder {
|
||||||
let (fanout, _) = broadcast::channel(FANOUT_CAPACITY);
|
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 channels = Arc::new(ChannelSet::default());
|
||||||
let (commands_tx, commands_rx) = mpsc::unbounded_channel();
|
let (commands_tx, commands_rx) = mpsc::unbounded_channel();
|
||||||
let (shutdown_tx, shutdown_rx) = watch::channel(false);
|
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
|
// The connection handed in was just established — the liveness
|
||||||
// flag starts true and the loop clears it at each generation's
|
// flag starts true and the loop clears it at each generation's
|
||||||
// death (the register path's honest transient state).
|
// death (the register path's honest transient state).
|
||||||
@@ -214,9 +240,10 @@ impl Forwarder {
|
|||||||
// lifetime).
|
// lifetime).
|
||||||
tokio::spawn(forwarder_loop(LoopCtx {
|
tokio::spawn(forwarder_loop(LoopCtx {
|
||||||
client,
|
client,
|
||||||
connection: Some(connection),
|
connection,
|
||||||
reconnect_config,
|
reconnect_config: reconnect_config_slot.clone(),
|
||||||
fanout: fanout.clone(),
|
fanout: fanout.clone(),
|
||||||
|
fanout_slot: fanout_slot.clone(),
|
||||||
channels: channels.clone(),
|
channels: channels.clone(),
|
||||||
commands: commands_rx,
|
commands: commands_rx,
|
||||||
connected: connected.clone(),
|
connected: connected.clone(),
|
||||||
@@ -224,18 +251,22 @@ impl Forwarder {
|
|||||||
}));
|
}));
|
||||||
|
|
||||||
Forwarder {
|
Forwarder {
|
||||||
fanout: std::sync::Mutex::new(Some(fanout)),
|
fanout: fanout_slot.clone(),
|
||||||
channels,
|
channels,
|
||||||
commands_tx,
|
commands_tx,
|
||||||
connected,
|
connected,
|
||||||
|
#[cfg(test)]
|
||||||
|
reconnect_config: reconnect_config_slot,
|
||||||
shutdown: shutdown_tx,
|
shutdown: shutdown_tx,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Subscribe to the raw fanout (the notify task bridges broadcast
|
/// Subscribe to the raw fanout (the notify task bridges broadcast
|
||||||
/// receivers into `WakeReceiver`s). `None` after shutdown (the
|
/// receivers into `WakeReceiver`s). `None` after shutdown — or
|
||||||
/// senders are being released) — a subscribe racing `close()` or
|
/// after *any* forwarder-loop exit (the loop's release guard takes
|
||||||
/// arriving past it fails closed rather than panicking.
|
/// 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<broadcast::Receiver<RawNotification>> {
|
pub(crate) fn subscribe(&self) -> Option<broadcast::Receiver<RawNotification>> {
|
||||||
self.fanout
|
self.fanout
|
||||||
.lock()
|
.lock()
|
||||||
@@ -302,9 +333,10 @@ impl Forwarder {
|
|||||||
/// Shutdown: flip the switch — the loop exits, dropping the
|
/// Shutdown: flip the switch — the loop exits, dropping the
|
||||||
/// listener `Client` and its session (the listener connection is
|
/// listener `Client` and its session (the listener connection is
|
||||||
/// gone) — and release this handle's fanout sender. The broadcast
|
/// gone) — and release this handle's fanout sender. The broadcast
|
||||||
/// closes once every sender is dropped (the loop's clone dies with
|
/// closes once every sender is dropped (the loop's exit-path
|
||||||
/// its exit; this handle's dies here — receivers unaffected by
|
/// release guard empties the shared slot too; this handle's take
|
||||||
/// `Arc<Forwarder>` lifetimes see the terminal arm). Idempotent.
|
/// fires the terminal arm at the flip — receivers unaffected by
|
||||||
|
/// `Arc<Forwarder>` lifetimes see it). Idempotent.
|
||||||
pub(crate) fn shutdown(&self) {
|
pub(crate) fn shutdown(&self) {
|
||||||
let _ = self.shutdown.send(true);
|
let _ = self.shutdown.send(true);
|
||||||
// Take the sender even if a receiver holds this Arc: the
|
// Take the sender even if a receiver holds this Arc: the
|
||||||
@@ -312,32 +344,114 @@ impl Forwarder {
|
|||||||
// whatever lifetime stands on the handle.
|
// whatever lifetime stands on the handle.
|
||||||
*self.fanout.lock().expect("fanout mutex poisoned") = None;
|
*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<F, Fut>(
|
||||||
|
mut attempt: F,
|
||||||
|
backoff_ms: &mut u64,
|
||||||
|
shutdown: &mut watch::Receiver<bool>,
|
||||||
|
) -> Option<GenerationPair>
|
||||||
|
where
|
||||||
|
F: FnMut() -> Fut,
|
||||||
|
Fut: Future<Output = Option<GenerationPair>>,
|
||||||
|
{
|
||||||
|
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
|
/// 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
|
/// Per connection generation: spawn the poll task (it owns the
|
||||||
/// `Connection` and fans out `poll_message`'s notifications) → issue
|
/// `Connection` and fans out `poll_message`'s notifications) → issue
|
||||||
/// LISTENs from the live channel-set snapshot → serve commands and
|
/// LISTENs from the live channel-set snapshot → serve commands and
|
||||||
/// wait for the connection's death → reconnect with exponential
|
/// wait for the connection's death → reconnect via
|
||||||
/// backoff (50 ms → 2 s cap). The loop consumes the handed-in
|
/// [`reconnect_with_retry`] (exponential backoff 50 ms → 2 s cap;
|
||||||
/// connection on its first generation and reconnects into the slot
|
/// failed connects never leave the arm). The loop-carried connection
|
||||||
/// afterwards.
|
/// 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
|
/// The receiver-facing close shape (ADR-021 §5's pg arm): the fanout
|
||||||
/// sender clones live here (loop) and on the handle (taken at
|
/// sender clones live here (loop) and on the shared slot (taken by the
|
||||||
/// shutdown); a connection death and its reconnect are invisible to
|
/// shutdown flip or the exit-path release guard — see
|
||||||
/// subscribers (the fanout never drops on a generation boundary —
|
/// [`FanoutRelease`]); a connection death and its reconnect are
|
||||||
/// receivers stay open across reconnects, they receive the reserved
|
/// invisible to subscribers (the fanout never drops on a generation
|
||||||
/// reconnect-wake instead); the senders drop only at shutdown — the
|
/// boundary — receivers stay open across reconnects, they receive the
|
||||||
/// loop's clone at its exit, the handle's at the shutdown take — the
|
/// reserved reconnect-wake instead); the senders drop only when the
|
||||||
/// receivers' terminal arm.
|
/// loop is gone — the invariant's shape.
|
||||||
struct LoopCtx {
|
struct LoopCtx {
|
||||||
client: tokio_postgres::Client,
|
client: tokio_postgres::Client,
|
||||||
connection: Option<ListenerConnection>,
|
connection: ListenerConnection,
|
||||||
reconnect_config: tokio_postgres::Config,
|
reconnect_config: ReconnectConfigSlot,
|
||||||
fanout: broadcast::Sender<RawNotification>,
|
fanout: broadcast::Sender<RawNotification>,
|
||||||
|
fanout_slot: FanoutSlot,
|
||||||
channels: Arc<ChannelSet>,
|
channels: Arc<ChannelSet>,
|
||||||
commands: mpsc::UnboundedReceiver<ListenCommand>,
|
commands: mpsc::UnboundedReceiver<ListenCommand>,
|
||||||
connected: Arc<AtomicBool>,
|
connected: Arc<AtomicBool>,
|
||||||
@@ -346,18 +460,22 @@ struct LoopCtx {
|
|||||||
|
|
||||||
async fn forwarder_loop(
|
async fn forwarder_loop(
|
||||||
LoopCtx {
|
LoopCtx {
|
||||||
mut client,
|
client,
|
||||||
mut connection,
|
connection,
|
||||||
reconnect_config,
|
reconnect_config,
|
||||||
fanout,
|
fanout,
|
||||||
|
fanout_slot,
|
||||||
channels,
|
channels,
|
||||||
mut commands,
|
mut commands,
|
||||||
connected,
|
connected,
|
||||||
mut shutdown,
|
mut shutdown,
|
||||||
}: LoopCtx,
|
}: LoopCtx,
|
||||||
) {
|
) {
|
||||||
|
let _fanout_release = FanoutRelease { slot: fanout_slot };
|
||||||
|
let mut client = client;
|
||||||
let mut backoff_ms = RECONNECT_BASE_MS;
|
let mut backoff_ms = RECONNECT_BASE_MS;
|
||||||
let mut first_listen = true;
|
let mut first_listen = true;
|
||||||
|
let mut connection = connection;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
// The dedicated poll loop, spawned BEFORE any client query can
|
// 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
|
// the poll task co-resident with client queries forever). The
|
||||||
// poll task owns the Connection; the loop drives queries
|
// poll task owns the Connection; the loop drives queries
|
||||||
// through the Client.
|
// through the Client.
|
||||||
let Some(poll_connection) = connection.take() else {
|
let mut poll_task = tokio::spawn(poll_loop(connection, fanout.clone()));
|
||||||
return;
|
|
||||||
};
|
|
||||||
let mut poll_task = tokio::spawn(poll_loop(poll_connection, fanout.clone()));
|
|
||||||
|
|
||||||
// Issue LISTENs from the dynamic channel set's snapshot (on
|
// Issue LISTENs from the dynamic channel set's snapshot (on
|
||||||
// first connect it carries pre-registered channels; on
|
// first connect it carries pre-registered channels; on
|
||||||
@@ -461,7 +576,8 @@ async fn forwarder_loop(
|
|||||||
// The command sender is gone (the store —
|
// The command sender is gone (the store —
|
||||||
// and its forwarder handle — went away
|
// and its forwarder handle — went away
|
||||||
// without a shutdown flip): serve no more,
|
// without a shutdown flip): serve no more,
|
||||||
// exit.
|
// exit (the release guard fires; the
|
||||||
|
// receivers' terminal arm).
|
||||||
None => {
|
None => {
|
||||||
poll_task.abort();
|
poll_task.abort();
|
||||||
return;
|
return;
|
||||||
@@ -486,27 +602,32 @@ async fn forwarder_loop(
|
|||||||
|
|
||||||
// Reconnect: the Client is dead with the session (dropping it
|
// Reconnect: the Client is dead with the session (dropping it
|
||||||
// here closes nothing server-side — pitfall 2 was satisfied
|
// here closes nothing server-side — pitfall 2 was satisfied
|
||||||
// for the connection's whole lifetime). Back off exponentially
|
// for the connection's whole lifetime). The retry core sleeps
|
||||||
// (50 ms → 2 s cap) and reconnect.
|
// 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);
|
connected.store(false, Ordering::Release);
|
||||||
tokio::time::sleep(Duration::from_millis(backoff_ms)).await;
|
let attempt = {
|
||||||
backoff_ms = backoff_ms.saturating_mul(2).min(RECONNECT_MAX_MS);
|
let slot = reconnect_config.clone();
|
||||||
|
move || {
|
||||||
if *shutdown.borrow() {
|
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;
|
return;
|
||||||
}
|
};
|
||||||
match reconnect_config.connect(NoTls).await {
|
client = next_client;
|
||||||
Ok((new_client, new_connection)) => {
|
connection = next_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;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -734,3 +734,221 @@ fn seam_mappings_preserve_the_source_chain() {
|
|||||||
let mapped = pool_error(deadpool_postgres::PoolError::Closed);
|
let mapped = pool_error(deadpool_postgres::PoolError::Closed);
|
||||||
assert!(matches!(mapped, alkstore::Error::Database(_)));
|
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::<crate::forwarder::GenerationPair>
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
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::<crate::forwarder::GenerationPair>
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
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::<crate::forwarder::RawNotification>(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;
|
||||||
|
}
|
||||||
@@ -1,7 +1,7 @@
|
|||||||
---
|
---
|
||||||
id: pg-fix-forwarder-reconnect
|
id: pg-fix-forwarder-reconnect
|
||||||
name: Fix forwarder permanent death after one failed reconnect (review 002 Finding 1)
|
name: Fix forwarder permanent death after one failed reconnect (review 002 Finding 1)
|
||||||
status: pending
|
status: completed
|
||||||
depends_on: []
|
depends_on: []
|
||||||
scope: moderate
|
scope: moderate
|
||||||
risk: medium
|
risk: medium
|
||||||
@@ -103,8 +103,92 @@ test-exercised, not reasoned-about.
|
|||||||
|
|
||||||
## Notes
|
## 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<Output =
|
||||||
|
Option<GenerationPair>>` 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<Mutex<Config>>`) 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<Mutex<Option<
|
||||||
|
Sender>>>`) 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
|
## Summary
|
||||||
|
|
||||||
> To be filled on completion
|
> 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.
|
||||||
Reference in new issue
Block a user