From 5014b870d6608c02d020f28c05ae15d720102a71 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 20 Aug 2026 08:57:19 +0000 Subject: [PATCH] =?UTF-8?q?feat(conn=5Fregistry):=20registry=20owns=20the?= =?UTF-8?q?=20drain=20=E2=80=94=20trapping=20gen=5Fserver,=20request=5Fshu?= =?UTF-8?q?tdown=20driven?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The drain protocol moves wholesale into the registry (smarm >=0.7: trap_exit + handle_shutdown + gen_server timers + StopHandle): - handle_shutdown: flip draining, stop idle conns, arm a drain_timeout timer, Continue; with no conns, Exit immediately. - ConnIdle while draining stops the conn (unchanged); ConnStarted while draining stops it on arrival (unchanged); ConnEnded that empties the set while draining = StopHandle::stop() — the registry's own normal exit is now the 'every connection is gone' barrier. - handle_timer (deadline): one force-stop sweep. The old re-sweep-every- 10ms loop existed to catch late registrants; stop-on-arrival already covers every post-sweep entry path, so one sweep suffices. - Cast::{BeginDrain,ForceStopConns} deleted (internal now); ConnCount stays as the introspection call. start() takes drain_timeout. - serve.rs: the root's whole drain/poll block collapses to registry.shutdown() (graceful, monitors until the server has stopped itself). SHUTDOWN_POLL + the listener flag are untouched here; they go with the endpoint refactor. - Known residual window documented in the module docs: a conn spawned by a dying listener that has not yet registered can outlive an already-empty registry; it is collected by smarm's root-exit sweep. Tests (in-lib, request_shutdown driven): empty-set immediate exit; idle-now/busy-at-deadline ordering with exit-after; stop-on-idle mid-drain; stop-on-arrival mid-drain. 35x hammer subset green. --- src/conn_registry.rs | 322 +++++++++++++++++++++++++++++++++++-------- src/serve.rs | 41 ++---- 2 files changed, 276 insertions(+), 87 deletions(-) diff --git a/src/conn_registry.rs b/src/conn_registry.rs index 267644e..eb7f70b 100644 --- a/src/conn_registry.rs +++ b/src/conn_registry.rs @@ -1,41 +1,49 @@ -//! Connection registry — the shutdown coordinator (v0.2 chunk 2). +//! Connection registry — the shutdown coordinator, now a trapping +//! `gen_server` that owns its own drain (v0.3). //! -//! A `gen_server` tracking live connection pids and their busy/idle -//! state. Connection actors self-register as their first action and -//! self-deregister via a drop guard (so panic unwinds deregister too); -//! both casts come from the same sender, so Started always precedes -//! Ended in the inbox. (The roadmap sketched the *listener* casting -//! `{Started, pid}`, but then a short-lived conn's Ended could overtake -//! its Started and leak a dead pid into the set forever. Self- -//! registration makes the order a per-sender FIFO guarantee instead of -//! a race.) +//! Tracks live connection pids and their busy/idle state. Connection +//! actors self-register as their first action and self-deregister via a +//! drop guard (so panic unwinds deregister too); both casts come from the +//! same sender, so Started always precedes Ended in the inbox. (The +//! roadmap sketched the *listener* casting `{Started, pid}`, but then a +//! short-lived conn's Ended could overtake its Started and leak a dead +//! pid into the set forever. Self-registration makes the order a +//! per-sender FIFO guarantee instead of a race.) //! -//! Listeners are deliberately NOT tracked here: they shut down via a -//! shared flag + timed accept-waits (see `serve`), never via -//! `request_stop` — stopping pids that announce themselves is racy (the -//! spawn-to-registration gap), and smarm's `request_stop` is lossy -//! against an actor that is QUEUED and then parks without passing an -//! observation point (see the v0.2 shutdown notes in the commit -//! message). Connections don't suffer this in practice: every stop the -//! registry issues targets a pid that just sent us a cast (so it is -//! running or parked, both covered), and the force-stop path re-sweeps -//! until the set empties. +//! Drain protocol (all internal since v0.3): a `request_shutdown` +//! reaching this server lands in `handle_shutdown`, which stops every +//! idle connection immediately and flips into draining mode — any +//! connection that *becomes* idle (finishes its in-flight request) is +//! stopped on the spot, and a connection that registers mid-drain (the +//! accept-just-before-listener-death race) is stopped on arrival. Busy +//! connections are left to finish up to `drain_timeout`, at which point a +//! timer fire force-stops whatever remains. The single force sweep +//! suffices: every path a conn can enter the set on after it (`ConnStarted` +//! while draining) is stopped at the point of entry. When the set empties +//! the server stops itself normally — under a supervisor with +//! `Shutdown::Infinity`, "the registry has exited" *is* the barrier "every +//! connection is gone". //! -//! Drain protocol: `BeginDrain` stops every idle connection immediately -//! and flips the registry into draining mode, in which any connection -//! that *becomes* idle (finishes its in-flight request) is stopped on the -//! spot. Busy connections are left to finish; `ForceStopConns` (sent by -//! `serve` at the drain deadline) stops whatever remains. `request_stop` -//! unwinds a conn actor parked in `wait_readable` safely (smarm's 06-10 -//! io fix) and `OwnedFd::drop` closes its socket on the way out. +//! `request_stop` unwinds a conn actor parked in `wait_readable` safely +//! (smarm's 06-10 io fix) and `OwnedFd::drop` closes its socket on the +//! way out. Every stop the registry issues targets a pid that just sent +//! it a cast (so it is running or parked, both covered). +//! +//! Known residual window (documented, not defended): a conn *spawned* by +//! a dying listener but not yet run has not registered; if the set was +//! already empty the registry may exit before that conn's `ConnStarted` +//! arrives, and the conn serves on unsupervised until the root-exit sweep +//! collects it. Inherent to self-registration; accepted for v0.3. //! //! This server is also the planned introspection point for ws/channels //! (roadmap v0.4+), which is why it exists as its own module rather than //! being inlined into `serve`. -use smarm::{GenServer, Pid, GenServerBuilder, GenServerRef}; +use smarm::gen_server::{ShutdownAction, StopHandle, TimerHandle}; +use smarm::{GenServer, GenServerBuilder, GenServerCtx, GenServerRef, Pid}; use std::collections::HashMap; +use std::time::Duration; // --------------------------------------------------------------------------- // Messages @@ -51,10 +59,6 @@ pub enum Cast { /// keep-alive request. ConnIdle(Pid), ConnEnded(Pid), - /// Stop idle conns now and stop each remaining conn as it goes idle. - BeginDrain, - /// Drain deadline passed: stop every remaining conn. - ForceStopConns, } pub enum Call { @@ -75,19 +79,51 @@ enum ConnState { Idle, } -#[derive(Default)] pub struct ConnRegistry { - conns: HashMap, + conns: HashMap, draining: bool, + drain_timeout: Duration, + stop: Option>, + timer: Option>, +} + +impl ConnRegistry { + pub fn new(drain_timeout: Duration) -> Self { + Self { + conns: HashMap::new(), + draining: false, + drain_timeout, + stop: None, + timer: None, + } + } + + /// Normal self-exit once draining and empty — the "every connection is + /// gone" barrier the supervisor's ordered shutdown waits on. + fn stop_if_drained(&self) { + if self.draining && self.conns.is_empty() { + self.stop + .as_ref() + .expect("init ran before any message") + .stop(); + } + } } impl GenServer for ConnRegistry { - type Call = Call; + type Call = Call; type Reply = Reply; - type Cast = Cast; - type Info = (); + type Cast = Cast; + type Info = (); + /// One meaning: the drain deadline elapsed — force-stop the stragglers. type Timer = (); + fn init(&mut self, ctx: &GenServerCtx) { + ctx.trap_exit(); // shutdown arrives as handle_shutdown, not a kill + self.stop = Some(ctx.stop_handle()); + self.timer = Some(ctx.timer()); + } + fn handle_call(&mut self, request: Call) -> Reply { match request { Call::ConnCount => Reply::ConnCount(self.conns.len()), @@ -100,7 +136,8 @@ impl GenServer for ConnRegistry { self.conns.insert(pid, ConnState::Idle); if self.draining { // Listener-stop race: this conn was accepted just - // before its listener died. Drain means no new work. + // before its listener died. Drain means no new work; + // it leaves the map via its guard's ConnEnded. smarm::request_stop(pid); } } @@ -114,32 +151,52 @@ impl GenServer for ConnRegistry { *s = ConnState::Idle; if self.draining { // Finished its in-flight request; nothing more is - // owed. The pid leaves the map via its drop - // guard's ConnEnded once the unwind completes. + // owed. smarm::request_stop(pid); } } } - Cast::ConnEnded(pid) => { self.conns.remove(&pid); } - Cast::BeginDrain => { - self.draining = true; - for (pid, state) in &self.conns { - if *state == ConnState::Idle { - smarm::request_stop(*pid); - } - } - } - Cast::ForceStopConns => { - for pid in self.conns.keys() { - smarm::request_stop(*pid); - } + Cast::ConnEnded(pid) => { + self.conns.remove(&pid); + self.stop_if_drained(); } } } + + fn handle_shutdown(&mut self) -> ShutdownAction { + self.draining = true; + if self.conns.is_empty() { + return ShutdownAction::Exit; + } + for (pid, state) in &self.conns { + if *state == ConnState::Idle { + smarm::request_stop(*pid); + } + } + // Busy conns get until the deadline; then handle_timer sweeps. + self.timer + .as_ref() + .expect("init ran before any message") + .arm_after(self.drain_timeout, ()); + ShutdownAction::Continue + } + + fn handle_timer(&mut self, _deadline: ()) { + // Drain deadline passed: force-stop every remaining conn. One + // sweep — see the module docs for why late registrants are + // already covered by the stop-on-arrival in ConnStarted. + for pid in self.conns.keys() { + smarm::request_stop(*pid); + } + // The map empties via each conn's guard ConnEnded; stop_if_drained + // fires on the last one. + } } -pub fn start() -> GenServerRef { - GenServerBuilder::new(ConnRegistry::default()).start() +/// Start an anonymous, ref-addressed registry (the pre-v0.3 shape, still +/// used by `serve` until the endpoint refactor lands). +pub fn start(drain_timeout: Duration) -> GenServerRef { + GenServerBuilder::new(ConnRegistry::new(drain_timeout)).start() } // --------------------------------------------------------------------------- @@ -151,13 +208,17 @@ pub fn start() -> GenServerRef { /// guard's perspective (a dead registry just returns an ignored Err). pub struct DeregisterGuard { registry: GenServerRef, - pid: Pid, - make: fn(Pid) -> Cast, + pid: Pid, + make: fn(Pid) -> Cast, } impl DeregisterGuard { pub fn new(registry: GenServerRef, pid: Pid, make: fn(Pid) -> Cast) -> Self { - Self { registry, pid, make } + Self { + registry, + pid, + make, + } } } @@ -166,3 +227,148 @@ impl Drop for DeregisterGuard { let _ = self.registry.cast((self.make)(self.pid)); } } + +// --------------------------------------------------------------------------- +// Tests — the drain protocol, driven purely by request_shutdown. +// --------------------------------------------------------------------------- + +#[cfg(test)] +mod tests { + use super::*; + use std::time::Instant; + + /// A stand-in connection actor: registers, optionally reports busy, + /// then parks forever. Only a `request_stop` ends it; the guard's + /// ConnEnded runs on the unwind. + fn fake_conn(registry: GenServerRef, busy: bool) { + let me = smarm::self_pid(); + let _ = registry.cast(Cast::ConnStarted(me)); + let _guard = DeregisterGuard::new(registry.clone(), me, Cast::ConnEnded); + if busy { + let _ = registry.cast(Cast::ConnBusy(me)); + } + loop { + smarm::sleep(Duration::from_secs(3600)); + } + } + + fn conn_count(reg: &GenServerRef) -> Result { + match reg.call(Call::ConnCount) { + Ok(Reply::ConnCount(n)) => Ok(n), + Err(_) => Err(()), + } + } + + /// Poll until `pred` holds or the deadline passes; panics with `what` + /// on timeout. Registry calls are the sync point, so no sleeps are + /// load-bearing for correctness — only for latency. + fn await_pred(what: &str, deadline: Duration, mut pred: impl FnMut() -> bool) { + let end = Instant::now() + deadline; + while !pred() { + assert!(Instant::now() < end, "timed out waiting for: {what}"); + smarm::sleep(Duration::from_millis(5)); + } + } + + #[test] + fn shutdown_with_no_conns_exits_immediately() { + smarm::run(|| { + let reg = start(Duration::from_secs(5)); + assert_eq!(conn_count(®), Ok(0)); + smarm::request_shutdown(reg.pid()); + // handle_shutdown returns Exit; the server is gone shortly. + await_pred("registry exit", Duration::from_secs(2), || { + conn_count(®).is_err() + }); + }); + } + + #[test] + fn drain_stops_idle_now_busy_at_deadline_then_exits() { + smarm::run(|| { + let drain = Duration::from_millis(300); + let reg = start(drain); + smarm::spawn({ + let r = reg.clone(); + move || fake_conn(r, false) // idle + }); + smarm::spawn({ + let r = reg.clone(); + move || fake_conn(r, true) // busy + }); + await_pred("both registered", Duration::from_secs(2), || { + conn_count(®) == Ok(2) + }); + + let t0 = Instant::now(); + smarm::request_shutdown(reg.pid()); + + // Idle conn goes promptly, well before the deadline. + await_pred("idle stopped", drain / 2, || conn_count(®) == Ok(1)); + + // Busy conn holds until the force sweep, then everything — + // registry included — winds down. + await_pred("busy swept + registry exit", drain * 4, || { + conn_count(®).is_err() + }); + assert!( + t0.elapsed() >= drain, + "busy conn must not be stopped before the drain deadline" + ); + }); + } + + #[test] + fn conn_finishing_mid_drain_is_stopped_on_idle() { + smarm::run(|| { + // Long deadline: the test passes only if the ConnIdle path + // stops the conn, not the force sweep. + let reg = start(Duration::from_secs(30)); + let conn = smarm::spawn({ + let r = reg.clone(); + move || fake_conn(r, true) + }); + await_pred("registered busy", Duration::from_secs(2), || { + conn_count(®) == Ok(1) + }); + smarm::request_shutdown(reg.pid()); + // Still busy: must survive the immediate idle sweep. + smarm::sleep(Duration::from_millis(50)); + assert_eq!(conn_count(®), Ok(1)); + // "Request finishes": the conn reports idle. + let _ = reg.cast(Cast::ConnIdle(conn.pid())); + await_pred( + "stopped on idle + registry exit", + Duration::from_secs(2), + || conn_count(®).is_err(), + ); + }); + } + + #[test] + fn conn_registering_mid_drain_is_stopped_on_arrival() { + smarm::run(|| { + let reg = start(Duration::from_secs(30)); + let holder = smarm::spawn({ + let r = reg.clone(); + move || fake_conn(r, true) // keeps the drain open + }); + await_pred("holder registered", Duration::from_secs(2), || { + conn_count(®) == Ok(1) + }); + smarm::request_shutdown(reg.pid()); + smarm::sleep(Duration::from_millis(20)); + // The listener-death race: a fresh conn registers mid-drain. + smarm::spawn({ + let r = reg.clone(); + move || fake_conn(r, false) + }); + // It is stopped on arrival: count returns to exactly the holder. + await_pred("late arrival stopped", Duration::from_secs(2), || { + conn_count(®) == Ok(1) + }); + // Cleanup: release the holder so the run can end. + let _ = reg.cast(Cast::ConnIdle(holder.pid())); + }); + } +} diff --git a/src/serve.rs b/src/serve.rs index b06e12e..ae0b081 100644 --- a/src/serve.rs +++ b/src/serve.rs @@ -13,7 +13,7 @@ //! waiting on "the same fd"). use crate::conn_actor::{run_connection, ConnLimits}; -use crate::conn_registry::{self, Call, Cast, ConnRegistry, Reply}; +use crate::conn_registry::{self, ConnRegistry}; use crate::net::{accept_nonblocking, bind_and_listen, OwnedFd}; use crate::plug::Pipeline; @@ -439,7 +439,7 @@ pub fn serve_with_shutdown( rt.run(move || { // Registry first: listeners and conns cast into it from birth. - let registry = conn_registry::start(); + let registry = conn_registry::start(drain_timeout); let mut sup = OneForOne::new().strategy(Strategy::OneForOne); for (i, lfd) in listener_fds.into_iter().enumerate() { @@ -488,34 +488,17 @@ pub fn serve_with_shutdown( shutdown_flag.store(true, Ordering::Relaxed); let _ = sup_h.join(); - // 3 + 4. Drain. Same sweep discipline as listeners on the force- - // stop path: a conn accepted just before its listener died may - // register after the deadline, so keep force-stopping until the - // set is empty (each pass kills everything registered; new - // registrants are a strictly shrinking population once listeners - // are gone). - let _ = registry.cast(Cast::BeginDrain); - let deadline = std::time::Instant::now() + drain_timeout; - let mut force = false; - loop { - match registry.call(Call::ConnCount) { - Ok(Reply::ConnCount(0)) => break, - Ok(_) => {} - Err(_) => break, // registry gone; nothing left to track - } - let now = std::time::Instant::now(); - if force || now >= deadline { - force = true; - let _ = registry.cast(Cast::ForceStopConns); - smarm::sleep(Duration::from_millis(10)); - } else { - smarm::sleep(Duration::from_millis(50).min(deadline - now)); - } - } + // 3 + 4. Drain: the registry owns the whole protocol now (idle + // stopped immediately, busy until its internal drain_timeout + // timer, late registrants stopped on arrival — see + // conn_registry docs). `GenServerRef::shutdown()` delivers the + // request and blocks on a monitor until the registry has + // stopped itself, which it does only once the conn set is + // empty: this line IS the "every connection is gone" barrier. + registry.shutdown(); - // 5. Our GenServerRef drops here. The registry's inbox closes once - // the last conn's clone drops with it, and the runtime winds - // down when the last actor exits. + // 5. Root returns; the runtime winds down when the last actor + // exits. }); Ok(())