feat(endpoint): urus is a supervisable child — endpoint gen_server owns listeners + conns
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.
This commit is contained in:
+3
-3
@@ -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<ConnRegistry>,
|
||||
registry: GenServerRef<Endpoint>,
|
||||
) {
|
||||
// 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);
|
||||
|
||||
@@ -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<Pid, ConnState>,
|
||||
draining: bool,
|
||||
drain_timeout: Duration,
|
||||
stop: Option<StopHandle<Self>>,
|
||||
timer: Option<TimerHandle<Self>>,
|
||||
}
|
||||
|
||||
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<Self>) {
|
||||
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<ConnRegistry> {
|
||||
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<ConnRegistry>,
|
||||
pid: Pid,
|
||||
make: fn(Pid) -> Cast,
|
||||
}
|
||||
|
||||
impl DeregisterGuard {
|
||||
pub fn new(registry: GenServerRef<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));
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// 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<ConnRegistry>, 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<ConnRegistry>) -> Result<usize, ()> {
|
||||
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()));
|
||||
});
|
||||
}
|
||||
}
|
||||
+656
@@ -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<Arc<OwnedFd>>,
|
||||
pipeline: Pipeline,
|
||||
limits: ConnLimits,
|
||||
conn_stack_reserve: usize,
|
||||
name: &'static str,
|
||||
}
|
||||
|
||||
pub struct Endpoint {
|
||||
boot: Arc<Boot>,
|
||||
conns: HashMap<Pid, ConnState>,
|
||||
draining: bool,
|
||||
drain_timeout: Duration,
|
||||
/// The listener supervisor, spawned in `init` and monitored. `None`
|
||||
/// once it is down.
|
||||
listener_sup: Option<Pid>,
|
||||
stop: Option<StopHandle<Self>>,
|
||||
timer: Option<TimerHandle<Self>>,
|
||||
}
|
||||
|
||||
impl Endpoint {
|
||||
fn name(&self) -> GenServerName<Self> {
|
||||
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<Self>) {
|
||||
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<Self> =
|
||||
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<impl Fn() + Clone + Send + Sync + 'static> {
|
||||
// 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<RawFd, Pid>`.
|
||||
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::<Endpoint>::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<GenServerRef<Endpoint>> {
|
||||
smarm::gen_server::whereis_server(GenServerName::<Endpoint>::new(name))
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Listener actors
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn dup_fd(fd: RawFd) -> io::Result<OwnedFd> {
|
||||
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<OwnedFd>,
|
||||
pipeline: Pipeline,
|
||||
limits: ConnLimits,
|
||||
conn_stack_reserve: usize,
|
||||
endpoint: GenServerRef<Endpoint>,
|
||||
) {
|
||||
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<Endpoint>,
|
||||
pid: Pid,
|
||||
make: fn(Pid) -> Cast,
|
||||
}
|
||||
|
||||
impl DeregisterGuard {
|
||||
pub fn new(endpoint: GenServerRef<Endpoint>, 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<Endpoint>, 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<usize, ()> {
|
||||
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));
|
||||
});
|
||||
}
|
||||
}
|
||||
+2
-1
@@ -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,
|
||||
};
|
||||
|
||||
+69
-294
@@ -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<RawFd, Pid>` 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<usize>,
|
||||
/// 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<usize>,
|
||||
}
|
||||
|
||||
/// 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<OwnedFd> {
|
||||
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<OwnedFd>,
|
||||
pipeline: Pipeline,
|
||||
limits: ConnLimits,
|
||||
conn_stack_reserve: usize,
|
||||
registry: GenServerRef<ConnRegistry>,
|
||||
shutdown: Arc<AtomicBool>,
|
||||
) {
|
||||
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<M>, Sender<M>)`, 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.
|
||||
// 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 {
|
||||
match signal.rx.try_recv() {
|
||||
Ok(Some(())) => break,
|
||||
Ok(None) => smarm::sleep(SHUTDOWN_POLL),
|
||||
Err(_) => smarm::sleep(Duration::from_secs(3600)),
|
||||
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::*;
|
||||
|
||||
+109
-16
@@ -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<AtomicBool>);
|
||||
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"
|
||||
);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user