From 099b6bc3208b555d881aeaa91d83fb11c26bd8b2 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 20 Aug 2026 12:59:57 +0000 Subject: [PATCH] =?UTF-8?q?feat(endpoint):=20urus=20is=20a=20supervisable?= =?UTF-8?q?=20child=20=E2=80=94=20endpoint=20gen=5Fserver=20owns=20listene?= =?UTF-8?q?rs=20+=20conns?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The v0.3 shape from the spec: an app owns its runtime and root supervisor and places urus in it as one ordered child among its own. your root sup └── ChildSpec(Permanent, urus::endpoint(cfg, pipeline)?) <- Endpoint └── listener_sup OneForOne over N listeners └── plain connection actors - src/conn_registry.rs -> src/endpoint.rs. The registry gains the listener pool it registers for and becomes the Endpoint gen_server; ConnRegistry -> Endpoint. It runs inline as the ChildSpec's actor (NamedGenServerBuilder::run), so supervisor shutdown arrives as handle_shutdown and a restart re-runs init on the same still-open fds. - Endpoint spawns its OWN listener sup in init rather than being its sibling: smarm's supervisor start order is not start *readiness* (spawn is fire-and-forget), so a sibling listener could whereis the name before the registry actor ran. Registrar-spawns-consumers makes it program order inside one init. Gap filed in smarm ROADMAP (readiness ack); making spawn blocking would only shrink the window, not close it — 'has begun executing' is not 'has bound its name'. - Listener sup is monitored: death outside shutdown = panic (loud, the app's supervisor decides) instead of a zombie on a dead port. Death during shutdown is the 'no new connections' barrier. - DELETED: the shutdown AtomicBool, LISTENER_TICK (250ms wake per listener per tick, now an untimed wait_readable park), SHUTDOWN_POLL (100ms root poll — the root parks on the signal channel now), Restart::Transient (listeners are Permanent: they only exit by supervisor action, so a self-exit always means breakage). Verified against current smarm: request_stop unwinds an untimed wait_readable park, is no longer lossy against a QUEUED actor, and supervisor shutdown joins in ~200us. - Config: scheduler_threads/max_actors removed (runtime knobs an endpoint-as-child cannot honour) -> serve_with(cfg, smarm::Config, pipe) and serve_with_shutdown(cfg, smarm::Config, pipe, signal). Added Config.name (default 'urus'): the endpoint's registry name, unique per endpoint, and the introspection handle via endpoint::whereis(name). - serve* keep their meaning as the batteries-included path: they build a one-child tree around endpoint() with Shutdown::Infinity. Handle stays (a serve* caller has no RuntimeHandle to reach for) and now backs a real park instead of a poll. Tests: 4 drain tests ported onto a real endpoint (bound socket, supervised child, request_shutdown driven); new integration test boots an app tree with an ordered sibling and asserts serve-then-drain, reverse-order teardown and a closed port. 106 lib + 50 integration green. --- src/conn_actor.rs | 6 +- src/conn_registry.rs | 374 ------------------------ src/endpoint.rs | 656 +++++++++++++++++++++++++++++++++++++++++++ src/lib.rs | 3 +- src/serve.rs | 365 +++++------------------- tests/integration.rs | 125 +++++++-- 6 files changed, 840 insertions(+), 689 deletions(-) delete mode 100644 src/conn_registry.rs create mode 100644 src/endpoint.rs diff --git a/src/conn_actor.rs b/src/conn_actor.rs index 6ca1633..e7e95b9 100644 --- a/src/conn_actor.rs +++ b/src/conn_actor.rs @@ -14,7 +14,7 @@ //! during those parks, other connection actors progress freely. use crate::conn::{Body, Conn, HttpVersion, RespBody, StreamBody}; -use crate::conn_registry::{Cast, ConnRegistry, DeregisterGuard}; +use crate::endpoint::{Cast, DeregisterGuard, Endpoint}; use crate::net::OwnedFd; use crate::parser::{self, ParseError}; use crate::plug::Pipeline; @@ -115,7 +115,7 @@ pub fn run_connection( fd: OwnedFd, pipeline: Pipeline, limits: ConnLimits, - registry: GenServerRef, + registry: GenServerRef, ) { // The OwnedFd cleans up via Drop on any exit path (panic, error, or // normal close). No explicit close calls below. @@ -125,7 +125,7 @@ pub fn run_connection( // Self-register (initially idle: no request head parsed yet) and arm // the deregistration guard. Both casts come from this actor, so // Started always precedes Ended in the registry's inbox — see - // conn_registry module docs for why the listener must not do this. + // endpoint module docs for why the listener must not do this. let me = smarm::self_pid(); let _ = registry.cast(Cast::ConnStarted(me)); let _guard = DeregisterGuard::new(registry.clone(), me, Cast::ConnEnded); diff --git a/src/conn_registry.rs b/src/conn_registry.rs deleted file mode 100644 index eb7f70b..0000000 --- a/src/conn_registry.rs +++ /dev/null @@ -1,374 +0,0 @@ -//! Connection registry — the shutdown coordinator, now a trapping -//! `gen_server` that owns its own drain (v0.3). -//! -//! 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.) -//! -//! 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". -//! -//! `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::gen_server::{ShutdownAction, StopHandle, TimerHandle}; -use smarm::{GenServer, GenServerBuilder, GenServerCtx, GenServerRef, Pid}; - -use std::collections::HashMap; -use std::time::Duration; - -// --------------------------------------------------------------------------- -// Messages -// --------------------------------------------------------------------------- - -pub enum Cast { - /// A connection actor started (self-registered, initially idle: it has - /// not parsed a request head yet). - ConnStarted(Pid), - /// Parsed a request head; a response is now owed. - ConnBusy(Pid), - /// Response written; parked (or about to park) waiting for the next - /// keep-alive request. - ConnIdle(Pid), - ConnEnded(Pid), -} - -pub enum Call { - ConnCount, -} - -pub enum Reply { - ConnCount(usize), -} - -// --------------------------------------------------------------------------- -// Server -// --------------------------------------------------------------------------- - -#[derive(Clone, Copy, PartialEq, Eq)] -enum ConnState { - Busy, - Idle, -} - -pub struct ConnRegistry { - 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 Reply = Reply; - 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()), - } - } - - fn handle_cast(&mut self, request: Cast) { - match request { - Cast::ConnStarted(pid) => { - 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; - // it leaves the map via its guard's ConnEnded. - smarm::request_stop(pid); - } - } - Cast::ConnBusy(pid) => { - if let Some(s) = self.conns.get_mut(&pid) { - *s = ConnState::Busy; - } - } - Cast::ConnIdle(pid) => { - if let Some(s) = self.conns.get_mut(&pid) { - *s = ConnState::Idle; - if self.draining { - // Finished its in-flight request; nothing more is - // owed. - 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. - } -} - -/// 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() -} - -// --------------------------------------------------------------------------- -// Drop guard — self-deregistration on any exit path. -// --------------------------------------------------------------------------- - -/// Casts `make(pid)` on drop. Runs on normal return, `request_stop` -/// unwind, and panic unwind alike; the cast is infallible from the -/// guard's perspective (a dead registry just returns an ignored Err). -pub struct DeregisterGuard { - registry: GenServerRef, - pid: Pid, - make: fn(Pid) -> Cast, -} - -impl DeregisterGuard { - pub fn new(registry: GenServerRef, pid: Pid, make: fn(Pid) -> Cast) -> Self { - Self { - registry, - pid, - make, - } - } -} - -impl Drop for DeregisterGuard { - fn drop(&mut self) { - 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/endpoint.rs b/src/endpoint.rs new file mode 100644 index 0000000..5ba51a4 --- /dev/null +++ b/src/endpoint.rs @@ -0,0 +1,656 @@ +//! The endpoint: one supervised gen_server that owns a listening socket, +//! the pool of listener actors accepting on it, and the set of live +//! connections. +//! +//! # Shape (v0.3) +//! +//! ```text +//! your root supervisor +//! └── ChildSpec(Permanent, urus::endpoint(config, pipeline)?) <- Endpoint gen_server +//! └── listener_sup OneForOne over N listener actors (spawned in init) +//! └── (listeners spawn plain connection actors) +//! ``` +//! +//! The endpoint runs *inline as the ChildSpec's actor* +//! (`NamedGenServerBuilder::run`), so your supervisor's ordered shutdown +//! reaches it as [`GenServer::handle_shutdown`] and a restart re-runs the +//! factory on the same still-open listen fds. +//! +//! ## Why the endpoint spawns its own listener supervisor +//! +//! The obvious tree is registry and listener-sup as *siblings* under a +//! `RestForOne`, with listeners resolving the registry by name. It has a +//! boot race: smarm's supervisor starts children with a fire-and-forget +//! spawn, so start *order* is not start *readiness* — a listener can look +//! the name up before the registry actor has run. Rather than paper over +//! that with a retry loop, the registrar spawns its consumers: everything +//! here is program order inside one actor's `init`, not a cross-actor +//! guarantee. (The general gap is filed in smarm's ROADMAP as a +//! supervisor readiness-ack item; when it lands, the sibling shape becomes +//! available, but this one costs nothing and is not waiting on it.) +//! +//! ## Connection accounting +//! +//! 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 `ConnStarted` always precedes +//! `ConnEnded` in the inbox. Were the *listener* to announce the pid +//! instead, a short-lived conn's `Ended` could overtake its `Started` and +//! leak a dead pid into the set forever. +//! +//! ## Shutdown +//! +//! `handle_shutdown` (i.e. a `request_shutdown` from your supervisor, or +//! from anywhere): +//! +//! 1. flip `draining` — from here on, a connection that registers or goes +//! idle is stopped on the spot; +//! 2. `request_shutdown` the listener supervisor, which stops its +//! listeners in reverse order and exits; its monitored death is the +//! "no new connections can ever be accepted" barrier; +//! 3. stop every currently-idle connection; +//! 4. arm one `drain_timeout` timer; +//! 5. return `Continue` — the endpoint exits normally once the listener +//! sup is down *and* the connection set is empty, so "the endpoint +//! actor has exited" is exactly "every connection is gone". With no +//! connections and (once the Down lands) nothing to wait for, that is +//! immediate. +//! +//! At the drain deadline one force sweep `request_stop`s whatever remains +//! (plus, defensively, the listener sup). A single sweep suffices: every +//! path a connection can enter the set on afterwards stops it at the point +//! of entry. Force-stopped conns unwind safely out of their fd waits and +//! close their sockets via `OwnedFd::drop`. +//! +//! Give the endpoint [`Shutdown::Infinity`](smarm::supervisor::Shutdown) in +//! your child spec: it bounds itself with `drain_timeout`, and a +//! supervisor-imposed deadline shorter than that would kill the drain +//! halfway and orphan connections. +//! +//! ## Known residual window +//! +//! A connection *spawned* by a listener in the instant before that +//! listener is stopped, which has not yet run, has not registered. If the +//! set was already empty the endpoint can exit before its `ConnStarted` +//! arrives, and that connection serves on as a forest root until smarm's +//! root-exit sweep collects it. Inherent to self-registration (the +//! alternative loses `Started`/`Ended` ordering, which is worse); +//! documented rather than defended. + +use crate::conn_actor::{run_connection, ConnLimits}; +use crate::net::{accept_nonblocking, bind_and_listen, OwnedFd}; +use crate::plug::Pipeline; +use crate::serve::Config; + +use smarm::gen_server::{GenServerName, ShutdownAction, StopHandle, TimerHandle}; +use smarm::{ + ChildSpec, GenServer, GenServerBuilder, GenServerCtx, GenServerRef, OneForOne, Pid, Restart, + Strategy, +}; + +use std::collections::HashMap; +use std::io::{self, ErrorKind}; +use std::os::fd::RawFd; +use std::sync::atomic::{AtomicU32, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +// --------------------------------------------------------------------------- +// Messages +// --------------------------------------------------------------------------- + +pub enum Cast { + /// A connection actor started (self-registered, initially idle: it has + /// not parsed a request head yet). + ConnStarted(Pid), + /// Parsed a request head; a response is now owed. + ConnBusy(Pid), + /// Response written; parked (or about to park) waiting for the next + /// keep-alive request. + ConnIdle(Pid), + ConnEnded(Pid), +} + +pub enum Call { + /// Live connection count — introspection (and the planned hook for + /// ws/channels stats). + ConnCount, +} + +pub enum Reply { + ConnCount(usize), +} + +// --------------------------------------------------------------------------- +// Endpoint +// --------------------------------------------------------------------------- + +#[derive(Clone, Copy, PartialEq, Eq)] +enum ConnState { + Busy, + Idle, +} + +/// Everything an endpoint incarnation needs to (re)build its listener +/// pool. Shared behind an `Arc` by the factory closure so a restart reuses +/// the same already-bound fds — no re-bind, no window where the port is +/// unclaimed. +struct Boot { + listener_fds: Vec>, + pipeline: Pipeline, + limits: ConnLimits, + conn_stack_reserve: usize, + name: &'static str, +} + +pub struct Endpoint { + boot: Arc, + conns: HashMap, + draining: bool, + drain_timeout: Duration, + /// The listener supervisor, spawned in `init` and monitored. `None` + /// once it is down. + listener_sup: Option, + stop: Option>, + timer: Option>, +} + +impl Endpoint { + fn name(&self) -> GenServerName { + GenServerName::new(self.boot.name) + } + + /// Exit normally once nothing can arrive and nothing is left: the + /// listener sup is down and the conn set is empty. This exit *is* the + /// barrier a supervisor's ordered shutdown waits on. + fn stop_if_drained(&self) { + if self.draining && self.listener_sup.is_none() && self.conns.is_empty() { + self.stop.as_ref().expect("init ran first").stop(); + } + } +} + +impl GenServer for Endpoint { + type Call = Call; + type Reply = Reply; + 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()); + + // Our own name is bound (from inside `run`) before `init` runs, so + // this resolves to us. Listeners are handed the resolved ref, not + // the name: no per-connection registry lookup on the accept path. + let me: GenServerRef = + smarm::gen_server::whereis_server(self.name()).expect("endpoint name bound before init"); + + let boot = self.boot.clone(); + let sup = smarm::spawn(move || { + let mut sup = OneForOne::new().strategy(Strategy::OneForOne); + for lfd in &boot.listener_fds { + let lfd = lfd.clone(); + let pipeline = boot.pipeline.clone(); + let limits = boot.limits; + let reserve = boot.conn_stack_reserve; + let me = me.clone(); + // Permanent: a listener only ever exits by supervisor + // action now (its normal-exit-as-shutdown flag is gone), + // so "exited on its own" always means something broke and + // always deserves a restart. + sup = sup.child(ChildSpec::new(Restart::Permanent, move || { + listener_loop(lfd.clone(), pipeline.clone(), limits, reserve, me.clone()); + })); + } + // Default intensity (3 restarts / 5s) applies; a listener + // crash-looping faster than that trips the cap and the pool + // tears down — which our monitor turns into a loud endpoint + // failure rather than a zombie server on a dead port. + sup.run(); + }); + ctx.watch(smarm::monitor(sup.pid())); + self.listener_sup = Some(sup.pid()); + } + + fn handle_call(&mut self, request: Call) -> Reply { + match request { + Call::ConnCount => Reply::ConnCount(self.conns.len()), + } + } + + fn handle_cast(&mut self, request: Cast) { + match request { + Cast::ConnStarted(pid) => { + self.conns.insert(pid, ConnState::Idle); + if self.draining { + // Accepted just before its listener was stopped. + // Draining means no new work; it leaves the map via + // its guard's ConnEnded. + smarm::request_stop(pid); + } + } + Cast::ConnBusy(pid) => { + if let Some(s) = self.conns.get_mut(&pid) { + *s = ConnState::Busy; + } + } + Cast::ConnIdle(pid) => { + if let Some(s) = self.conns.get_mut(&pid) { + *s = ConnState::Idle; + if self.draining { + smarm::request_stop(pid); // in-flight request finished + } + } + } + Cast::ConnEnded(pid) => { + self.conns.remove(&pid); + self.stop_if_drained(); + } + } + } + + fn handle_down(&mut self, down: smarm::Down) { + if self.listener_sup != Some(down.pid) { + return; + } + self.listener_sup = None; + if !self.draining { + // The pool died on its own (restart intensity exceeded, or the + // supervisor itself failed). Nothing is listening any more, so + // this endpoint is a zombie: fail loudly and let the caller's + // supervisor decide (restart re-runs init on the same fds). + panic!("urus endpoint '{}': listener pool died: {:?}", self.boot.name, down.reason); + } + // Expected during shutdown: the accept side is now provably gone. + self.stop_if_drained(); + } + + fn handle_shutdown(&mut self) -> ShutdownAction { + self.draining = true; + if let Some(sup) = self.listener_sup { + // Stops listeners in reverse start order and exits; the + // monitored Down is our "no new connections" barrier. + smarm::request_shutdown(sup); + } + 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 first") + .arm_after(self.drain_timeout, ()); + ShutdownAction::Continue + } + + fn handle_timer(&mut self, _deadline: ()) { + // Drain deadline: force-stop everything left. One sweep — late + // registrants are already stopped on arrival (see module docs). + for pid in self.conns.keys() { + smarm::request_stop(*pid); + } + if let Some(sup) = self.listener_sup { + // Defensive: a listener wedged past its supervisor's own grace + // period must not hold the whole shutdown open. + smarm::request_stop(sup); + } + // The map empties via each conn's guard ConnEnded; stop_if_drained + // fires on the last one (or on the sup's Down, whichever is last). + } +} + +// --------------------------------------------------------------------------- +// endpoint() — the public constructor +// --------------------------------------------------------------------------- + +/// Bind `config.addr` and return the endpoint's supervisable body: pass it +/// to [`ChildSpec::new`] under your own supervisor, alongside your +/// application's other children. +/// +/// The bind happens **here**, eagerly, so an address-in-use error surfaces +/// on the caller's thread rather than inside an actor — and the fds +/// outlive any restart of the child. +/// +/// ```ignore +/// let endpoint = urus::endpoint(Config::new(addr), pipeline)?; +/// let rt = smarm::init(smarm::Config::default()); +/// rt.run(move || { +/// smarm::OneForOne::new() +/// .child(ChildSpec::new(Restart::Permanent, my_app_state)) +/// .child(ChildSpec::new(Restart::Permanent, endpoint) +/// .shutdown(Shutdown::Infinity)) +/// .run() +/// }); +/// ``` +/// +/// The returned closure is `Clone`, so one endpoint definition can be +/// handed to more than one place; each *invocation* is one running +/// endpoint, and two live at once under the same [`Config::name`] is a +/// name clash (panic on the second). +pub fn endpoint( + config: Config, + pipeline: Pipeline, +) -> io::Result { + // One dup'd fd per listener: `accept4` is thread-safe on a single fd, + // but a per-listener RawFd keeps each actor's epoll registration + // distinct in smarm's `waiters: HashMap`. + let mut listener_fds = Vec::with_capacity(config.listener_pool); + listener_fds.push(Arc::new(bind_and_listen(config.addr)?)); + for _ in 1..config.listener_pool { + let dup = dup_fd(listener_fds[0].as_raw())?; + listener_fds.push(Arc::new(dup)); + } + + let boot = Arc::new(Boot { + listener_fds, + pipeline, + limits: config.to_conn_limits(), + conn_stack_reserve: config.conn_stack_reserve, + name: config.name, + }); + let drain_timeout = config.drain_timeout; + + Ok(move || { + let state = Endpoint { + boot: boot.clone(), + conns: HashMap::new(), + draining: false, + drain_timeout, + listener_sup: None, + stop: None, + timer: None, + }; + let name = GenServerName::::new(state.boot.name); + // Inline: this actor *is* the server, so a supervisor's shutdown + // reaches handle_shutdown and a restart re-binds the name. + if let Err(e) = GenServerBuilder::new(state).named(name).run() { + panic!("urus endpoint '{}': name already taken: {e:?}", name.as_str()); + } + }) +} + +/// Resolve a running endpoint by [`Config::name`] — for `ConnCount` and +/// other introspection from elsewhere in the app. +pub fn whereis(name: &'static str) -> Option> { + smarm::gen_server::whereis_server(GenServerName::::new(name)) +} + +// --------------------------------------------------------------------------- +// Listener actors +// --------------------------------------------------------------------------- + +fn dup_fd(fd: RawFd) -> io::Result { + let new_fd = unsafe { libc::fcntl(fd, libc::F_DUPFD_CLOEXEC, 0) }; + if new_fd < 0 { + return Err(io::Error::last_os_error()); + } + Ok(OwnedFd::from_raw(new_fd)) +} + +/// Test-only fault injection. When nonzero, the next accept-loop iteration +/// of whichever listener gets there first decrements this and panics — +/// *before* calling `accept`, so a pending connection stays in the kernel +/// backlog and must be picked up by the restarted listener. Cost when idle +/// is one relaxed load per accept-loop iteration (each of which already +/// pays a syscall). Not public API. +#[doc(hidden)] +pub static INJECT_LISTENER_PANICS: AtomicU32 = AtomicU32::new(0); + +/// Accept forever, spawning one connection actor per connection. Exits +/// only by supervisor action: `request_stop` unwinds the untimed +/// `wait_readable` park below (smarm's 06-10 io fix), which is why there +/// is no shutdown flag and no tick — the park is genuinely open-ended and +/// costs nothing while idle. +fn listener_loop( + listener: Arc, + pipeline: Pipeline, + limits: ConnLimits, + conn_stack_reserve: usize, + endpoint: GenServerRef, +) { + let fd = listener.as_raw(); + + loop { + if INJECT_LISTENER_PANICS + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| n.checked_sub(1)) + .is_ok() + { + panic!("urus: injected listener panic (test hook)"); + } + match accept_nonblocking(fd) { + Ok(client) => { + // Hand the fd off to a new connection actor. spawn() is + // cheap on smarm — it's a single Vec push under the + // shared lock. + let p = pipeline.clone(); + let l = limits; + let e = endpoint.clone(); + let opts = smarm::SpawnOpts { + stack_reserve: Some(conn_stack_reserve), + ..smarm::SpawnOpts::default() + }; + smarm::spawn_with(opts, move || run_connection(client, p, l, e)); + } + Err(e) if e.kind() == ErrorKind::WouldBlock => { + if let Err(we) = smarm::wait_readable(fd) { + // epoll registration failed — abnormal, so panic: a + // transient failure (e.g. EMFILE on the epoll set) + // heals by restart instead of silently shrinking the + // pool. smarm catches actor panics in the trampoline; + // this is a Signal::Panic to the supervisor, not + // process noise. + panic!("urus: listener wait_readable failed: {we}"); + } + } + Err(e) if e.kind() == ErrorKind::Interrupted => continue, + Err(e) => { + // EMFILE / ENFILE / ECONNABORTED etc. Log and back off + // briefly in case the error is sticky; the system may + // recover. + eprintln!("urus: accept error: {e}"); + smarm::sleep(Duration::from_millis(10)); + } + } + } + // The Arc clone we were started with drops on unwind, but the + // ChildSpec factory holds another — the fd outlives any one + // incarnation of this listener. +} + +// --------------------------------------------------------------------------- +// Deregistration guard +// --------------------------------------------------------------------------- + +/// Casts `make(pid)` on drop. Runs on normal return, `request_stop` +/// unwind, and panic unwind alike; the cast is infallible from the +/// guard's perspective (a dead endpoint just returns an ignored Err). +pub struct DeregisterGuard { + endpoint: GenServerRef, + pid: Pid, + make: fn(Pid) -> Cast, +} + +impl DeregisterGuard { + pub fn new(endpoint: GenServerRef, pid: Pid, make: fn(Pid) -> Cast) -> Self { + Self { endpoint, pid, make } + } +} + +impl Drop for DeregisterGuard { + fn drop(&mut self) { + let _ = self.endpoint.cast((self.make)(self.pid)); + } +} + +// --------------------------------------------------------------------------- +// Tests — the drain protocol, driven purely by request_shutdown. +// --------------------------------------------------------------------------- + +#[cfg(test)] +mod tests { + use super::*; + use crate::plug::Pipeline; + use std::time::Instant; + + const NAME: &str = "urus.test.endpoint"; + + /// Spawn a real endpoint (bound to an ephemeral port) as a supervised + /// child, exactly as an application would, and hand back its + /// supervisor's pid. + fn spawn_endpoint(drain: Duration) -> Pid { + let cfg = Config { + listener_pool: 1, + drain_timeout: drain, + name: NAME, + ..Config::new("127.0.0.1:0".parse().unwrap()) + }; + let body = endpoint(cfg, Pipeline::new()).expect("bind"); + let sup = smarm::spawn(move || { + OneForOne::new() + .child( + ChildSpec::new(Restart::Permanent, body) + .shutdown(smarm::supervisor::Shutdown::Infinity), + ) + .run() + }); + await_pred("endpoint up", Duration::from_secs(2), || { + whereis(NAME).is_some() + }); + sup.pid() + } + + /// A stand-in connection actor: registers with the endpoint, optionally + /// reports busy, then parks forever. Only a `request_stop` ends it; the + /// guard's ConnEnded runs on the unwind. + fn fake_conn(ep: GenServerRef, busy: bool) { + let me = smarm::self_pid(); + let _ = ep.cast(Cast::ConnStarted(me)); + let _guard = DeregisterGuard::new(ep.clone(), me, Cast::ConnEnded); + if busy { + let _ = ep.cast(Cast::ConnBusy(me)); + } + loop { + smarm::sleep(Duration::from_secs(3600)); + } + } + + fn spawn_fake_conn(busy: bool) -> Pid { + let ep = whereis(NAME).expect("endpoint live"); + smarm::spawn(move || fake_conn(ep, busy)).pid() + } + + /// Live conn count, or `Err` once the endpoint is gone. + fn conn_count() -> Result { + match whereis(NAME) { + Some(ep) => match ep.call(Call::ConnCount) { + Ok(Reply::ConnCount(n)) => Ok(n), + Err(_) => Err(()), + }, + None => Err(()), + } + } + + 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_promptly() { + smarm::run(|| { + let sup = spawn_endpoint(Duration::from_secs(30)); + assert_eq!(conn_count(), Ok(0)); + let t0 = Instant::now(); + smarm::request_shutdown(sup); + // Nothing to drain: gone well inside the (huge) drain window, + // i.e. the idle path does not run the clock out. + await_pred("endpoint exit", Duration::from_secs(3), || { + conn_count().is_err() + }); + assert!(t0.elapsed() < Duration::from_secs(3)); + }); + } + + #[test] + fn drain_stops_idle_now_busy_at_deadline_then_exits() { + smarm::run(|| { + let drain = Duration::from_millis(300); + let sup = spawn_endpoint(drain); + spawn_fake_conn(false); // idle + spawn_fake_conn(true); // busy + await_pred("both registered", Duration::from_secs(2), || { + conn_count() == Ok(2) + }); + + let t0 = Instant::now(); + smarm::request_shutdown(sup); + + // 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 the endpoint winds up. + await_pred("busy swept + endpoint exit", drain * 6, || { + 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: this passes only via the ConnIdle path, not + // the force sweep. + let sup = spawn_endpoint(Duration::from_secs(30)); + let conn = spawn_fake_conn(true); + await_pred("registered busy", Duration::from_secs(2), || { + conn_count() == Ok(1) + }); + smarm::request_shutdown(sup); + smarm::sleep(Duration::from_millis(50)); + assert_eq!(conn_count(), Ok(1), "busy conn survives the idle sweep"); + + // "Request finishes": the conn reports idle. + let ep = whereis(NAME).expect("endpoint still draining"); + let _ = ep.cast(Cast::ConnIdle(conn)); + await_pred("stopped on idle + endpoint exit", Duration::from_secs(3), || { + conn_count().is_err() + }); + }); + } + + #[test] + fn conn_registering_mid_drain_is_stopped_on_arrival() { + smarm::run(|| { + let sup = spawn_endpoint(Duration::from_secs(30)); + let holder = spawn_fake_conn(true); // keeps the drain open + await_pred("holder registered", Duration::from_secs(2), || { + conn_count() == Ok(1) + }); + smarm::request_shutdown(sup); + smarm::sleep(Duration::from_millis(20)); + + // The listener-death race: a fresh conn registers mid-drain. + spawn_fake_conn(false); + // It is stopped on arrival — the count returns to just 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 ep = whereis(NAME).expect("endpoint still draining"); + let _ = ep.cast(Cast::ConnIdle(holder)); + }); + } +} diff --git a/src/lib.rs b/src/lib.rs index 0b7d88c..708d20a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -24,7 +24,7 @@ pub mod router; pub mod parser; pub mod net; pub mod conn_actor; -pub mod conn_registry; +pub mod endpoint; pub mod serve; pub mod sse; pub mod ws; @@ -41,6 +41,7 @@ pub use pubsub::{PubSub, PubSubDown}; #[cfg(feature = "channels")] pub use channels::{Channel, ChannelHub, ChannelSession, ChannelSocket, PrefixRouter, Status, TopicRouter}; pub use ws::{Message, WsClosed, WsHandler, WsSender}; +pub use endpoint::endpoint; pub use serve::{ serve, serve_with, serve_with_shutdown, shutdown_handle, Config, Handle, ShutdownSignal, }; diff --git a/src/serve.rs b/src/serve.rs index ae0b081..1fc0417 100644 --- a/src/serve.rs +++ b/src/serve.rs @@ -1,29 +1,19 @@ -//! Listener pool and the `serve` entry point. +//! [`Config`] and the `serve*` entry points. //! -//! A small fixed pool of listener actors share the same TCP listen fd (via -//! `dup`) — each blocks in non-blocking `accept4` + `wait_readable` on its -//! own copy. When a connection arrives the listener spawns a connection -//! actor with the `OwnedFd` and immediately returns to `accept`. No -//! coordination needed; the kernel serialises `accept` calls across the fds. +//! The listener pool itself lives in [`crate::endpoint`], which is the +//! real API: an endpoint is a supervisable child you place in your own +//! tree. `serve*` is the batteries-included path for a process whose only +//! job is serving HTTP — it owns the runtime and wraps one endpoint in a +//! one-child supervisor. //! -//! Sharing via `dup` rather than the same fd is deliberate — Linux's -//! `accept4` is thread-safe on a single fd, but dup'ing per-listener keeps -//! each actor's epoll registration local to its own RawFd value (so smarm's -//! `waiters: HashMap` doesn't see collisions between listeners -//! waiting on "the same fd"). - -use crate::conn_actor::{run_connection, ConnLimits}; -use crate::conn_registry::{self, ConnRegistry}; -use crate::net::{accept_nonblocking, bind_and_listen, OwnedFd}; +use crate::conn_actor::ConnLimits; use crate::plug::Pipeline; -use smarm::{ChildSpec, OneForOne, Restart, GenServerRef, Strategy}; +use smarm::supervisor::Shutdown; +use smarm::{ChildSpec, OneForOne, Restart}; use std::io::{self, ErrorKind}; use std::net::{SocketAddr, ToSocketAddrs}; -use std::os::fd::RawFd; -use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; -use std::sync::Arc; use std::time::Duration; // --------------------------------------------------------------------------- @@ -67,10 +57,11 @@ pub struct Config { /// WebSocket: cap on a complete reassembled message (spans /// fragments; violation closes 1009). pub max_message_bytes: usize, - /// Number of smarm scheduler OS threads. `None` means smarm's default - /// (one per CPU). Set this to a small fixed number in tests so multiple - /// concurrent test servers don't oversubscribe the host. - pub scheduler_threads: Option, + /// Registry name for this endpoint's gen_server — how it is addressed + /// from elsewhere in the app ([`crate::endpoint::whereis`]), and what + /// must be unique between two endpoints in one process (a public and + /// an admin port, say). Default `"urus"`. + pub name: &'static str, /// Stack reserve (RFC 019 `smarm::SpawnOpts::stack_reserve`) given to /// each per-connection actor. Request handlers routinely pull in /// application code — DB drivers, (de)compression, templating — whose @@ -79,12 +70,6 @@ pub struct Config { /// Default: 256 KiB. The reserve is virtual/demand-paged, so raising it /// costs address space, not RSS, until a handler actually uses it. pub conn_stack_reserve: usize, - /// Maximum concurrently-live actors — smarm's fixed slot slab, allocated - /// once at init. Each connection is one actor, so this is also the hard - /// cap on concurrent connections. `None` uses smarm's default (16_384). - /// Slots are ~256 B, so raising this is cheap relative to per-connection - /// stacks; size it to peak concurrent connections. - pub max_actors: Option, } /// Default per-connection actor stack reserve (see [`Config::conn_stack_reserve`]). @@ -111,13 +96,12 @@ impl Config { drain_timeout: Duration::from_secs(30), max_frame_payload: 1024 * 1024, max_message_bytes: 4 * 1024 * 1024, - scheduler_threads: None, + name: "urus", conn_stack_reserve: DEFAULT_CONN_STACK_RESERVE, - max_actors: None, } } - fn to_conn_limits(&self) -> ConnLimits { + pub(crate) fn to_conn_limits(&self) -> ConnLimits { ConnLimits { max_headers: self.max_header_count, initial_read_buf: self.read_buf_size, @@ -209,129 +193,14 @@ impl Config { } // --------------------------------------------------------------------------- -// dup helper +// Handle / ShutdownSignal — graceful shutdown plumbing for the serve* entries. // --------------------------------------------------------------------------- -fn dup_fd(fd: RawFd) -> io::Result { - let new_fd = unsafe { libc::fcntl(fd, libc::F_DUPFD_CLOEXEC, 0) }; - if new_fd < 0 { - return Err(io::Error::last_os_error()); - } - Ok(OwnedFd::from_raw(new_fd)) -} - -// --------------------------------------------------------------------------- -// listener actor body -// --------------------------------------------------------------------------- - -/// Test-only fault injection. When nonzero, the next accept-loop iteration -/// of whichever listener gets there first decrements this and panics — -/// *before* calling `accept`, so a pending connection stays in the kernel -/// backlog and must be picked up by the restarted listener. Cost when idle -/// is one relaxed load per accept-loop iteration (each of which already -/// pays a syscall). Not public API. -#[doc(hidden)] -pub static INJECT_LISTENER_PANICS: AtomicU32 = AtomicU32::new(0); - -fn listener_loop( - listener: Arc, - pipeline: Pipeline, - limits: ConnLimits, - conn_stack_reserve: usize, - registry: GenServerRef, - shutdown: Arc, -) { - let fd = listener.as_raw(); - - loop { - // Self-termination on shutdown — the ONLY way a listener exits at - // shutdown, and deliberately a normal return: under - // `Restart::Transient` a normal exit is terminal, so the - // supervisor's active-count drains and `sup.run()` returns on its - // own. No pid is ever `request_stop`ped, which sidesteps both the - // spawn-to-registration race of self-announced pids and smarm's - // lossy stop-while-QUEUED window (a flag is wake-free and - // race-free; a freshly spawned listener observes it on its very - // first iteration, a parked one within LISTENER_TICK). - if shutdown.load(Ordering::Relaxed) { - return; - } - if INJECT_LISTENER_PANICS - .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| n.checked_sub(1)) - .is_ok() - { - panic!("urus: injected listener panic (test hook)"); - } - match accept_nonblocking(fd) { - Ok(client) => { - // Hand the fd off to a new connection actor. spawn() is - // cheap on smarm — it's a single Vec push under the - // shared lock. - let p = pipeline.clone(); - let l = limits; - let r = registry.clone(); - let opts = smarm::SpawnOpts { - stack_reserve: Some(conn_stack_reserve), - ..smarm::SpawnOpts::default() - }; - smarm::spawn_with(opts, move || run_connection(client, p, l, r)); - } - Err(e) if e.kind() == ErrorKind::WouldBlock => { - // No pending connection. Park until the listener is - // readable or the tick elapses; either way we come back - // around through the shutdown-flag check above. - match smarm::wait_readable_timeout(fd, LISTENER_TICK) { - Ok(_ready) => {} // ready or tick — loop re-checks, retries accept - Err(we) => { - // epoll registration failed. Under Transient - // supervision a normal return is terminal but a - // panic restarts us — and a failed wait IS - // abnormal, so panic: a transient failure (e.g. - // EMFILE on the epoll set) heals by restart - // instead of silently shrinking the pool. (smarm - // catches actor panics in the trampoline; this is - // a Signal::Panic to the supervisor, not process - // noise.) - panic!("urus: listener wait_readable failed: {we}"); - } - } - } - Err(e) if e.kind() == ErrorKind::Interrupted => { - continue; - } - Err(e) => { - // EMFILE / ENFILE / ECONNABORTED etc. Log and continue; - // the system may recover. - eprintln!("urus: accept error: {e}"); - // Small backoff via smarm's sleep to avoid spinning if - // the error is sticky. - smarm::sleep(Duration::from_millis(10)); - } - } - } - // The Arc clone we were started with drops here (normal exit or - // unwind), but the ChildSpec factory holds another — the fd outlives - // any one incarnation of this listener. -} - -// --------------------------------------------------------------------------- -// Handle / ShutdownSignal — graceful shutdown plumbing. -// --------------------------------------------------------------------------- - -/// How often the root actor polls for a shutdown signal (see -/// `serve_with_shutdown` for why this is a poll); bounds shutdown latency. -const SHUTDOWN_POLL: Duration = Duration::from_millis(100); - -/// Listener accept-waits are timed at this tick; the shutdown flag is -/// observed at the top of every accept-loop iteration, so this bounds how -/// long a fully idle listener takes to notice shutdown. It is the SOLE -/// listener-shutdown mechanism (no `request_stop` — see the shutdown -/// sequence notes). Idle cost: one timer wake per listener per tick. -const LISTENER_TICK: Duration = Duration::from_millis(250); - -/// A clonable trigger for graceful shutdown. Safe to use from any OS -/// thread (the send only enqueues; the serving side polls), e.g. from a -/// signal-handling thread. +/// A clonable trigger for graceful shutdown, usable from any OS thread +/// (e.g. a signal-handling thread) — the shutdown path for callers who let +/// `serve*` own the runtime and so have no [`smarm::RuntimeHandle`] of +/// their own. If you build the tree yourself with [`crate::endpoint`], use +/// `rt.handle().request_shutdown(root_sup)` instead and ignore this. #[derive(Clone)] pub struct Handle { tx: smarm::Sender<()>, @@ -340,8 +209,8 @@ pub struct Handle { impl Handle { /// Begin graceful shutdown: stop accepting, close idle keep-alive /// connections, drain in-flight requests up to `Config.drain_timeout`, - /// then force-stop stragglers. `serve_with_shutdown` returns once the - /// runtime has wound down. Idempotent; extra calls are no-ops. + /// then force-stop stragglers. `serve*` returns once the runtime has + /// wound down. Idempotent; extra calls are no-ops. pub fn shutdown(&self) { let _ = self.tx.send(()); } @@ -359,174 +228,80 @@ pub fn shutdown_handle() -> (Handle, ShutdownSignal) { } // --------------------------------------------------------------------------- -// serve_with_shutdown — main entry. Boots smarm, supervises listeners, -// blocks until shutdown. +// serve* — batteries-included entries for apps whose only job is serving. // --------------------------------------------------------------------------- // -// Boots an smarm runtime (one OS thread per CPU by default — see smarm's -// `Config::default()`). The root actor starts the connection registry and -// a one-for-one supervisor over the listener pool, then blocks on the -// shutdown signal. A panicking listener is restarted on the same -// (still-open) fd instead of silently shrinking the accept pool. -// Connection actors stay unsupervised bare spawns — per spec, a connection -// is cheap and its failure is local: a 500 path, not a restart path. -// (`spawn` from a listener parents the conn actor under that listener, -// which never registers a supervisor channel, so conn deaths are invisible -// to the pool supervisor by construction.) -// -// Shutdown sequence (drain-then-stop, the v0.2 chunk 2 decision): -// 1. Set the shared shutdown flag. Every listener observes it at the -// top of its accept loop (parks are timed at LISTENER_TICK) and -// returns normally. Under `Restart::Transient` a normal exit is -// terminal, so no restart happens and the supervisor's active count -// drains to zero. -// 2. Join the supervisor: `sup.run()` returns on its own once every -// listener has exited. The join is therefore the barrier "no new -// connections can ever be accepted" — listener fds are closed (the -// last Arc clones drop with the supervisor's ChildSpecs), and the -// kernel refuses new connects. -// -// Why a flag and not `request_stop`: stopping pids that announce -// themselves races the spawn-to-registration gap (a fast shutdown -// CAN beat a fresh listener to the registry), and smarm's -// `request_stop` is lossy against a QUEUED actor that then parks -// without passing an observation point — found the hard way; see -// the chunk 2 commit message. The flag is wake-free and race-free. -// 3. BeginDrain: the registry stops idle conns now and each remaining -// conn the moment it finishes its in-flight request. -// 4. Poll until no conns remain or the drain deadline passes; past the -// deadline, ForceStopConns every tick until the set empties (a conn -// accepted just before listener death may register late). -// 5. Root returns. `rt.run` itself returns only when every actor has -// exited (force-stopped conns unwind through their fd waits safely — -// smarm's 06-10 io fix — and close their sockets via OwnedFd::drop). +// These own the smarm runtime and build a one-child tree around +// `endpoint()`. An app with its own actors should call `endpoint()` +// directly and put it in its own supervision tree — that is the real API; +// everything here is a convenience wrapper over it. +/// Boot a runtime, serve until `signal` fires, then drain and return. +/// +/// The tree is `root sup -> endpoint`, with the endpoint on +/// [`Shutdown::Infinity`] so its `drain_timeout` — not a supervisor +/// deadline — bounds the drain. The root actor parks on the signal +/// channel; a `Handle::shutdown` from a foreign OS thread wakes it, it +/// shuts the supervisor down (ordered, so the endpoint drains) and +/// `rt.run` returns when the last actor is gone. pub fn serve_with_shutdown( config: Config, + rt_config: smarm::Config, pipeline: Pipeline, signal: ShutdownSignal, ) -> io::Result<()> { - let listener = bind_and_listen(config.addr)?; - println!("urus: listening on {}", config.addr); - - // One connection-actor-spawning loop per listener pool slot. Each gets - // its own dup'd fd so epoll registrations don't collide. Each fd is - // owned by its ChildSpec's factory closure (via `Arc`): a restarted - // listener re-enters `accept`/`wait_readable` on the same fd — no - // re-dup, no window where the slot has no fd. - let mut listener_fds = Vec::with_capacity(config.listener_pool); - listener_fds.push(Arc::new(listener)); // primary keeps the original - - for _ in 1..config.listener_pool { - let dup = dup_fd(listener_fds[0].as_raw())?; - listener_fds.push(Arc::new(dup)); - } - - let limits = config.to_conn_limits(); - let conn_stack_reserve = config.conn_stack_reserve; - let drain_timeout = config.drain_timeout; - - let mut smarm_cfg = match config.scheduler_threads { - Some(n) => smarm::Config::exact(n), - None => smarm::Config::default(), - }; - if let Some(m) = config.max_actors { - smarm_cfg = smarm_cfg.max_actors(m); - } - let rt = smarm::init(smarm_cfg); - // Listener self-termination flag — see the shutdown sequence below. - let shutdown_flag = Arc::new(AtomicBool::new(false)); + let addr = config.addr; + let endpoint = crate::endpoint::endpoint(config, pipeline)?; + println!("urus: listening on {addr}"); + let rt = smarm::init(rt_config); rt.run(move || { - // Registry first: listeners and conns cast into it from birth. - let registry = conn_registry::start(drain_timeout); + let sup = smarm::spawn(move || { + OneForOne::new() + .child( + ChildSpec::new(Restart::Permanent, endpoint).shutdown(Shutdown::Infinity), + ) + .run() + }); - let mut sup = OneForOne::new().strategy(Strategy::OneForOne); - for (i, lfd) in listener_fds.into_iter().enumerate() { - let p = pipeline.clone(); - let r = registry.clone(); - let sf = shutdown_flag.clone(); - sup = sup.child(ChildSpec::new(Restart::Transient, move || { - println!("urus: listener {} starting", i); - listener_loop(lfd.clone(), p.clone(), limits, conn_stack_reserve, r.clone(), sf.clone()); - })); - } - // Default intensity (3 per 5s) applies; a listener crash-looping - // faster than that trips the cap and tears the pool down — loud - // failure over a zombie server. - let sup_h = smarm::spawn(move || sup.run()); - // The old `urus.server` / `urus.listener.{i}` name registrations - // are gone with smarm's RFC 014 registry rework: `register` is now - // `(Name, Sender)`, self-only — a name is a typed messaging - // endpoint, not a pid tag. urus's bindings were introspection-only - // with no channel behind them, so they were dropped rather than - // faked with a unit channel. A real messageable `urus.server` - // name is in the icebox (ROADMAP.md). - - // Block until told to shut down. We poll `try_recv` + `sleep` - // rather than parking in `recv`: a smarm `Sender::send` from a - // foreign OS thread (no runtime in its TLS) enqueues fine but its - // unpark is a `try_with_runtime` no-op — a parked receiver would - // never wake. The timer wake comes from inside the runtime, so the - // poll sees the message within one interval. (A cross-thread-safe - // unpark is a smarm roadmap candidate; this poll dies with it.) - // If every Handle was dropped the channel closes and no shutdown - // can ever arrive: serve forever, exactly v1's semantics. - loop { - match signal.rx.try_recv() { - Ok(Some(())) => break, - Ok(None) => smarm::sleep(SHUTDOWN_POLL), - Err(_) => smarm::sleep(Duration::from_secs(3600)), + // Park until told to shut down. If every Handle was dropped the + // channel closes and no shutdown can ever arrive: serve forever, + // exactly v1's semantics. + match signal.rx.recv() { + Ok(()) => smarm::request_shutdown(sup.pid()), + Err(_) => { + // Sender side gone. Park indefinitely; the process is + // expected to be killed externally. + loop { + smarm::sleep(Duration::from_secs(3600)); + } } } - - // ----- Shutdown. ----- - // 1 + 2. Flag the listeners down and join the supervisor; the - // join returns once every listener has exited normally - // (Transient: normal exit is terminal). After this point - // no connection can ever be accepted again. - shutdown_flag.store(true, Ordering::Relaxed); - let _ = sup_h.join(); - - // 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. Root returns; the runtime winds down when the last actor - // exits. + let _ = sup.join(); }); Ok(()) } -// --------------------------------------------------------------------------- -// serve_with / serve — convenience entries without a shutdown handle. -// --------------------------------------------------------------------------- - -pub fn serve_with(config: Config, pipeline: Pipeline) -> io::Result<()> { - // The Handle is dropped immediately: shutdown can never be signalled - // and the server runs until externally killed (v1 semantics). +/// [`serve_with_shutdown`] without a shutdown handle: serves until the +/// process is killed. +pub fn serve_with(config: Config, rt_config: smarm::Config, pipeline: Pipeline) -> io::Result<()> { + // The Handle is dropped immediately: shutdown can never be signalled. let (_handle, signal) = shutdown_handle(); - serve_with_shutdown(config, pipeline, signal) + serve_with_shutdown(config, rt_config, pipeline, signal) } -// --------------------------------------------------------------------------- -// serve — convenience over serve_with. -// --------------------------------------------------------------------------- - +/// Defaults all round: default [`Config`], default smarm runtime (one +/// scheduler thread per CPU), serve until killed. pub fn serve(addr: impl ToSocketAddrs, pipeline: Pipeline) -> io::Result<()> { let addr = addr .to_socket_addrs()? .next() .ok_or_else(|| io::Error::new(ErrorKind::InvalidInput, "no addresses resolved"))?; - serve_with(Config::new(addr), pipeline) + serve_with(Config::new(addr), smarm::Config::default(), pipeline) } + #[cfg(all(test, feature = "config-file"))] mod config_file_tests { use super::*; diff --git a/tests/integration.rs b/tests/integration.rs index 65b3c15..a196cd5 100644 --- a/tests/integration.rs +++ b/tests/integration.rs @@ -31,10 +31,9 @@ fn spawn_server(pipeline: Pipeline) -> u16 { std::thread::spawn(move || { let cfg = Config { listener_pool: 2, - scheduler_threads: Some(2), ..Config::new(addr) }; - serve_with(cfg, pipeline).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipeline).unwrap(); }); // Wait for the server to actually be listening. for _ in 0..50 { @@ -252,10 +251,9 @@ fn panicking_listener_restarts() { std::thread::spawn(move || { let cfg = Config { listener_pool: 1, - scheduler_threads: Some(2), ..Config::new(addr) }; - serve_with(cfg, pipe).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipe).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { @@ -271,13 +269,13 @@ fn panicking_listener_restarts() { // Arm the fault: the listener's next accept-loop iteration panics // *before* accepting, so our connection waits in the kernel backlog // until the restarted listener picks it up. - urus::serve::INJECT_LISTENER_PANICS.store(1, std::sync::atomic::Ordering::Relaxed); + urus::endpoint::INJECT_LISTENER_PANICS.store(1, std::sync::atomic::Ordering::Relaxed); let resp = send_request(port, b"GET /ping HTTP/1.1\r\nHost: x\r\nConnection: close\r\n\r\n"); assert_eq!(http_status(&resp), 200, "request pending across the panic was not served"); assert_eq!(http_body(&resp), b"pong"); assert_eq!( - urus::serve::INJECT_LISTENER_PANICS.load(std::sync::atomic::Ordering::Relaxed), + urus::endpoint::INJECT_LISTENER_PANICS.load(std::sync::atomic::Ordering::Relaxed), 0, "fault was never consumed — listener didn't wake for the connection" ); @@ -306,11 +304,10 @@ fn spawn_server_with_handle( std::thread::spawn(move || { let cfg = Config { listener_pool: 2, - scheduler_threads: Some(2), drain_timeout: drain, ..Config::new(addr) }; - urus::serve_with_shutdown(cfg, pipeline, signal).unwrap(); + urus::serve_with_shutdown(cfg, smarm::Config::exact(2), pipeline, signal).unwrap(); let _ = done_tx.send(()); }); for _ in 0..50 { @@ -445,13 +442,12 @@ fn spawn_server_with_timeouts( std::thread::spawn(move || { let cfg = Config { listener_pool: 2, - scheduler_threads: Some(2), keep_alive_timeout: keep_alive, head_timeout: head, body_timeout: body, ..Config::new(addr) }; - serve_with(cfg, pipeline).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipeline).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { @@ -476,7 +472,6 @@ fn spawn_server_with_body_gate( std::thread::spawn(move || { let cfg = Config { listener_pool: 2, - scheduler_threads: Some(2), keep_alive_timeout: Duration::from_secs(30), head_timeout: head, body_timeout: body, @@ -484,7 +479,7 @@ fn spawn_server_with_body_gate( body_stall_timeout: stall, ..Config::new(addr) }; - serve_with(cfg, pipeline).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipeline).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { @@ -813,11 +808,10 @@ fn stalled_reader_killed_at_write_timeout() { std::thread::spawn(move || { let cfg = Config { listener_pool: 2, - scheduler_threads: Some(2), write_timeout: Duration::from_millis(300), ..Config::new(addr) }; - serve_with(cfg, pipe).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipe).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { @@ -903,11 +897,10 @@ fn chunked_request_over_limit_413() { std::thread::spawn(move || { let cfg = Config { listener_pool: 2, - scheduler_threads: Some(2), max_body_bytes: 8, // tiny ..Config::new(addr) }; - serve_with(cfg, pipe).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipe).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { break; } @@ -1816,3 +1809,103 @@ mod channels_wire { .expect("detached session outlived the drain: the registry-drop chain is broken"); } } + +// --------------------------------------------------------------------------- +// Endpoint as a supervised child (v0.3) — the spec §6 tree shape. +// --------------------------------------------------------------------------- + +/// The whole point of v0.3: the APP owns the runtime and the root +/// supervisor; urus is one ordered child among the app's own. A +/// `RuntimeHandle::request_shutdown` on the root sup (the SIGTERM shape) +/// winds the tree down in reverse start order and `rt.run` returns. +#[test] +fn endpoint_as_supervised_child_serves_and_drains() { + use std::sync::atomic::{AtomicBool, Ordering}; + use std::sync::Arc; + + let port = free_port(); + let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); + let app_child_shut_down = Arc::new(AtomicBool::new(false)); + + let pipeline = Pipeline::new() + .plug(Router::new().get("/", |c: Conn, _n: Next| c.put_status(200).put_body("app+urus"))); + + // Endpoint construction is eager about the bind: errors surface here, + // on the app's thread, not inside some actor. + let endpoint = urus::endpoint( + Config { + listener_pool: 2, + drain_timeout: Duration::from_secs(5), + ..Config::new(addr) + }, + pipeline, + ) + .expect("bind"); + + let (handle_tx, handle_rx) = std::sync::mpsc::channel(); + let (pid_tx, pid_rx) = std::sync::mpsc::channel(); + let (done_tx, done_rx) = std::sync::mpsc::channel(); + let flag = app_child_shut_down.clone(); + std::thread::spawn(move || { + let rt = smarm::init(smarm::Config::exact(2)); + handle_tx.send(rt.handle()).unwrap(); + rt.run(move || { + let sup = smarm::spawn(move || { + smarm::OneForOne::new() + .strategy(smarm::Strategy::RestForOne) + // An app child started BEFORE the endpoint: reverse-order + // shutdown must take the endpoint down first, so when + // this child's guard runs the port must already be dead. + .child(smarm::ChildSpec::new(smarm::Restart::Permanent, { + let flag = flag.clone(); + move || { + struct G(Arc); + impl Drop for G { + fn drop(&mut self) { + self.0.store(true, Ordering::SeqCst); + } + } + let _g = G(flag.clone()); + loop { + smarm::sleep(Duration::from_secs(3600)); + } + } + })) + .child(smarm::ChildSpec::new(smarm::Restart::Permanent, endpoint.clone())) + .run() + }); + // Export the sup pid for the "signal thread" below. + pid_tx.send(sup.pid()).unwrap(); + sup.join().expect("root sup returns normally after shutdown"); + }); + let _ = done_tx.send(()); + }); + + // Stand-in signal thread state. + let handle = handle_rx.recv().unwrap(); + let sup_pid = pid_rx.recv().unwrap(); + + // Server is up and serving through the supervised endpoint. + for _ in 0..50 { + if TcpStream::connect(addr).is_ok() { + break; + } + std::thread::sleep(Duration::from_millis(50)); + } + let resp = send_request(port, b"GET / HTTP/1.1\r\nhost: x\r\nconnection: close\r\n\r\n"); + assert!(resp.windows(8).any(|w| w == b"app+urus"), "endpoint serves"); + + // SIGTERM shape: an outside thread shuts the root sup down. + handle.request_shutdown(sup_pid); + done_rx + .recv_timeout(Duration::from_secs(5)) + .expect("rt.run returned after root-sup shutdown"); + assert!( + app_child_shut_down.load(Ordering::SeqCst), + "app sibling wound down" + ); + assert!( + TcpStream::connect(addr).is_err(), + "listen fds closed: no new connections after shutdown" + ); +}