From b37888ec2cfc7e1c5593a11097ec067a700b3d80 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 20 Aug 2026 13:20:28 +0000 Subject: [PATCH] docs+examples: v0.3 endpoint API; crud converted to the app-owns-the-tree shape - crud is now the demonstrator: app owns smarm runtime + root supervisor, store actor and urus::endpoint as ordered siblings (store first, so reverse-order shutdown drains HTTP before stopping the store), Shutdown::Infinity on the endpoint, shutdown via rt.handle().request_shutdown(root_sup) from the stdin thread. Deletes two kludges the old shape forced: * static OnceLock spawn-on-first-use -> a supervised child that self-registers a typed Name; handlers use smarm::send per request and turn 'between incarnations' into a 503 instead of panicking on a dropped store. * static SHUTTING_DOWN AtomicBool + 250ms recv_timeout poll in the store loop and in the SSE ticker -> a plain park; the tree stops both. Smoke-tested live: CRUD round-trips, SSE stream, clean drain with the stream open, port closed after. - Other examples stay short and on serve*, updated for the split config (serve_with(cfg, smarm::Config, pipe) / serve_with_shutdown(..., signal)). plain_serve's URUS_SCHED_THREADS now builds a smarm::Config. ws_chat's doc block explains the OnceLock is a serve*-only workaround and points at crud for the clean shape. - README: new 'Your Own Supervision Tree' section (endpoint as the real API), graceful shutdown reframed as the serve*-only path, Config table loses scheduler_threads and gains name, PubSub rule 1 notes the supervised-sibling alternative. 111 lib + 50 integration + 2 doc tests green; clippy clean. --- README.md | 68 ++++++++++++++--- examples/causal_bench.rs | 2 +- examples/channels_chat.rs | 7 +- examples/crud.rs | 156 +++++++++++++++++++++----------------- examples/load_profile.rs | 2 +- examples/plain_serve.rs | 11 ++- examples/serve_toml.rs | 2 +- examples/ws_chat.rs | 26 ++++--- examples/ws_echo.rs | 2 +- 9 files changed, 181 insertions(+), 95 deletions(-) diff --git a/README.md b/README.md index 6dc71fa..d0bd9af 100644 --- a/README.md +++ b/README.md @@ -122,7 +122,11 @@ Handlers are closures taking `(Conn, Next) -> Conn`. Call `Next::call(c)` to con ### Server Configuration -`serve(addr, pipeline)` binds and listens on the given address. For more control, use `serve_with(config, pipeline)`: +`serve(addr, pipeline)` binds and listens on the given address. For more +control, use `serve_with(config, runtime_config, pipeline)` — the urus +`Config` holds endpoint knobs, the `smarm::Config` holds runtime knobs +(they are separate because an endpoint placed in someone else's tree +cannot dictate the runtime): ```rust use urus::{serve_with, Config}; @@ -130,7 +134,7 @@ use std::time::Duration; let cfg = Config { listener_pool: 2, // Supervised accept-loop actors - scheduler_threads: Some(2), // smarm worker threads (None = one per CPU) + name: "urus", // Endpoint registry name (unique per endpoint) keep_alive_timeout: Duration::from_secs(60), // Idle budget between requests request_timeout: Duration::from_secs(30), // Whole-request READ deadline write_timeout: Duration::from_secs(30), // Per-write response budget @@ -141,7 +145,7 @@ let cfg = Config { ..Config::new("127.0.0.1:8080".parse().unwrap()) }; -serve_with(cfg, pipeline).unwrap(); +serve_with(cfg, smarm::Config::exact(2), pipeline).unwrap(); ``` **Timeout semantics:** @@ -162,10 +166,51 @@ serve_with(cfg, pipeline).unwrap(); write stalls past the budget (the write-side twin of slowloris). Streamed chunks each get a fresh budget; a stream as a whole has no deadline. -### Graceful Shutdown +### Your Own Supervision Tree -`serve_with_shutdown` takes a `ShutdownSignal`; the paired `Handle` can be -triggered from anywhere (another thread, a signal handler): +The real API is `urus::endpoint(config, pipeline)`, which binds the socket +and hands back a supervisable child body. Your application owns the +runtime and the root supervisor; urus is one child among your own actors: + +```rust +use smarm::{ChildSpec, OneForOne, Restart, supervisor::Shutdown}; + +// Binds here: an address-in-use error is a startup failure, not an actor +// crash. The fds outlive any restart of the endpoint child. +let endpoint = urus::endpoint(cfg, pipeline)?; + +let rt = smarm::init(smarm::Config::default()); +rt.run(move || { + let sup = smarm::spawn(move || { + OneForOne::new() + // Your state actor FIRST: reverse-order shutdown therefore + // stops it LAST, after HTTP has finished draining. + .child(ChildSpec::new(Restart::Permanent, my_store)) + .child(ChildSpec::new(Restart::Permanent, endpoint) + // The endpoint bounds its own drain with drain_timeout; + // a shorter supervisor deadline would cut it in half. + .shutdown(Shutdown::Infinity)) + .run() + }); + let _ = sup.join(); +}); +``` + +Shutdown is then whatever your app already does — `request_shutdown` on +the root supervisor (from a SIGTERM handler via `rt.handle()`, say). The +endpoint's `handle_shutdown` stops the listener pool, closes idle +keep-alive connections, drains in-flight requests up to `drain_timeout`, +force-stops stragglers past it, and exits only when the last connection is +gone — so "the endpoint child has stopped" *is* "every connection is +gone". See `examples/crud.rs` for a complete app in this shape. + +Address a running endpoint by name for introspection: +`urus::endpoint::whereis("urus")`. + +### Graceful Shutdown Without a Tree + +If serving is all your process does, let `serve*` own the runtime and use +`Handle`, which triggers the same sequence from any thread: ```rust use urus::{serve_with_shutdown, shutdown_handle, Config}; @@ -177,13 +222,13 @@ std::thread::spawn(move || { handle.shutdown(); }); -serve_with_shutdown(cfg, pipeline, signal).unwrap(); +serve_with_shutdown(cfg, smarm::Config::default(), pipeline, signal).unwrap(); // Returns once the runtime has fully wound down. ``` `Handle::shutdown()` is idempotent and performs, in order: -1. Stop accepting — every listener exits; no new connections. +1. Stop accepting — the listener pool is stopped; no new connections. 2. Close idle keep-alive connections immediately. 3. Drain in-flight requests for up to `Config.drain_timeout`. 4. Force-stop any stragglers past the deadline (sockets close cleanly on @@ -331,9 +376,12 @@ by the two cleanup paths above. Two composition rules that matter (both enforced by `shutdown_with_open_chat_terminates` in the integration suite): -1. `PubSub::new()` spawns an actor, so it must run in-runtime — build it +1. `PubSub::new()` spawns an actor, so it must run in-runtime. If your + app owns its tree, start the table as a supervised sibling of the + endpoint and address it by name — the clean shape. Under `serve*` + there is no in-runtime moment before the first request, so build it lazily from a handler via a **non-static** `Arc>>` - captured by the route closure. A `static` cell pins the table actor + captured by the route closure; a `static` cell pins the table actor forever and graceful shutdown never returns. 2. Relay/producer actors hold the `Receiver` (plus e.g. a `WsSender` clone) — **never a `PubSub` clone**, or relay and table keep each diff --git a/examples/causal_bench.rs b/examples/causal_bench.rs index 9a84fea..290d139 100644 --- a/examples/causal_bench.rs +++ b/examples/causal_bench.rs @@ -240,7 +240,7 @@ fn main() { let (handle, signal) = shutdown_handle(); let server = std::thread::spawn(move || { let pipe = Pipeline::new().plug(Router::new().get("/order/:id", order)); - serve_with_shutdown(Config::new(addr), pipe, signal).expect("serve"); + serve_with_shutdown(Config::new(addr), smarm::Config::default(), pipe, signal).expect("serve"); }); // Wait until it's accepting. diff --git a/examples/channels_chat.rs b/examples/channels_chat.rs index fa96dd6..48ef5f0 100644 --- a/examples/channels_chat.rs +++ b/examples/channels_chat.rs @@ -111,7 +111,12 @@ fn main() { handle.shutdown(); }); - serve_with_shutdown(Config::new("0.0.0.0:8080".parse().unwrap()), Pipeline::new().plug(router), signal) + serve_with_shutdown( + Config::new("0.0.0.0:8080".parse().unwrap()), + smarm::Config::default(), + Pipeline::new().plug(router), + signal, + ) .unwrap(); println!("channels_chat: drained, bye"); } diff --git a/examples/crud.rs b/examples/crud.rs index 55de4ae..b024881 100644 --- a/examples/crud.rs +++ b/examples/crud.rs @@ -1,11 +1,18 @@ //! CRUD example: a tiny user database with JSON persistence. //! -//! Demonstrates urus and the actor model together: -//! - The pipeline is shared (Arc) across all connection actors. -//! - Handlers do NOT take a lock or share mutable state directly. -//! - A single "store" actor owns the data; handlers send it a request -//! via a channel and block on the reply. Serialization is structural — -//! the store processes one request at a time, no Mutex needed. +//! Demonstrates urus and the actor model together, in the shape a real +//! application should use (v0.3): +//! - The APP owns the smarm runtime and the root supervisor. urus is one +//! child in that tree — `urus::endpoint(...)` — and the store actor is +//! an ordered sibling started BEFORE it, so the supervisor's +//! reverse-order shutdown drains HTTP first and only then stops the +//! store. No handler can be mid-request against a store that is gone. +//! - Handlers do NOT take a lock or share mutable state directly. A +//! single "store" actor owns the data; handlers address it by +//! registered name and block on a reply channel. Serialization is +//! structural — one request at a time, no Mutex. +//! - Both actors are supervised: kill the store (or let it panic) and it +//! restarts from the JSON file, with the endpoint left alone. //! - On every mutating request the store writes the JSON file. Read //! requests don't touch disk. //! @@ -23,9 +30,8 @@ //! curl -s http://localhost:8080/users/1 use serde::{Deserialize, Serialize}; -use smarm::{channel, Sender}; -use std::sync::OnceLock; -use urus::{serve_with_shutdown, shutdown_handle, Config, Conn, Next, Pipeline, Router}; +use smarm::{channel, ChildSpec, Name, OneForOne, Restart, Sender}; +use urus::{Config, Conn, Next, Pipeline, Router}; // --------------------------------------------------------------------------- // Domain @@ -66,7 +72,12 @@ const DB_PATH: &str = "/tmp/urus-crud.json"; // Store actor body // --------------------------------------------------------------------------- -fn store_loop(rx: smarm::Receiver) { +fn store_loop() { + let (tx, rx) = channel::(); + // Self-registration: the name is bound before the first recv, and it + // is re-bound automatically on every restart. + smarm::register(STORE, tx).expect("crud.store name already taken"); + // Load on start. Missing file = empty store. Corrupt file = panic; we // don't auto-rebuild because silently losing data is worse than failing // loud. @@ -78,19 +89,10 @@ fn store_loop(rx: smarm::Receiver) { let mut next_id: u64 = users.iter().map(|u| u.id).max().unwrap_or(0) + 1; loop { - // recv with a timeout rather than a bare recv: the store must be - // stoppable at shutdown, but a cross-thread Sender::send (from the - // stdin thread) can't wake a parked actor — same smarm limitation - // that motivates urus's SHUTDOWN_POLL. So we wake on our own timer - // and poll the flag. This poll dies with that limitation too. - let req = match rx.recv_timeout(std::time::Duration::from_millis(250)) { + // A plain park. The supervisor stops this actor at shutdown (after + // the endpoint has drained), so there is nothing to poll for. + let req = match rx.recv() { Ok(r) => r, - Err(smarm::channel::RecvTimeoutError::Timeout) => { - if SHUTTING_DOWN.load(std::sync::atomic::Ordering::Relaxed) { - return; - } - continue; - } Err(_) => return, // all senders dropped }; match req { @@ -169,25 +171,29 @@ fn persist(users: &[User]) { // Handler helpers // --------------------------------------------------------------------------- // -// Once-cell trick: the store actor is spawned the first time a handler -// runs (smarm requires `spawn` to be called from inside an actor — which -// connection actors are). After that all handlers share the same Sender. -// Simpler than threading the Sender through the pipeline at startup. +// The store is a supervised child that registers its own inbox under a +// typed name; handlers resolve it per send. That replaces the old +// `OnceLock` spawn-on-first-use trick — which had no supervisor, +// no restart, and no defined shutdown point — with a plain actor whose +// lifecycle the tree owns. A restart re-registers the same name, so +// in-flight handlers heal on their next send. -static STORE_TX: OnceLock> = OnceLock::new(); +const STORE: Name = Name::new("crud.store"); -// Set by the stdin thread at shutdown; the store actor polls it (see -// store_loop). An always-on app actor that never returns would otherwise -// block smarm's AllDone and keep serve_with_shutdown from returning. -static SHUTTING_DOWN: std::sync::atomic::AtomicBool = - std::sync::atomic::AtomicBool::new(false); +/// Send to the store and wait for its reply. `Err` only if the store is +/// between incarnations (restarting); the handler turns that into a 503 +/// rather than pretending. +fn ask(make: impl FnOnce(Sender<(u16, Vec)>) -> Request) -> Option<(u16, Vec)> { + let (tx, rx) = channel::<(u16, Vec)>(); + smarm::send(STORE, make(tx)).ok()?; + rx.recv().ok() +} -fn store() -> &'static Sender { - STORE_TX.get_or_init(|| { - let (tx, rx) = channel::(); - smarm::spawn(move || store_loop(rx)); - tx - }) +fn reply(conn: Conn, r: Option<(u16, Vec)>) -> Conn { + match r { + Some((status, body)) => json(conn, status, body), + None => json(conn, 503, b"{\"error\":\"store unavailable\"}".to_vec()), + } } fn json(conn: Conn, status: u16, body: Vec) -> Conn { @@ -206,16 +212,13 @@ fn parse_id(s: &str) -> Option { /// SSE demo: `curl -N localhost:8080/ticker` streams a tick every second /// (with `: keep-alive` comments if it ever goes quiet). The producer -/// exits on SseClosed (client gone / write timeout / shutdown drain) or -/// when the example is shutting down. +/// exits on SseClosed — client gone, write timeout, or the drain stopping +/// its connection actor. Nothing to flag: the tree's shutdown reaches it. fn ticker(conn: Conn, _next: Next) -> Conn { let (conn, events) = conn.sse(); smarm::spawn(move || { let mut n: u64 = 0; loop { - if SHUTTING_DOWN.load(std::sync::atomic::Ordering::Relaxed) { - return; // dropping `events` ends the stream cleanly - } if events.send("tick", &n.to_string()).is_err() { return; // SseClosed } @@ -227,18 +230,14 @@ fn ticker(conn: Conn, _next: Next) -> Conn { } fn list(conn: Conn, _next: Next) -> Conn { - let (tx, rx) = channel::<(u16, Vec)>(); - store().send(Request::List { reply: tx }).ok(); - let (status, body) = rx.recv().expect("store dropped"); - json(conn, status, body) + let r = ask(|reply| Request::List { reply }); + reply(conn, r) } fn create(conn: Conn, _next: Next) -> Conn { let body = conn.body.as_bytes().to_vec(); - let (tx, rx) = channel::<(u16, Vec)>(); - store().send(Request::Create { body, reply: tx }).ok(); - let (status, body) = rx.recv().expect("store dropped"); - json(conn, status, body) + let r = ask(|reply| Request::Create { body, reply }); + reply(conn, r) } fn get_one(conn: Conn, _next: Next) -> Conn { @@ -246,10 +245,8 @@ fn get_one(conn: Conn, _next: Next) -> Conn { Some(id) => id, None => return json(conn, 400, b"{\"error\":\"bad id\"}".to_vec()), }; - let (tx, rx) = channel::<(u16, Vec)>(); - store().send(Request::Get { id, reply: tx }).ok(); - let (status, body) = rx.recv().expect("store dropped"); - json(conn, status, body) + let r = ask(|reply| Request::Get { id, reply }); + reply(conn, r) } fn update(conn: Conn, _next: Next) -> Conn { @@ -258,10 +255,8 @@ fn update(conn: Conn, _next: Next) -> Conn { None => return json(conn, 400, b"{\"error\":\"bad id\"}".to_vec()), }; let body = conn.body.as_bytes().to_vec(); - let (tx, rx) = channel::<(u16, Vec)>(); - store().send(Request::Update { id, body, reply: tx }).ok(); - let (status, body) = rx.recv().expect("store dropped"); - json(conn, status, body) + let r = ask(|reply| Request::Update { id, body, reply }); + reply(conn, r) } fn delete(conn: Conn, _next: Next) -> Conn { @@ -269,10 +264,8 @@ fn delete(conn: Conn, _next: Next) -> Conn { Some(id) => id, None => return json(conn, 400, b"{\"error\":\"bad id\"}".to_vec()), }; - let (tx, rx) = channel::<(u16, Vec)>(); - store().send(Request::Delete { id, reply: tx }).ok(); - let (status, body) = rx.recv().expect("store dropped"); - json(conn, status, body) + let r = ask(|reply| Request::Delete { id, reply }); + reply(conn, r) } // --------------------------------------------------------------------------- @@ -305,20 +298,47 @@ fn main() { ); let cfg = Config::new("127.0.0.1:8080".parse().unwrap()); + // Bind here, on this thread: an address-in-use error is a startup + // failure, not an actor crash. The fds outlive any restart of the + // endpoint child. + let endpoint = urus::endpoint(cfg, pipeline).expect("bind 127.0.0.1:8080"); + println!("urus-crud: DB at {DB_PATH}"); println!("urus-crud: listening on 127.0.0.1:8080 — press Enter to shut down"); + let rt = smarm::init(smarm::Config::default()); + let handle = rt.handle(); + let (sup_tx, sup_rx) = std::sync::mpsc::channel(); + // Graceful shutdown on stdin-Enter: a plain OS thread blocks on - // read_line and fires the handle. No signal handling crate needed. - let (handle, signal) = shutdown_handle(); + // read_line and shuts the ROOT SUPERVISOR down. No signal-handling + // crate needed, and no urus-specific shutdown plumbing — this is + // exactly what a SIGTERM handler would do. std::thread::spawn(move || { + let sup: smarm::Pid = sup_rx.recv().expect("supervisor pid"); let mut line = String::new(); let _ = std::io::stdin().read_line(&mut line); println!("urus-crud: shutting down (draining in-flight requests)…"); - SHUTTING_DOWN.store(true, std::sync::atomic::Ordering::Relaxed); - handle.shutdown(); + handle.request_shutdown(sup); }); - serve_with_shutdown(cfg, pipeline, signal).unwrap(); + rt.run(move || { + let sup = smarm::spawn(move || { + OneForOne::new() + // Store FIRST: reverse-order shutdown therefore stops it + // LAST, after the endpoint has finished draining. + .child(ChildSpec::new(Restart::Permanent, store_loop)) + .child( + ChildSpec::new(Restart::Permanent, endpoint) + // The endpoint bounds its own drain with + // Config::drain_timeout; a shorter supervisor + // deadline would cut that drain in half. + .shutdown(smarm::supervisor::Shutdown::Infinity), + ) + .run() + }); + let _ = sup_tx.send(sup.pid()); + let _ = sup.join(); + }); println!("urus-crud: bye"); } diff --git a/examples/load_profile.rs b/examples/load_profile.rs index 095f2e1..77ee30f 100644 --- a/examples/load_profile.rs +++ b/examples/load_profile.rs @@ -52,7 +52,7 @@ fn main() { let (handle, signal) = shutdown_handle(); let server = std::thread::spawn(move || { let pipe = Pipeline::new().plug(Router::new().get("/json/:id", json_id)); - serve_with_shutdown(Config::new(addr), pipe, signal).expect("serve"); + serve_with_shutdown(Config::new(addr), smarm::Config::default(), pipe, signal).expect("serve"); }); // Readiness is the orchestrator's job (TCP probe); ours is to not diff --git a/examples/plain_serve.rs b/examples/plain_serve.rs index f7c41ef..3e0fea1 100644 --- a/examples/plain_serve.rs +++ b/examples/plain_serve.rs @@ -41,8 +41,13 @@ fn main() { // Audit line: lands in each bench cell's server.log so the effective // scheduler count is recorded per cell, same discipline as mode-verify. eprintln!("plain_serve: scheduler_threads={sched_threads:?}"); - let mut cfg = Config::new(addr); - cfg.scheduler_threads = sched_threads; + // Scheduler threads are a RUNTIME knob, so they live in smarm::Config, + // not urus's — an endpoint placed in someone else's tree could not + // honour them anyway. + let rt_cfg = match sched_threads { + Some(n) => smarm::Config::exact(n), + None => smarm::Config::default(), + }; let pipe = Pipeline::new().plug(Router::new().get("/json/:id", json_id)); - serve_with(cfg, pipe).expect("serve"); + serve_with(Config::new(addr), rt_cfg, pipe).expect("serve"); } diff --git a/examples/serve_toml.rs b/examples/serve_toml.rs index adbac16..9d1845c 100644 --- a/examples/serve_toml.rs +++ b/examples/serve_toml.rs @@ -44,5 +44,5 @@ fn main() { } let pipe = Pipeline::new().plug(Router::new().get("/json/:id", json_id)); - serve_with(cfg, pipe).expect("serve"); + serve_with(cfg, smarm::Config::default(), pipe).expect("serve"); } diff --git a/examples/ws_chat.rs b/examples/ws_chat.rs index 975ffaa..56c8ff1 100644 --- a/examples/ws_chat.rs +++ b/examples/ws_chat.rs @@ -7,14 +7,16 @@ //! //! - **One `PubSub` for the whole app**, topics are rooms //! (`room:{name}`). Built lazily on the first connection via a -//! NON-static `Arc>` captured by the route closure: -//! `PubSub::new()` spawns the table actor, which smarm only allows -//! in-runtime — and keeping the cell non-static means the table's -//! last handle drops in-runtime when the drained pipeline drops, so -//! `serve_with_shutdown` actually returns. A `static` cell would pin -//! the table forever and block smarm's all-done. (Same constraint -//! crud's store solves with its OnceLock; the non-static refinement -//! is what makes graceful shutdown compose.) +//! NON-static `Arc>` captured by the route closure, +//! because `PubSub::new()` spawns the table actor and smarm only +//! allows `spawn` in-runtime — while `serve_with_shutdown` owns the +//! runtime, there is no in-runtime moment before the first request. +//! Keep the cell non-static so the table is dropped with the drained +//! pipeline rather than pinned for the life of the process. +//! +//! If your app owns its own runtime and tree (see `crud`), skip this +//! entirely: start the table as a supervised sibling of +//! `urus::endpoint(...)` and address it by name. //! //! - **`on_open` subscribes and spawns the relay** — a listen-only //! client receives the room without ever sending. The subscription is @@ -123,6 +125,12 @@ fn main() { handle.shutdown(); }); - serve_with_shutdown(Config::new("0.0.0.0:8080".parse().unwrap()), pipeline, signal).unwrap(); + serve_with_shutdown( + Config::new("0.0.0.0:8080".parse().unwrap()), + smarm::Config::default(), + pipeline, + signal, + ) + .unwrap(); println!("ws_chat: drained, bye"); } diff --git a/examples/ws_echo.rs b/examples/ws_echo.rs index 0a96e44..427c888 100644 --- a/examples/ws_echo.rs +++ b/examples/ws_echo.rs @@ -98,6 +98,6 @@ fn main() { handle.shutdown(); }); - serve_with_shutdown(cfg, pipeline, signal).unwrap(); + serve_with_shutdown(cfg, smarm::Config::default(), pipeline, signal).unwrap(); println!("ws_echo: bye"); }