feat(serve): graceful shutdown + connection registry (v0.2 chunk 2)

Shutdown is drain-then-force-stop: flag.store(true) -> sup_h.join() (the
no-more-accepts barrier) -> cast BeginDrain -> poll ConnCount every 50ms
until 0 or drain_timeout; past the deadline each tick casts ForceStopConns
and re-sweeps, so conns that registered after the deadline are still
caught.

Listener shutdown is redesigned vs chunk 1: Restart::Permanent ->
Transient, wait-failure return -> panic (Transient restarts it), and a
shared Arc<AtomicBool> flag checked at loop top with a 250ms
wait_readable_timeout tick instead of request_stop — no request_stop on
the supervisor or listeners anywhere. Historical note: the redesign was
originally forced by smarm's then-lossy stop against QUEUED actors (fixed
upstream in 7bab4d2); it is kept because flag-based listener shutdown is
simpler and stop-semantics-free.

Conns self-register with the registry (ConnStarted from the conn actor,
not a listener cast): per-sender FIFO ordering is the only thing that
prevents ConnEnded overtaking ConnStarted; see conn_registry module docs.

conn_actor: in the catch_unwind(pipeline.run) Err branch,
smarm::preempt::check_cancelled() runs BEFORE composing the 500 — smarm's
stop is an undowncastable panic_any(StopSentinel), so the catch_unwind
would otherwise swallow a stop and 500-and-keep-running. The stop flag is
persistent, so check_cancelled re-raises cleanly.

Handle.shutdown() wakes the root via a 100ms try_recv poll
(SHUTDOWN_POLL): a cross-thread Sender::send's unpark is a
try_with_runtime no-op without runtime TLS on the sending thread (still
true on smarm 8e5b754; cross-thread unpark is a recorded smarm roadmap
candidate — this poll dies with it).

Two smarm bugs were found during this chunk and fixed upstream: lossy
stop against QUEUED actors (7bab4d2) and the terminal-wake shutdown
stall/hang (eddf3fe); post-mortem in artefact smarm-bug-terminal-wake.md.
This commit is contained in:
Claude
2026-06-11 21:45:11 +00:00
parent 5fe696992a
commit c658ac06b2
5 changed files with 564 additions and 36 deletions
+38 -1
View File
@@ -14,10 +14,13 @@
//! during those parks, other connection actors progress freely. //! during those parks, other connection actors progress freely.
use crate::conn::{Body, Conn, RespBody}; use crate::conn::{Body, Conn, RespBody};
use crate::conn_registry::{Cast, ConnRegistry, DeregisterGuard};
use crate::net::OwnedFd; use crate::net::OwnedFd;
use crate::parser::{self, ParseError}; use crate::parser::{self, ParseError};
use crate::plug::Pipeline; use crate::plug::Pipeline;
use smarm::ServerRef;
use std::io::{self, ErrorKind}; use std::io::{self, ErrorKind};
use std::os::fd::RawFd; use std::os::fd::RawFd;
use std::time::Duration; use std::time::Duration;
@@ -57,14 +60,29 @@ impl Default for ConnLimits {
// run_connection — entry point spawned by the listener actor. // run_connection — entry point spawned by the listener actor.
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
pub fn run_connection(fd: OwnedFd, pipeline: Pipeline, limits: ConnLimits) { pub fn run_connection(
fd: OwnedFd,
pipeline: Pipeline,
limits: ConnLimits,
registry: ServerRef<ConnRegistry>,
) {
// The OwnedFd cleans up via Drop on any exit path (panic, error, or // The OwnedFd cleans up via Drop on any exit path (panic, error, or
// normal close). No explicit close calls below. // normal close). No explicit close calls below.
let raw = fd.as_raw(); let raw = fd.as_raw();
let mut buf: Vec<u8> = Vec::with_capacity(limits.initial_read_buf); let mut buf: Vec<u8> = Vec::with_capacity(limits.initial_read_buf);
// 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.
let me = smarm::self_pid();
let _ = registry.cast(Cast::ConnStarted(me));
let _guard = DeregisterGuard::new(registry.clone(), me, Cast::ConnEnded);
loop { loop {
// ----- 1. Read until we have a full request head. ----- // ----- 1. Read until we have a full request head. -----
// We are idle until a head parses: stoppable by a draining
// registry while parked here.
let parsed = match read_head(raw, &mut buf, &limits) { let parsed = match read_head(raw, &mut buf, &limits) {
Ok(p) => p, Ok(p) => p,
Err(ReadHeadErr::ClientClosed) => { Err(ReadHeadErr::ClientClosed) => {
@@ -80,6 +98,7 @@ pub fn run_connection(fd: OwnedFd, pipeline: Pipeline, limits: ConnLimits) {
return; return;
} }
}; };
let _ = registry.cast(Cast::ConnBusy(me));
// ----- 2. Read body. ----- // ----- 2. Read body. -----
let body_len = parsed.content_length.unwrap_or(0); let body_len = parsed.content_length.unwrap_or(0);
@@ -117,6 +136,16 @@ pub fn run_connection(fd: OwnedFd, pipeline: Pipeline, limits: ConnLimits) {
let mut response_conn = match result { let mut response_conn = match result {
Ok(c) => c, Ok(c) => c,
Err(_) => { Err(_) => {
// Distinguish a genuine handler panic from smarm's stop
// sentinel, which is also a panic payload and which this
// catch_unwind would otherwise swallow — turning a
// graceful stop into a 500-and-keep-running. The stop
// flag is persistent (not consumed by raising the
// sentinel), so if we were stopped this re-raises it
// here, outside the catch, and we unwind properly (the
// fd and registry guards clean up).
smarm::preempt::check_cancelled();
// Compose a 500 manually; the original Conn was moved into // Compose a 500 manually; the original Conn was moved into
// the closure. // the closure.
let mut c = Conn::new(); let mut c = Conn::new();
@@ -143,6 +172,14 @@ pub fn run_connection(fd: OwnedFd, pipeline: Pipeline, limits: ConnLimits) {
return; return;
} }
// Response is on the wire; nothing is owed. Going idle here makes
// us stoppable by a draining registry while we park for the next
// keep-alive request. (If the next request is already pipelined in
// `buf`, the very next read_head parses it without parking and we
// go Busy again — a draining registry's request_stop may still
// catch us, which is acceptable: drain means no new work.)
let _ = registry.cast(Cast::ConnIdle(me));
// Drop the request bytes (head + body) from `buf`; anything past // Drop the request bytes (head + body) from `buf`; anything past
// them is the start of the next pipelined request. // them is the start of the next pipelined request.
let consumed = head_len + body_len; let consumed = head_len + body_len;
+167
View File
@@ -0,0 +1,167 @@
//! Connection registry — the shutdown coordinator (v0.2 chunk 2).
//!
//! 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.)
//!
//! 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: `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.
//!
//! 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, ServerBuilder, ServerRef};
use std::collections::HashMap;
// ---------------------------------------------------------------------------
// 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),
/// 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 {
ConnCount,
}
pub enum Reply {
ConnCount(usize),
}
// ---------------------------------------------------------------------------
// Server
// ---------------------------------------------------------------------------
#[derive(Clone, Copy, PartialEq, Eq)]
enum ConnState {
Busy,
Idle,
}
#[derive(Default)]
pub struct ConnRegistry {
conns: HashMap<Pid, ConnState>,
draining: bool,
}
impl GenServer for ConnRegistry {
type Call = Call;
type Reply = Reply;
type Cast = Cast;
type Info = ();
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.
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. The pid leaves the map via its drop
// guard's ConnEnded once the unwind completes.
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);
}
}
}
}
}
pub fn start() -> ServerRef<ConnRegistry> {
ServerBuilder::new(ConnRegistry::default()).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: ServerRef<ConnRegistry>,
pid: Pid,
make: fn(Pid) -> Cast,
}
impl DeregisterGuard {
pub fn new(registry: ServerRef<ConnRegistry>, 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));
}
}
+4 -1
View File
@@ -24,10 +24,13 @@ pub mod router;
pub mod parser; pub mod parser;
pub mod net; pub mod net;
pub mod conn_actor; pub mod conn_actor;
pub mod conn_registry;
pub mod serve; pub mod serve;
// Re-exports — what most users want at the crate root. // Re-exports — what most users want at the crate root.
pub use conn::{Assigns, Body, Conn, HeaderMap, HttpVersion, Method, Params, RespBody}; pub use conn::{Assigns, Body, Conn, HeaderMap, HttpVersion, Method, Params, RespBody};
pub use plug::{Next, Pipeline, Plug}; pub use plug::{Next, Pipeline, Plug};
pub use router::Router; pub use router::Router;
pub use serve::{serve, serve_with, Config}; pub use serve::{
serve, serve_with, serve_with_shutdown, shutdown_handle, Config, Handle, ShutdownSignal,
};
+214 -34
View File
@@ -13,15 +13,16 @@
//! waiting on "the same fd"). //! waiting on "the same fd").
use crate::conn_actor::{run_connection, ConnLimits}; use crate::conn_actor::{run_connection, ConnLimits};
use crate::conn_registry::{self, Call, Cast, ConnRegistry, Reply};
use crate::net::{accept_nonblocking, bind_and_listen, OwnedFd}; use crate::net::{accept_nonblocking, bind_and_listen, OwnedFd};
use crate::plug::Pipeline; use crate::plug::Pipeline;
use smarm::{ChildSpec, OneForOne, Restart, Strategy}; use smarm::{ChildSpec, OneForOne, Restart, ServerRef, Strategy};
use std::io::{self, ErrorKind}; use std::io::{self, ErrorKind};
use std::net::{SocketAddr, ToSocketAddrs}; use std::net::{SocketAddr, ToSocketAddrs};
use std::os::fd::RawFd; use std::os::fd::RawFd;
use std::sync::atomic::{AtomicU32, Ordering}; use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::sync::Arc; use std::sync::Arc;
use std::time::Duration; use std::time::Duration;
@@ -38,6 +39,11 @@ pub struct Config {
pub read_buf_size: usize, pub read_buf_size: usize,
pub request_timeout: Duration, pub request_timeout: Duration,
pub max_body_bytes: usize, pub max_body_bytes: usize,
/// How long a graceful shutdown waits for in-flight requests before
/// force-stopping the remaining connections. Idle keep-alive
/// connections are closed immediately on shutdown and do not run the
/// clock out.
pub drain_timeout: Duration,
/// Number of smarm scheduler OS threads. `None` means smarm's default /// 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 /// (one per CPU). Set this to a small fixed number in tests so multiple
/// concurrent test servers don't oversubscribe the host. /// concurrent test servers don't oversubscribe the host.
@@ -58,6 +64,7 @@ impl Config {
read_buf_size: 8 * 1024, read_buf_size: 8 * 1024,
request_timeout: Duration::from_secs(30), request_timeout: Duration::from_secs(30),
max_body_bytes: 16 * 1024 * 1024, max_body_bytes: 16 * 1024 * 1024,
drain_timeout: Duration::from_secs(30),
scheduler_threads: None, scheduler_threads: None,
} }
} }
@@ -98,9 +105,28 @@ fn dup_fd(fd: RawFd) -> io::Result<OwnedFd> {
#[doc(hidden)] #[doc(hidden)]
pub static INJECT_LISTENER_PANICS: AtomicU32 = AtomicU32::new(0); pub static INJECT_LISTENER_PANICS: AtomicU32 = AtomicU32::new(0);
fn listener_loop(listener: Arc<OwnedFd>, pipeline: Pipeline, limits: ConnLimits) { fn listener_loop(
listener: Arc<OwnedFd>,
pipeline: Pipeline,
limits: ConnLimits,
registry: ServerRef<ConnRegistry>,
shutdown: Arc<AtomicBool>,
) {
let fd = listener.as_raw(); let fd = listener.as_raw();
loop { 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 if INJECT_LISTENER_PANICS
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| n.checked_sub(1)) .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| n.checked_sub(1))
.is_ok() .is_ok()
@@ -114,19 +140,27 @@ fn listener_loop(listener: Arc<OwnedFd>, pipeline: Pipeline, limits: ConnLimits)
// shared lock. // shared lock.
let p = pipeline.clone(); let p = pipeline.clone();
let l = limits; let l = limits;
smarm::spawn(move || run_connection(client, p, l)); let r = registry.clone();
smarm::spawn(move || run_connection(client, p, l, r));
} }
Err(e) if e.kind() == ErrorKind::WouldBlock => { Err(e) if e.kind() == ErrorKind::WouldBlock => {
// No pending connection. Park until the listener is // No pending connection. Park until the listener is
// readable again, then retry. // readable or the tick elapses; either way we come back
if let Err(we) = smarm::wait_readable(fd) { // around through the shutdown-flag check above.
// epoll registration failed. Under supervision a match smarm::wait_readable_timeout(fd, LISTENER_TICK) {
// `return` is a `Signal::Exit` and `Restart::Permanent` Ok(_ready) => {} // ready or tick — loop re-checks, retries accept
// restarts us — a transient failure (e.g. EMFILE on Err(we) => {
// the epoll set) now heals instead of silently // epoll registration failed. Under Transient
// shrinking the pool. // supervision a normal return is terminal but a
eprintln!("urus: listener wait_readable failed: {we}"); // panic restarts us — and a failed wait IS
return; // 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 => { Err(e) if e.kind() == ErrorKind::Interrupted => {
@@ -148,24 +182,97 @@ fn listener_loop(listener: Arc<OwnedFd>, pipeline: Pipeline, limits: ConnLimits)
} }
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// serve_with — main entry. Boots smarm, supervises listeners, blocks. // 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.
#[derive(Clone)]
pub struct Handle {
tx: smarm::Sender<()>,
}
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.
pub fn shutdown(&self) {
let _ = self.tx.send(());
}
}
/// The receiving half consumed by [`serve_with_shutdown`].
pub struct ShutdownSignal {
rx: smarm::Receiver<()>,
}
/// Create a connected [`Handle`]/[`ShutdownSignal`] pair.
pub fn shutdown_handle() -> (Handle, ShutdownSignal) {
let (tx, rx) = smarm::channel();
(Handle { tx }, ShutdownSignal { rx })
}
// ---------------------------------------------------------------------------
// serve_with_shutdown — main entry. Boots smarm, supervises listeners,
// blocks until shutdown.
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// //
// Boots an smarm runtime (one OS thread per CPU by default — see smarm's // Boots an smarm runtime (one OS thread per CPU by default — see smarm's
// `Config::default()`) and runs a one-for-one supervisor over the listener // `Config::default()`). The root actor starts the connection registry and
// pool on the root actor. A panicking listener is restarted on the same // a one-for-one supervisor over the listener pool, then blocks on the
// (still-open) fd instead of silently shrinking the accept pool. Connection // shutdown signal. A panicking listener is restarted on the same
// actors stay unsupervised bare spawns — per spec, a connection is cheap // (still-open) fd instead of silently shrinking the accept pool.
// and its failure is local: a 500 path, not a restart path. (`spawn` from a // Connection actors stay unsupervised bare spawns — per spec, a connection
// listener parents the conn actor under that listener, which never // is cheap and its failure is local: a 500 path, not a restart path.
// registers a supervisor channel, so conn deaths are invisible to the pool // (`spawn` from a listener parents the conn actor under that listener,
// supervisor by construction.) // which never registers a supervisor channel, so conn deaths are invisible
// to the pool supervisor by construction.)
// //
// No graceful-shutdown signal yet (v0.2 chunk 2); the supervisor loop never // Shutdown sequence (drain-then-stop, the v0.2 chunk 2 decision):
// returns because listeners never reach a terminal state. Ctrl-C is your // 1. Set the shared shutdown flag. Every listener observes it at the
// friend. // 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).
pub fn serve_with(config: Config, pipeline: Pipeline) -> io::Result<()> { pub fn serve_with_shutdown(
config: Config,
pipeline: Pipeline,
signal: ShutdownSignal,
) -> io::Result<()> {
let listener = bind_and_listen(config.addr)?; let listener = bind_and_listen(config.addr)?;
println!("urus: listening on {}", config.addr); println!("urus: listening on {}", config.addr);
@@ -182,32 +289,105 @@ pub fn serve_with(config: Config, pipeline: Pipeline) -> io::Result<()> {
listener_fds.push(Arc::new(dup)); listener_fds.push(Arc::new(dup));
} }
let limits = config.to_conn_limits(); let limits = config.to_conn_limits();
let drain_timeout = config.drain_timeout;
let smarm_cfg = match config.scheduler_threads { let smarm_cfg = match config.scheduler_threads {
Some(n) => smarm::Config::exact(n), Some(n) => smarm::Config::exact(n),
None => smarm::Config::default(), None => smarm::Config::default(),
}; };
let rt = smarm::init(smarm_cfg); let rt = smarm::init(smarm_cfg);
// Listener self-termination flag — see the shutdown sequence below.
let shutdown_flag = Arc::new(AtomicBool::new(false));
rt.run(move || { rt.run(move || {
// Registry first: listeners and conns cast into it from birth.
let registry = conn_registry::start();
let mut sup = OneForOne::new().strategy(Strategy::OneForOne); let mut sup = OneForOne::new().strategy(Strategy::OneForOne);
for (i, lfd) in listener_fds.into_iter().enumerate() { for (i, lfd) in listener_fds.into_iter().enumerate() {
let p = pipeline.clone(); let p = pipeline.clone();
sup = sup.child(ChildSpec::new(Restart::Permanent, move || { let r = registry.clone();
let sf = shutdown_flag.clone();
sup = sup.child(ChildSpec::new(Restart::Transient, move || {
println!("urus: listener {} starting", i); println!("urus: listener {} starting", i);
listener_loop(lfd.clone(), p.clone(), limits); listener_loop(lfd.clone(), p.clone(), limits, r.clone(), sf.clone());
})); }));
} }
// Supervision loop on the root actor. Default intensity (3 per 5s) // Default intensity (3 per 5s) applies; a listener crash-looping
// applies per the runtime; a listener crash-looping faster than // faster than that trips the cap and tears the pool down — loud
// that trips the cap and tears the pool down — loud failure over a // failure over a zombie server.
// zombie server. let sup_h = smarm::spawn(move || sup.run());
sup.run();
// 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)),
}
}
// ----- 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. 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));
}
}
// 5. Our ServerRef 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.
}); });
Ok(()) 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).
let (_handle, signal) = shutdown_handle();
serve_with_shutdown(config, pipeline, signal)
}
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// serve — convenience over serve_with. // serve — convenience over serve_with.
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
+141
View File
@@ -284,3 +284,144 @@ fn panicking_listener_restarts() {
assert_eq!(http_status(&resp), 200); assert_eq!(http_status(&resp), 200);
} }
} }
// ---------------------------------------------------------------------------
// Graceful shutdown (v0.2 chunk 2)
// ---------------------------------------------------------------------------
/// Boot a server with a shutdown handle. Returns (port, handle, done_rx)
/// where done_rx fires when serve_with_shutdown returns.
fn spawn_server_with_handle(
pipeline: Pipeline,
drain: Duration,
) -> (u16, urus::Handle, std::sync::mpsc::Receiver<()>) {
let port = free_port();
let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap();
let (handle, signal) = urus::shutdown_handle();
let (done_tx, done_rx) = std::sync::mpsc::channel();
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();
let _ = done_tx.send(());
});
for _ in 0..50 {
if TcpStream::connect(addr).is_ok() {
return (port, handle, done_rx);
}
std::thread::sleep(Duration::from_millis(50));
}
panic!("server didn't come up on {addr}");
}
/// Shutdown with nothing in flight returns promptly — well inside the
/// (deliberately huge) drain window, proving idle conns don't run the
/// clock out.
#[test]
fn shutdown_with_no_connections_returns() {
let pipe = Pipeline::new().plug(
Router::new().get("/", |c: Conn, _n: Next| c.put_status(200))
);
let (_port, handle, done_rx) = spawn_server_with_handle(pipe, Duration::from_secs(30));
handle.shutdown();
done_rx
.recv_timeout(Duration::from_secs(5))
.expect("serve_with_shutdown did not return after shutdown()");
}
/// An idle keep-alive connection is closed immediately on shutdown (the
/// drain deadline of 30s must NOT be what gates the return), and the
/// server refuses new connections afterwards.
#[test]
fn shutdown_closes_idle_keepalive_promptly() {
let pipe = Pipeline::new().plug(
Router::new().get("/", |c: Conn, _n: Next| c.put_status(200).put_body("ok"))
);
let (port, handle, done_rx) = spawn_server_with_handle(pipe, Duration::from_secs(30));
// Complete one request on a keep-alive connection and leave it open.
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
s.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
s.write_all(b"GET / HTTP/1.1\r\nHost: x\r\n\r\n").unwrap();
let mut buf = [0u8; 1024];
let n = s.read(&mut buf).unwrap();
assert!(n > 0 && buf.starts_with(b"HTTP/1.1 200"));
let t0 = std::time::Instant::now();
handle.shutdown();
// The parked conn actor is stopped; its socket closes → we see EOF.
let n = s.read(&mut buf).expect("read after shutdown");
assert_eq!(n, 0, "expected EOF on the idle keep-alive connection");
done_rx
.recv_timeout(Duration::from_secs(5))
.expect("serve did not return");
assert!(
t0.elapsed() < Duration::from_secs(5),
"shutdown waited for the drain deadline instead of stopping the idle conn"
);
// Listeners are gone: new connections are refused.
assert!(
TcpStream::connect(("127.0.0.1", port)).is_err(),
"server still accepting after shutdown"
);
}
/// A request in flight when shutdown is signalled completes with a real
/// response before its connection is closed.
#[test]
fn shutdown_drains_in_flight_request() {
let pipe = Pipeline::new().plug(
Router::new().get("/slow", |c: Conn, _n: Next| {
smarm::sleep(Duration::from_millis(400));
c.put_status(200).put_body("made it")
})
);
let (port, handle, done_rx) = spawn_server_with_handle(pipe, Duration::from_secs(10));
// Fire the slow request, then shut down while it is in flight.
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
s.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
s.write_all(b"GET /slow HTTP/1.1\r\nHost: x\r\nConnection: close\r\n\r\n").unwrap();
std::thread::sleep(Duration::from_millis(100)); // let the head land
handle.shutdown();
let mut resp = Vec::new();
s.read_to_end(&mut resp).unwrap();
assert_eq!(http_status(&resp), 200, "in-flight request was cut off by shutdown");
assert_eq!(http_body(&resp), b"made it");
done_rx
.recv_timeout(Duration::from_secs(5))
.expect("serve did not return after drain");
}
/// A request still in flight at the drain deadline is force-stopped: the
/// connection dies without a response, but serve returns near the
/// deadline rather than hanging on the stuck handler's full duration.
#[test]
fn shutdown_force_stops_at_drain_deadline() {
let pipe = Pipeline::new().plug(
Router::new().get("/stuck", |c: Conn, _n: Next| {
smarm::sleep(Duration::from_secs(60));
c.put_status(200)
})
);
let (port, handle, done_rx) = spawn_server_with_handle(pipe, Duration::from_millis(300));
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
s.set_read_timeout(Some(Duration::from_secs(10))).unwrap();
s.write_all(b"GET /stuck HTTP/1.1\r\nHost: x\r\nConnection: close\r\n\r\n").unwrap();
std::thread::sleep(Duration::from_millis(100));
handle.shutdown();
done_rx
.recv_timeout(Duration::from_secs(5))
.expect("serve did not return after force-stop deadline");
}