From 8f0da2a80664e634d898a8151b781e2eab51741f Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 20 Aug 2026 14:49:23 +0000 Subject: [PATCH] =?UTF-8?q?feat(pubsub,channels)!:=20handles=20are=20addre?= =?UTF-8?q?sses=20=E2=80=94=20bus,=20hub=20and=20session=20registries=20ar?= =?UTF-8?q?e=20supervised=20children?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit smarm 0.7 (415effb, "lifetime is the actor's — refs are addresses") removed the rule these three actors were built on: a GenServerRef no longer owns the server, the loop holds its own inbox sender, and the inbox never closes when the last ref drops. urus's pubsub table, channel hub bus and session registries were still governed by that deleted rule — PubSub::new() spawned the table and the handle owned its life — so nothing commanded them to stop. What still terminated a run was the root-exit sweep, racing the drain: shutdown_with_open_chat_terminates and channels_wire::shutdown_with_open_channel_terminates failed 4 times in 25 --all-features runs with "serve did not return: ... outlived the drain: Timeout". Zero in 10 full-suite runs after this change. The fix is not a supervisor wrapped around the old shape. Every gotcha in this area descended from constructors that spawn: PubSub::new(), ChannelHub::new() and PrefixRouter::channel_session() all started actors, which forced in-runtime-only construction, which forced the Arc> lazy-init from the first handler, which forced the "cell must not be static" and "a relay must never hold a PubSub clone" rules. Five documented rules propping up one inverted dependency. So: description is separated from instantiation. - PubSub is a name, not a GenServerRef: const-constructible, Copy, spawns nothing, valid outside the runtime and in a static. Operations resolve through the registry per call, so a table restarted by its supervisor is reached transparently (one lookup per broadcast — bench before caching a ref, which would go stale across exactly the restart the supervisor exists to perform). PubSub::new() is gone; PubSub::new(name) + PubSub::child() replace it. - ChannelHub::new(bus, router) returns (hub, Vec) — the bus table plus one registry per session route. Returning both is the point: a hub whose children were never started compiles and fails on the first join, so the vec is not left behind a method you can forget to call. #[must_use]. - channel_session gains a registry name; each session registry is separately named and separately supervised. - serve_with/serve_with_shutdown take a Vec of app children and build the root as RestForOne[..app children, endpoint]. They start before the endpoint and, shutdown being ordered in reverse, stop after it has drained, so a request still in flight can reach the bus. RestForOne because a bus crash leaves live sockets addressing a table that no longer knows them. - Deleted: the Arc idiom from both examples and both test pipelines, and the module rules that existed only to hand-manage a refcount. Known cost, not fixed here: channel_session("session:*", "chat-sessions", f) puts two unrelated string literals side by side and nothing catches a transposition — a RegistryName newtype is the obvious follow-up. Tests: 111 lib + 50 integration + 2 doc green, clippy clean, 10/10 full-suite runs. Unit tests poll for name binding before use — smarm's start-order-is-not- start-readiness gap; real apps don't hit it, since a handler only runs once a connection has been accepted. --- examples/causal_bench.rs | 2 +- examples/channels_chat.rs | 39 ++++----- examples/load_profile.rs | 2 +- examples/plain_serve.rs | 2 +- examples/serve_toml.rs | 2 +- examples/ws_chat.rs | 67 +++++--------- examples/ws_echo.rs | 2 +- src/channels/mod.rs | 154 ++++++++++++++++++++++++++------- src/channels/session.rs | 80 ++++++++++++----- src/channels/testkit.rs | 2 +- src/pubsub.rs | 178 ++++++++++++++++++++++++++------------ src/serve.rs | 48 ++++++---- tests/integration.rs | 139 +++++++++++++++++------------ 13 files changed, 466 insertions(+), 251 deletions(-) diff --git a/examples/causal_bench.rs b/examples/causal_bench.rs index 290d139..b7c0411 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), smarm::Config::default(), pipe, signal).expect("serve"); + serve_with_shutdown(Config::new(addr), smarm::Config::default(), pipe, Vec::new(), signal).expect("serve"); }); // Wait until it's accepting. diff --git a/examples/channels_chat.rs b/examples/channels_chat.rs index 48ef5f0..27244ab 100644 --- a/examples/channels_chat.rs +++ b/examples/channels_chat.rs @@ -10,10 +10,13 @@ //! //! The shapes this demonstrates: //! -//! - **One `ChannelHub` for the whole app**, built lazily on the first -//! connection via the NON-static `Arc>` pattern (the -//! hub spawns the pubsub table; same in-runtime + shutdown laws as -//! `examples/ws_chat.rs`, see the docs there). +//! - **One `ChannelHub` for the whole app**, built up front. The hub is +//! a description — a bus name plus the routing table — and spawns +//! nothing; `hub.children()` is the vec of actors it needs (the bus +//! table, plus one registry per session route), handed to +//! `serve_with_shutdown` so the supervisor starts them ahead of the +//! endpoint and stops them after it has drained. Same shape as +//! `examples/ws_chat.rs`, see the docs there. //! //! - **A channel per joined topic, not per socket.** `Room` never sees //! frames, refs, or the transport heartbeat — the conn-side handler @@ -31,8 +34,6 @@ //! the buffer holds the relay path's own `Arc`s, nothing is copied. //! The session is keyed by the `"name"` in the join payload. -use std::sync::{Arc, OnceLock}; - use serde_json::{json, Value}; use urus::channels::phoenix::Json; use urus::{ @@ -85,23 +86,16 @@ impl ChannelSession

for BySessionName { } fn main() { - let hub: Arc>> = Arc::new(OnceLock::new()); + let (hub, children) = ChannelHub::new( + "chat-bus", + PrefixRouter::new() + .channel_default::("room:*") + .channel_session::("session:*", "chat-sessions", |_: &str| { + Box::new(Room::default()) as Box> + }), + ); - let router = Router::new().get("/socket", move |c: Conn, _n: Next| { - // Hub construction is in-runtime only (it spawns the pubsub - // table and, for session patterns, their registries); the - // non-static cell is what lets shutdown drain them all. - let hub = hub.get_or_init(|| { - ChannelHub::new( - PrefixRouter::new() - .channel_default::("room:*") - .channel_session::("session:*", |_: &str| { - Box::new(Room::default()) as Box> - }), - ) - }); - hub.upgrade(c) - }); + let router = Router::new().get("/socket", move |c: Conn, _n: Next| hub.upgrade(c)); let (handle, signal) = shutdown_handle(); std::thread::spawn(move || { @@ -115,6 +109,7 @@ fn main() { Config::new("0.0.0.0:8080".parse().unwrap()), smarm::Config::default(), Pipeline::new().plug(router), + children, signal, ) .unwrap(); diff --git a/examples/load_profile.rs b/examples/load_profile.rs index 77ee30f..75124d9 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), smarm::Config::default(), pipe, signal).expect("serve"); + serve_with_shutdown(Config::new(addr), smarm::Config::default(), pipe, Vec::new(), 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 3e0fea1..55ab87a 100644 --- a/examples/plain_serve.rs +++ b/examples/plain_serve.rs @@ -49,5 +49,5 @@ fn main() { None => smarm::Config::default(), }; let pipe = Pipeline::new().plug(Router::new().get("/json/:id", json_id)); - serve_with(Config::new(addr), rt_cfg, pipe).expect("serve"); + serve_with(Config::new(addr), rt_cfg, pipe, Vec::new()).expect("serve"); } diff --git a/examples/serve_toml.rs b/examples/serve_toml.rs index 9d1845c..aba09b4 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, smarm::Config::default(), pipe).expect("serve"); + serve_with(cfg, smarm::Config::default(), pipe, Vec::new()).expect("serve"); } diff --git a/examples/ws_chat.rs b/examples/ws_chat.rs index 56c8ff1..d72f6fd 100644 --- a/examples/ws_chat.rs +++ b/examples/ws_chat.rs @@ -1,4 +1,4 @@ -//! WebSocket chat rooms (v0.5): `urus::pubsub` wired to the v0.4 duplex. +//! WebSocket chat rooms (v0.7): `urus::pubsub` wired to the v0.4 duplex. //! //! cargo run --example ws_chat //! websocat ws://127.0.0.1:8080/chat/lobby (run two of these) @@ -6,41 +6,35 @@ //! The shapes this demonstrates: //! //! - **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, -//! 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. +//! (`room:{name}`). `BUS` is a `const`: the handle is an address, not +//! the table. The table actor is `BUS.child()`, handed to +//! `serve_with_shutdown` as an app child, so it starts before the +//! endpoint and — shutdown being ordered in reverse — stops after the +//! endpoint has drained. //! -//! 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. +//! The v0.5 idiom this replaces was a non-static `Arc>` +//! lazily initialised from the first connection, because `PubSub::new()` +//! used to spawn the table and the handle used to own its life. Neither +//! is true any more. //! //! - **`on_open` subscribes and spawns the relay** — a listen-only //! client receives the room without ever sending. The subscription is //! pinned to the CONNECTION actor (`on_open` runs inside it), so the //! monitor cleans up exactly when the connection dies. //! -//! - **The relay holds the `Receiver` and a `WsSender` clone — and -//! deliberately NOT a `PubSub` handle.** A relay holding the handle -//! keeps the table's inbox open while the table keeps the relay's -//! receiver open: neither ever exits, and shutdown hangs. Receiver -//! only: conn dies → monitor prunes → sender drops → relay's recv -//! errs → relay exits. Every link in that chain is in-runtime. - -use std::sync::{Arc, OnceLock}; +//! - **The relay holds the `Receiver` and a `WsSender` clone.** It may +//! also hold `BUS` — a handle pins nothing now — but it has no use for +//! one. Conn dies → monitor prunes → sender drops → relay's recv errs +//! → relay exits. use urus::{ serve_with_shutdown, shutdown_handle, Config, Conn, Message, Next, Pipeline, PubSub, Router, WsHandler, WsSender, }; -type Bus = PubSub; +const BUS: PubSub = PubSub::new("chat"); struct ChatHandler { - bus: Arc>, room: String, } @@ -48,32 +42,23 @@ impl ChatHandler { fn topic(&self) -> String { format!("room:{}", self.room) } - - fn bus(&self) -> &Bus { - // First connection anywhere spawns the table; we are inside the - // connection actor here, so the spawn is legal. - self.bus.get_or_init(PubSub::new) - } } impl WsHandler for ChatHandler { fn on_open(&mut self, sender: &WsSender) { let topic = self.topic(); - let rx = match self.bus().subscribe(&topic) { + let rx = match BUS.subscribe(&topic) { Ok(rx) => rx, Err(_) => { let _ = sender.close(1011, "chat bus down"); return; } }; - let _ = self - .bus() - .broadcast_from(smarm::self_pid(), &topic, format!("* someone joined {topic}")); + let _ = BUS.broadcast_from(smarm::self_pid(), &topic, format!("* someone joined {topic}")); // The relay: room messages -> this socket. Exits when the // subscription is pruned (conn death / unsubscribe) or the - // socket is gone (WsClosed). See the module docs for why it - // must not capture a Bus handle. + // socket is gone (WsClosed). let out = sender.clone(); smarm::spawn(move || { while let Ok(msg) = rx.recv() { @@ -92,16 +77,12 @@ impl WsHandler for ChatHandler { // broadcast_from: the sender's own relay is skipped — no echo. // self_pid() here is the connection actor, the pid on_open // subscribed as. - let _ = self - .bus() - .broadcast_from(smarm::self_pid(), self.topic(), text); + let _ = BUS.broadcast_from(smarm::self_pid(), self.topic(), text); } fn on_close(&mut self, _code: Option, _reason: &str) { let topic = self.topic(); - let _ = self - .bus() - .broadcast_from(smarm::self_pid(), &topic, format!("* someone left {topic}")); + let _ = BUS.broadcast_from(smarm::self_pid(), &topic, format!("* someone left {topic}")); // No explicit unsubscribe: the connection actor is about to // exit and the monitor prunes the subscription (which is also // what stops the relay). @@ -109,12 +90,9 @@ impl WsHandler for ChatHandler { } fn main() { - let bus: Arc> = Arc::new(OnceLock::new()); - - let pipeline = Pipeline::new().plug(Router::new().get("/chat/:room", move |c: Conn, _n: Next| { + let pipeline = Pipeline::new().plug(Router::new().get("/chat/:room", |c: Conn, _n: Next| { let room = c.params.get("room").unwrap_or("lobby").to_string(); - let bus = bus.clone(); - c.upgrade(ChatHandler { bus, room }) + c.upgrade(ChatHandler { room }) })); let (handle, signal) = shutdown_handle(); @@ -129,6 +107,7 @@ fn main() { Config::new("0.0.0.0:8080".parse().unwrap()), smarm::Config::default(), pipeline, + vec![BUS.child()], signal, ) .unwrap(); diff --git a/examples/ws_echo.rs b/examples/ws_echo.rs index 427c888..a4f2788 100644 --- a/examples/ws_echo.rs +++ b/examples/ws_echo.rs @@ -98,6 +98,6 @@ fn main() { handle.shutdown(); }); - serve_with_shutdown(cfg, smarm::Config::default(), pipeline, signal).unwrap(); + serve_with_shutdown(cfg, smarm::Config::default(), pipeline, Vec::new(), signal).unwrap(); println!("ws_echo: bye"); } diff --git a/src/channels/mod.rs b/src/channels/mod.rs index b85e162..3986dbb 100644 --- a/src/channels/mod.rs +++ b/src/channels/mod.rs @@ -63,7 +63,7 @@ use crate::pubsub::PubSub; use crate::ws::{Message, WsHandler, WsSender}; -use smarm::{Pid, Receiver, Sender}; +use smarm::{ChildSpec, Pid, Receiver, Sender}; use std::cell::RefCell; use std::collections::HashMap; @@ -199,6 +199,16 @@ pub trait Channel: Send + 'static { pub trait ChannelFactory: Send + Sync { fn create(&self, topic: &str) -> Box>; + /// Supervised actors this factory needs running before it can + /// deploy. Empty for the default (ephemeral) factory, which spawns + /// its channel actor per join under the connection; the session + /// factory returns its registry's spec. Internal seam, collected by + /// [`ChannelHub::children`]. + #[doc(hidden)] + fn children(&self) -> Vec { + Vec::new() + } + /// How an accepted-routing join becomes a running channel actor. /// Internal seam — the default (ephemeral actor, linked to the /// connection, cold start per join) is the contract; only the @@ -258,6 +268,13 @@ impl + Default> ChannelFactory

for De /// rejects the join (status error). pub trait TopicRouter: Send + Sync + 'static { fn route(&self, topic: &str) -> Option>>; + + /// Every supervised actor the routes need, gathered for + /// [`ChannelHub::children`]. Override only if your router holds + /// factories it does not surface through [`route`](Self::route). + fn children(&self) -> Vec { + Vec::new() + } } /// The shipped router: exact topics and `head:*` prefix patterns. @@ -301,31 +318,46 @@ impl PrefixRouter

{ /// [`channel`](Self::channel), with opt-in session persistence /// keyed and configured by `S`'s [`ChannelSession`] impl (usually - /// `S` is the channel type itself). **In-runtime only**: this - /// spawns the pattern's session-registry actor, same law as - /// [`ChannelHub::new`]. - pub fn channel_session(self, pattern: &str, factory: impl ChannelFactory

+ 'static) -> Self + /// `S` is the channel type itself). + /// + /// `registry` names the pattern's session-registry actor, which is + /// supervised: it appears in [`ChannelHub::children`] and must be + /// unique across the app. Spawns nothing here. + pub fn channel_session( + self, + pattern: &str, + registry: &'static str, + factory: impl ChannelFactory

+ 'static, + ) -> Self where S: ChannelSession

, S::Key: Clone, P: Encode + Decode, { - self.channel(pattern, session::SessionFactory::new::(Arc::new(factory))) + self.channel(pattern, session::SessionFactory::new::(registry, Arc::new(factory))) } /// [`channel_session`](Self::channel_session) with a /// [`Default`]-built impl that is its own session config. - pub fn channel_session_default(self, pattern: &str) -> Self + pub fn channel_session_default(self, pattern: &str, registry: &'static str) -> Self where C: Channel

+ ChannelSession

+ Default, C::Key: Clone, P: Encode + Decode, { - self.channel_session::(pattern, DefaultFactory::(PhantomData)) + self.channel_session::(pattern, registry, DefaultFactory::(PhantomData)) } } impl TopicRouter

for PrefixRouter

{ + fn children(&self) -> Vec { + self.exact + .values() + .chain(self.prefix.values()) + .flat_map(|f| f.children()) + .collect() + } + fn route(&self, topic: &str) -> Option>> { if let Some(f) = self.exact.get(topic) { return Some(f.clone()); @@ -427,13 +459,20 @@ impl ChannelSocket

{ // ChannelHub — the app-facing entry point // --------------------------------------------------------------------------- -/// One per app (or per channel namespace): the router plus the pubsub -/// bus every channel broadcasts on. +/// One per app (or per channel namespace): the router plus the address +/// of the pubsub bus every channel broadcasts on. /// -/// **In-runtime only** (spawns the pubsub table) and **must not live in -/// a `static`** — use the non-static `Arc>>` -/// captured by the route closure, exactly the v0.5 pubsub pattern, or -/// graceful shutdown will hang waiting for the table to exit. +/// Spawns nothing — build it wherever you like, including out of the +/// runtime and alongside the pipeline. The actors it needs (the bus +/// table, plus one registry per session route) come from +/// [`children`](Self::children); hand that vec to `serve_with*` or splice +/// it into your own supervision tree ahead of the endpoint. +/// +/// ```ignore +/// let (hub, children) = +/// ChannelHub::new("chat-bus", PrefixRouter::new().channel_default::("room:*")); +/// serve_with(cfg, rt_cfg, pipeline, children)?; +/// ``` pub struct ChannelHub { bus: PubSub>, router: Arc>, @@ -441,13 +480,27 @@ pub struct ChannelHub { impl Clone for ChannelHub

{ fn clone(&self) -> Self { - Self { bus: self.bus.clone(), router: self.router.clone() } + Self { bus: self.bus, router: self.router.clone() } } } impl ChannelHub

{ - pub fn new(router: impl TopicRouter

) -> Self { - Self { bus: PubSub::new(), router: Arc::new(router) } + /// Address a hub whose bus table is registered under `bus`, together + /// with every actor it needs running: the bus table first, then one + /// registry per session route. Spawns nothing. + /// + /// The two come back together on purpose. A hub whose children were + /// never started compiles fine and fails on the first join, so the + /// constructor hands you the vec rather than leaving it behind a + /// method you can forget to call. Start them ahead of the endpoint — + /// `RestForOne` in that order means a bus crash also restarts the + /// endpoint, dropping connections whose subscriptions died with it. + #[must_use] + pub fn new(bus: &'static str, router: impl TopicRouter

) -> (Self, Vec) { + let hub = Self { bus: PubSub::new(bus), router: Arc::new(router) }; + let mut children = vec![hub.bus.child()]; + children.extend(hub.router.children()); + (hub, children) } /// Accept the WebSocket upgrade on `conn` and speak channels over @@ -455,7 +508,7 @@ impl ChannelHub

{ /// [`Conn::upgrade`](crate::Conn::upgrade) apply unchanged. pub fn upgrade(&self, conn: crate::Conn) -> crate::Conn { conn.upgrade(SocketHandler { - bus: self.bus.clone(), + bus: self.bus, router: self.router.clone(), joined: HashMap::new(), }) @@ -567,7 +620,7 @@ impl WsHandler for SocketHandler

reference: frame.reference, payload, ws: sender.clone(), - bus: self.bus.clone(), + bus: self.bus, }); self.joined .insert(frame.topic, Joined { join_ref: frame.join_ref, tx: inbox.0 }); @@ -649,7 +702,7 @@ fn run_channel( let socket = ChannelSocket { ws, - bus: bus.clone(), + bus, topic: topic.clone(), join_ref, cur_ref: RefCell::new(None), @@ -799,12 +852,39 @@ mod tests { } - fn hub(terminated: &Arc) -> ChannelHub { + const TEST_BUS: &str = "chan-unit-bus"; + + fn hub(terminated: &Arc) -> (ChannelHub, Vec) { ChannelHub::new( + TEST_BUS, PrefixRouter::new().channel("room:*", RoomFactory { terminated: terminated.clone() }), ) } + /// Start a hub's children under a supervisor and wait for the bus to + /// bind its name (smarm's start-order-is-not-start-readiness gap: + /// `start_child` spawns and moves on). Returns the sup's pid. + pub(crate) fn start_hub( + hub: &ChannelHub

, + children: Vec, + ) -> Pid { + let sup = smarm::spawn(move || { + let mut sup = smarm::OneForOne::new(); + for c in children { + sup = sup.child(c); + } + sup.run() + }); + let bus = hub.bus; + for _ in 0..200 { + if bus.pid().is_some() { + return sup.pid(); + } + smarm::sleep(std::time::Duration::from_millis(10)); + } + panic!("hub children never came up"); + } + #[test] fn prefix_router_exact_and_prefix_and_miss() { @@ -825,7 +905,8 @@ mod tests { let out = Arc::new(Mutex::new(Vec::::new())); let (out2, t) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = hub(&t); + let (hub, children) = hub(&t); + let _sup = start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("j1|r1|room:a|event|phx_join|hi"); out2.lock().unwrap().push(h.recv()); @@ -843,7 +924,8 @@ mod tests { let (out2, t) = (out.clone(), Arc::new(AtomicBool::new(false))); let t2 = t.clone(); smarm::run(move || { - let hub = hub(&t2); + let (hub, children) = hub(&t2); + let _sup = start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("j1|r1|room:locked|event|phx_join|hi"); *out2.lock().unwrap() = h.recv(); @@ -859,7 +941,8 @@ mod tests { let out = Arc::new(Mutex::new(String::new())); let (out2, t) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = hub(&t); + let (hub, children) = hub(&t); + let _sup = start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("j1|r1|hall:a|event|phx_join|hi"); *out2.lock().unwrap() = h.recv(); @@ -872,7 +955,8 @@ mod tests { let out = Arc::new(Mutex::new(String::new())); let (out2, t) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = hub(&t); + let (hub, children) = hub(&t); + let _sup = start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("-|hb1|phoenix|event|heartbeat|"); *out2.lock().unwrap() = h.recv(); @@ -885,7 +969,8 @@ mod tests { let out = Arc::new(Mutex::new(String::new())); let (out2, t) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = hub(&t); + let (hub, children) = hub(&t); + let _sup = start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("j1|r1|room:a|event|ping|x"); *out2.lock().unwrap() = h.recv(); @@ -898,7 +983,8 @@ mod tests { let got = Arc::new(Mutex::new(Vec::<(u8, String)>::new())); let (got2, t) = (got.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = hub(&t); + let (hub, children) = hub(&t); + let _sup = start_hub(&hub, children); let mut a = Harness::new(&hub); let mut b = Harness::new(&hub); a.send("j1|r1|room:a|event|phx_join|A"); @@ -930,7 +1016,8 @@ mod tests { let (out2, t) = (out.clone(), Arc::new(AtomicBool::new(false))); let t2 = t.clone(); smarm::run(move || { - let hub = hub(&t2); + let (hub, children) = hub(&t2); + let _sup = start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("j1|r1|room:a|event|phx_join|hi"); out2.lock().unwrap().push(h.recv()); @@ -954,8 +1041,9 @@ mod tests { let t = Arc::new(AtomicBool::new(false)); let t2 = t.clone(); smarm::run(move || { - let hub = hub(&t2); - let bus = hub.bus.clone(); + let (hub, children) = hub(&t2); + let _sup = start_hub(&hub, children); + let bus = hub.bus; { let mut h = Harness::new(&hub); h.send("j1|r1|room:a|event|phx_join|hi"); @@ -974,7 +1062,8 @@ mod tests { let (out2, t) = (out.clone(), Arc::new(AtomicBool::new(false))); let t2 = t.clone(); smarm::run(move || { - let hub = hub(&t2); + let (hub, children) = hub(&t2); + let _sup = start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("j1|r1|room:a|event|phx_join|one"); out2.lock().unwrap().push(h.recv()); @@ -998,7 +1087,8 @@ mod tests { let out = Arc::new(Mutex::new(String::new())); let (out2, t) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = hub(&t); + let (hub, children) = hub(&t); + let _sup = start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("j1|r1|room:a|event|phx_join|hi"); h.recv(); diff --git a/src/channels/session.rs b/src/channels/session.rs index 9dce353..84f1b21 100644 --- a/src/channels/session.rs +++ b/src/channels/session.rs @@ -48,8 +48,8 @@ use super::*; -use smarm::gen_server::{self, GenServer, GenServerCtx}; -use smarm::{Down, GenServerRef, Watcher}; +use smarm::gen_server::{self, GenServer, GenServerBuilder, GenServerCtx, GenServerName}; +use smarm::{Down, Restart, Watcher}; use std::collections::VecDeque; use std::hash::Hash; @@ -87,7 +87,14 @@ pub trait ChannelSession

: Send + 'static { pub(super) struct SessionFactory { inner: Arc>, keyfn: fn(&str, &P) -> K, - registry: GenServerRef>, + /// The registry's registered name. Resolved per deploy, so a + /// registry restarted by the supervisor is reached transparently — + /// with its session map empty, which is the honest outcome: the + /// session actors it tracked died with their control senders. + registry: GenServerName>, + /// Everything the registry's `ChildSpec` needs to build it again. + cap: usize, + ttl: Duration, } /// The registry's working bounds for a session key. @@ -95,17 +102,19 @@ pub(super) trait SessionKey: Eq + Hash + Clone + Send + 'static {} impl SessionKey for K {} impl SessionFactory { - /// **In-runtime only** — spawns the registry gen_server. - pub(super) fn new>(inner: Arc>) -> Self { - let registry = gen_server::start(Registry { - factory: inner.clone(), + /// Describe the session route. Spawns nothing; the registry actor is + /// started from [`children`](ChannelFactory::children). + pub(super) fn new>( + registry: &'static str, + inner: Arc>, + ) -> Self { + SessionFactory { + inner, + keyfn: S::session_key, + registry: GenServerName::new(registry), cap: S::buffer_cap(), ttl: S::ttl(), - sessions: HashMap::new(), - pid_key: HashMap::new(), - watcher: None, - }); - SessionFactory { inner, keyfn: S::session_key, registry } + } } } @@ -116,6 +125,23 @@ impl ChannelFactory

Vec { + let (name, factory, cap, ttl) = (self.registry, self.inner.clone(), self.cap, self.ttl); + vec![ChildSpec::new(Restart::Permanent, move || { + let state = Registry:: { + factory: factory.clone(), + cap, + ttl, + sessions: HashMap::new(), + pid_key: HashMap::new(), + watcher: None, + }; + if let Err(e) = GenServerBuilder::new(state).named(name).run() { + panic!("urus session registry '{}': name already taken: {e:?}", name.as_str()); + } + })] + } + fn deploy(&self, hs: JoinHandshake

) -> ChannelInbox

where P: Encode + Decode, @@ -125,7 +151,7 @@ impl ChannelFactory

>( term: &Arc, reject_rejoin: bool, - ) -> ChannelHub { + ) -> (ChannelHub, Vec) { ChannelHub::new( - PrefixRouter::new() - .channel_session::("room:*", counter_factory(term, reject_rejoin)), + "sess-unit-bus", + PrefixRouter::new().channel_session::( + "room:*", + "sess-unit-registry", + counter_factory(term, reject_rejoin), + ), ) } @@ -580,7 +610,8 @@ mod tests { let out = Arc::new(Mutex::new(Vec::::new())); let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = session_hub::(&term, false); + let (hub, children) = session_hub::(&term, false); + let _sup = crate::channels::tests::start_hub(&hub, children); let mut h1 = Harness::new(&hub); h1.send("j1|r1|room:a|event|phx_join|u"); out2.lock().unwrap().push(h1.recv()); @@ -611,7 +642,8 @@ mod tests { let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false))); let term2 = term.clone(); smarm::run(move || { - let hub = session_hub::(&term2, false); + let (hub, children) = session_hub::(&term2, false); + let _sup = crate::channels::tests::start_hub(&hub, children); let mut h1 = Harness::new(&hub); h1.send("j1|r1|room:a|event|phx_join|u"); out2.lock().unwrap().push(h1.recv()); @@ -632,7 +664,8 @@ mod tests { let out = Arc::new(Mutex::new(Vec::::new())); let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = session_hub::(&term, false); + let (hub, children) = session_hub::(&term, false); + let _sup = crate::channels::tests::start_hub(&hub, children); let mut h1 = Harness::new(&hub); h1.send("j1|r1|room:a|event|phx_join|u"); out2.lock().unwrap().push(h1.recv()); @@ -657,7 +690,8 @@ mod tests { let out = Arc::new(Mutex::new(Vec::::new())); let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = session_hub::(&term, false); + let (hub, children) = session_hub::(&term, false); + let _sup = crate::channels::tests::start_hub(&hub, children); let mut h = Harness::new(&hub); h.send("j1|r1|room:a|event|phx_join|u"); out2.lock().unwrap().push(h.recv()); @@ -681,7 +715,8 @@ mod tests { let out = Arc::new(Mutex::new(Vec::::new())); let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = session_hub::(&term, false); + let (hub, children) = session_hub::(&term, false); + let _sup = crate::channels::tests::start_hub(&hub, children); let mut h1 = Harness::new(&hub); h1.send("j1|r1|room:a|event|phx_join|u"); out2.lock().unwrap().push(h1.recv()); @@ -714,7 +749,8 @@ mod tests { let out = Arc::new(Mutex::new(Vec::::new())); let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false))); smarm::run(move || { - let hub = session_hub::(&term, true); + let (hub, children) = session_hub::(&term, true); + let _sup = crate::channels::tests::start_hub(&hub, children); let mut h1 = Harness::new(&hub); h1.send("j1|r1|room:a|event|phx_join|u"); out2.lock().unwrap().push(h1.recv()); diff --git a/src/channels/testkit.rs b/src/channels/testkit.rs index 7cc2dc5..9966091 100644 --- a/src/channels/testkit.rs +++ b/src/channels/testkit.rs @@ -75,7 +75,7 @@ impl Harness { let (tx, out_rx) = smarm::channel(); Self { handler: SocketHandler { - bus: hub.bus.clone(), + bus: hub.bus, router: hub.router.clone(), joined: HashMap::new(), }, diff --git a/src/pubsub.rs b/src/pubsub.rs index d8b6664..41f1601 100644 --- a/src/pubsub.rs +++ b/src/pubsub.rs @@ -38,37 +38,46 @@ //! but stays alive keeps its (one-shot, inert) monitor until death; //! that's a bounded bookkeeping entry, not a leak. //! -//! # Construction must happen in-runtime +//! # The handle is an address, not an owner (v0.7) //! -//! [`PubSub::new`] spawns the table actor, which smarm only permits from -//! inside `Runtime::run`. Pipelines are built *before* `serve*` boots the -//! runtime, so the working pattern (same constraint crud's store hits) is -//! lazy init from the first handler — but with a **non-static** -//! `Arc>>` captured by the route closure, NOT a -//! `static`: when the drained pipeline drops at graceful shutdown, the -//! cell (and thus the last handle) drops in-runtime, the table's inbox -//! closes, and the table exits — `serve_with_shutdown` returns. A -//! `static` pins the table forever and blocks smarm's all-done. The same -//! reasoning forbids long-lived consumer actors (relays, producers) from -//! holding a `PubSub` clone: relay holds table's inbox open, table holds -//! relay's receiver open, neither exits. Relays take the `Receiver` only. -//! See `examples/ws_chat.rs` and the `shutdown_with_open_chat_terminates` -//! integration test for the full chain. +//! [`PubSub::new`] takes a **name** and spawns nothing. It is `const`, +//! costs a `&'static str`, and may be built anywhere — out of the +//! runtime, in a `static`, in a route closure, held by a long-lived +//! relay. Every operation resolves the name through smarm's registry, so +//! a table restarted by its supervisor is reached transparently. //! -//! # Why there is no `register(name)` helper (yet) +//! The table actor is started by [`PubSub::child`], a `ChildSpec` you put +//! in your supervision tree (or hand to `serve_with*`, which puts it in +//! the root ahead of the endpoint). Its lifetime is the supervisor's: +//! it stops when the supervisor shuts it down, in reverse start order, +//! after the endpoint has drained. //! -//! smarm's registry maps `name → Pid`, but a `Pid` cannot be turned back -//! into a `GenServerRef` (the ref *is* the inbox sender). A useful named -//! lookup therefore needs either smarm support (registry-held senders) -//! or a process-global type-erased map here — both against the grain of -//! the ratified design. Deferred; pass the handle. +//! ```ignore +//! const BUS: PubSub = PubSub::new("events"); +//! serve_with(cfg, rt_cfg, pipeline, vec![BUS.child()])?; +//! ``` +//! +//! This replaces the v0.5–v0.6 idiom where `PubSub::new()` spawned the +//! table and the handle owned its life — a non-static +//! `Arc>>` lazily initialised from the first handler, +//! with matching rules that the cell must not be `static` and that +//! relays must never hold a clone. All of that existed to hand-manage a +//! refcount. smarm 0.7 made a server's lifetime its own (refs are +//! addresses; the inbox no longer closes when the last ref drops), which +//! both removed the mechanism those rules relied on and made the named +//! lookup below possible. +//! +//! Cost: one registry resolution per operation, including per broadcast. +//! Caching a `GenServerRef` in the handle would save it and go stale +//! across exactly the restart the supervisor exists to perform. Measure +//! before optimising — see the bench item in `ROADMAP.md`. use std::collections::{HashMap, HashSet}; use std::fmt; use std::sync::Arc; -use smarm::gen_server::{self, GenServer, GenServerCtx}; -use smarm::{channel, Down, Pid, Receiver, Sender, GenServerRef, Watcher}; +use smarm::gen_server::{self, GenServer, GenServerBuilder, GenServerCtx, GenServerName}; +use smarm::{channel, ChildSpec, Down, GenServerRef, Pid, Receiver, Restart, Sender, Watcher}; // --------------------------------------------------------------------------- // Public handle @@ -87,32 +96,58 @@ impl fmt::Display for PubSubDown { impl std::error::Error for PubSubDown {} -/// A clonable handle to one pub/sub instance (one topic table actor). +/// The address of one pub/sub instance: a name, resolved per operation. +/// +/// Cheap (`Copy`, a `&'static str`), constructible anywhere including +/// `const` context, and owns nothing — the table actor it addresses is +/// started by [`child`](Self::child) under a supervisor. Every operation +/// returns [`PubSubDown`] if no live table currently holds the name. /// /// Payloads are broadcast as `Arc`: one allocation per broadcast, not /// per subscriber. pub struct PubSub { - server: GenServerRef>, + name: GenServerName>, } impl Clone for PubSub { fn clone(&self) -> Self { - PubSub { server: self.server.clone() } + *self } } -impl Default for PubSub { - fn default() -> Self { - Self::new() - } -} +impl Copy for PubSub {} impl PubSub { - /// Start a fresh topic table. Must run inside the smarm runtime (it - /// spawns the table's gen_server actor). The table lives until the - /// last handle is dropped. - pub fn new() -> Self { - PubSub { server: gen_server::start(Table::new()) } + /// Address the topic table registered under `name`. Spawns nothing + /// and never fails: the name is resolved at each use. + pub const fn new(name: &'static str) -> Self { + PubSub { name: GenServerName::new(name) } + } + + /// The `ChildSpec` that runs this instance's table actor. Put it in + /// your supervision tree ahead of anything that broadcasts — + /// `serve_with*` takes a `Vec` for exactly this. + /// + /// `Permanent`: a table that dies is a bug, and its subscribers' + /// receivers died with it, so the restart is only half a repair — + /// pair it with a `RestForOne` parent (as `serve_with*` does) so the + /// endpoint restarts behind it and connections re-subscribe. + pub fn child(&self) -> ChildSpec { + let name = self.name; + ChildSpec::new(Restart::Permanent, move || { + if let Err(e) = GenServerBuilder::new(Table::::new()).named(name).run() { + panic!("urus pubsub '{}': name already taken: {e:?}", name.as_str()); + } + }) + } + + /// The registry key this handle resolves. + pub const fn name(&self) -> &'static str { + self.name.as_str() + } + + fn server(&self) -> Result>, PubSubDown> { + gen_server::whereis_server(self.name).ok_or(PubSubDown) } /// Subscribe the **calling actor** to `topic`. Returns the receiving @@ -139,7 +174,7 @@ impl PubSub { let (tx, rx) = channel(); // A call, not a cast: when this returns the table is updated, so // a broadcast issued right after by the same caller is seen. - match self.server.call(Call::Subscribe { topic: topic.into(), pid, tx }) { + match self.server()?.call(Call::Subscribe { topic: topic.into(), pid, tx }) { Ok(_) => Ok(rx), Err(_) => Err(PubSubDown), } @@ -153,7 +188,7 @@ impl PubSub { /// [`unsubscribe`](Self::unsubscribe) for an explicit pid. pub fn unsubscribe_as(&self, pid: Pid, topic: impl Into) -> Result<(), PubSubDown> { - self.server + self.server()? .cast(Cast::Unsubscribe { topic: topic.into(), pid }) .map_err(|_| PubSubDown) } @@ -181,20 +216,21 @@ impl PubSub { /// table — a subscriber whose receiver was dropped but hasn't been /// pruned yet (no broadcast since, still alive) is still counted. pub fn subscriber_count(&self, topic: impl Into) -> Result { - match self.server.call(Call::Count { topic: topic.into() }) { + match self.server()?.call(Call::Count { topic: topic.into() }) { Ok(Reply::Count(n)) => Ok(n), Ok(Reply::Subscribed) => unreachable!("Count call answered with Subscribed"), Err(_) => Err(PubSubDown), } } - /// The topic table actor's pid — for introspection / registry use. - pub fn pid(&self) -> Pid { - self.server.pid() + /// The topic table actor's pid, or `None` if no live table holds the + /// name (not started yet, or between a crash and its restart). + pub fn pid(&self) -> Option { + self.server().ok().map(|s| s.pid()) } fn cast_broadcast(&self, topic: String, msg: M, skip: Option) -> Result<(), PubSubDown> { - self.server + self.server()? .cast(Cast::Broadcast { topic, msg: Arc::new(msg), skip }) .map_err(|_| PubSubDown) } @@ -328,15 +364,33 @@ mod tests { false } + /// Start `ps`'s table under a supervisor and hand back the sup's pid. + /// + /// The readiness poll is smarm's "start order is not start readiness" + /// gap (see smarm ROADMAP): `start_child` spawns and moves on, so the + /// name may not be bound when this returns. Real apps don't hit it — + /// a handler only runs once a connection has been accepted, long + /// after the tree is up — but a test that broadcasts immediately does. + fn start_table(ps: PubSub) -> Pid { + let sup = smarm::spawn(move || smarm::OneForOne::new().child(ps.child()).run()); + assert!(eventually(|| ps.pid().is_some()), "table never bound its name"); + sup.pid() + } + + /// Ordered shutdown of the tree `start_table` built. + fn stop_table(sup: Pid) { + smarm::request_shutdown(sup); + } + #[test] fn broadcast_reaches_all_subscribers_once() { let got = Arc::new(Mutex::new(Vec::<(u32, String)>::new())); let got2 = got.clone(); smarm::run(move || { - let ps = PubSub::::new(); + let ps = PubSub::::new("test-bus"); + let sup = start_table(ps); let mut handles = Vec::new(); for i in 0..2u32 { - let ps = ps.clone(); let got = got2.clone(); handles.push(smarm::spawn(move || { let rx = ps.subscribe("room:a").unwrap(); @@ -351,6 +405,7 @@ mod tests { for h in handles { h.join().unwrap(); } + stop_table(sup); }); let mut v = got.lock().unwrap().clone(); v.sort(); @@ -362,10 +417,10 @@ mod tests { let ptrs = Arc::new(Mutex::new(Vec::::new())); let ptrs2 = ptrs.clone(); smarm::run(move || { - let ps = PubSub::>::new(); + let ps = PubSub::>::new("test-bus"); + let sup = start_table(ps); let mut handles = Vec::new(); for _ in 0..2 { - let ps = ps.clone(); let ptrs = ptrs2.clone(); handles.push(smarm::spawn(move || { let rx = ps.subscribe("t").unwrap(); @@ -378,6 +433,7 @@ mod tests { for h in handles { h.join().unwrap(); } + stop_table(sup); }); let v = ptrs.lock().unwrap(); assert_eq!(v.len(), 2); @@ -389,8 +445,9 @@ mod tests { let got = Arc::new(Mutex::new(Vec::::new())); let got2 = got.clone(); smarm::run(move || { - let ps = PubSub::::new(); - let ps_loud = ps.clone(); + let ps = PubSub::::new("test-bus"); + let sup = start_table(ps); + let ps_loud = ps; let got = got2.clone(); let loud = smarm::spawn(move || { let rx = ps_loud.subscribe("room").unwrap(); @@ -410,6 +467,7 @@ mod tests { .unwrap() .push(format!("root got {} then {}", *first, *second)); loud.join().unwrap(); + stop_table(sup); }); let v = got.lock().unwrap(); assert!(v.contains(&"loud got for everyone".to_string()), "{v:?}"); @@ -426,7 +484,8 @@ mod tests { let ok = Arc::new(Mutex::new(false)); let ok2 = ok.clone(); smarm::run(move || { - let ps = PubSub::::new(); + let ps = PubSub::::new("test-bus"); + let sup = start_table(ps); let rx = ps.subscribe("t").unwrap(); ps.unsubscribe("t").unwrap(); // Unsubscribe is a cast; the table holds the only sender, so @@ -435,6 +494,7 @@ mod tests { assert!(rx.recv().is_err()); ps.broadcast("t", 7).unwrap(); // no subscribers: no-op, no panic *ok2.lock().unwrap() = true; + stop_table(sup); }); assert!(*ok.lock().unwrap()); } @@ -444,7 +504,8 @@ mod tests { let got = Arc::new(Mutex::new((0u32, false))); let got2 = got.clone(); smarm::run(move || { - let ps = PubSub::::new(); + let ps = PubSub::::new("test-bus"); + let sup = start_table(ps); let rx_old = ps.subscribe("t").unwrap(); let rx_new = ps.subscribe("t").unwrap(); assert_eq!(ps.subscriber_count("t").unwrap(), 1, "idempotent per (pid, topic)"); @@ -452,6 +513,7 @@ mod tests { let v = *rx_new.recv().unwrap(); let old_closed = rx_old.recv().is_err(); *got2.lock().unwrap() = (v, old_closed); + stop_table(sup); }); assert_eq!(*got.lock().unwrap(), (42, true)); } @@ -461,8 +523,9 @@ mod tests { let ok = Arc::new(Mutex::new(false)); let ok2 = ok.clone(); smarm::run(move || { - let ps = PubSub::::new(); - let ps2 = ps.clone(); + let ps = PubSub::::new("test-bus"); + let sup = start_table(ps); + let ps2 = ps; let h = smarm::spawn(move || { let _rx = ps2.subscribe("t").unwrap(); // Exit without ever receiving: rx drops with the stack. @@ -471,6 +534,7 @@ mod tests { // No broadcast issued — this MUST be the monitor path, not // prune-on-send-failure. *ok2.lock().unwrap() = eventually(|| ps.subscriber_count("t").unwrap() == 0); + stop_table(sup); }); assert!(*ok.lock().unwrap(), "monitor Down never pruned the dead subscriber"); } @@ -480,7 +544,8 @@ mod tests { let counts = Arc::new(Mutex::new((0usize, 0usize))); let counts2 = counts.clone(); smarm::run(move || { - let ps = PubSub::::new(); + let ps = PubSub::::new("test-bus"); + let sup = start_table(ps); let rx = ps.subscribe("t").unwrap(); drop(rx); let before = ps.subscriber_count("t").unwrap(); @@ -493,6 +558,7 @@ mod tests { after = if eventually(|| ps.subscriber_count("t").unwrap() == 0) { 0 } else { after }; } *counts2.lock().unwrap() = (before, after); + stop_table(sup); }); assert_eq!(*counts.lock().unwrap(), (1, 0)); } @@ -502,11 +568,13 @@ mod tests { let got = Arc::new(Mutex::new(Vec::::new())); let got2 = got.clone(); smarm::run(move || { - let ps = PubSub::::new(); + let ps = PubSub::::new("test-bus"); + let sup = start_table(ps); let rx_a = ps.subscribe("a").unwrap(); ps.broadcast("b", 99).unwrap(); // nobody on b; must not reach a ps.broadcast("a", 1).unwrap(); got2.lock().unwrap().push(*rx_a.recv().unwrap()); + stop_table(sup); }); assert_eq!(*got.lock().unwrap(), vec![1]); } diff --git a/src/serve.rs b/src/serve.rs index 1fc0417..24ccb1e 100644 --- a/src/serve.rs +++ b/src/serve.rs @@ -10,7 +10,7 @@ use crate::conn_actor::ConnLimits; use crate::plug::Pipeline; use smarm::supervisor::Shutdown; -use smarm::{ChildSpec, OneForOne, Restart}; +use smarm::{ChildSpec, OneForOne, Restart, Strategy}; use std::io::{self, ErrorKind}; use std::net::{SocketAddr, ToSocketAddrs}; @@ -238,16 +238,28 @@ pub fn shutdown_handle() -> (Handle, ShutdownSignal) { /// 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. +/// The tree is `root sup -> [..app_children, endpoint]` on +/// [`Strategy::RestForOne`], with the endpoint on [`Shutdown::Infinity`] +/// so its `drain_timeout` — not a supervisor deadline — bounds the drain. +/// +/// `app_children` is your application's actors: a +/// [`PubSub::child`](crate::PubSub::child), a +/// [`ChannelHub::children`](crate::channels::ChannelHub::children), your +/// own state servers. They start **before** the endpoint and, because +/// shutdown is ordered in reverse, stop **after** it has drained — so a +/// request still in flight can still reach the bus. `RestForOne` in that +/// order also means one of them crashing restarts the endpoint behind it, +/// dropping connections whose subscriptions died with it, rather than +/// leaving live sockets addressing a table that no longer knows them. +/// +/// The root actor parks on the signal channel; a `Handle::shutdown` from a +/// foreign OS thread wakes it, it shuts the supervisor down and `rt.run` +/// returns when the last actor is gone. pub fn serve_with_shutdown( config: Config, rt_config: smarm::Config, pipeline: Pipeline, + app_children: Vec, signal: ShutdownSignal, ) -> io::Result<()> { let addr = config.addr; @@ -257,10 +269,11 @@ pub fn serve_with_shutdown( let rt = smarm::init(rt_config); rt.run(move || { let sup = smarm::spawn(move || { - OneForOne::new() - .child( - ChildSpec::new(Restart::Permanent, endpoint).shutdown(Shutdown::Infinity), - ) + let mut sup = OneForOne::new().strategy(Strategy::RestForOne); + for child in app_children { + sup = sup.child(child); + } + sup.child(ChildSpec::new(Restart::Permanent, endpoint).shutdown(Shutdown::Infinity)) .run() }); @@ -285,20 +298,25 @@ pub fn serve_with_shutdown( /// [`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<()> { +pub fn serve_with( + config: Config, + rt_config: smarm::Config, + pipeline: Pipeline, + app_children: Vec, +) -> io::Result<()> { // The Handle is dropped immediately: shutdown can never be signalled. let (_handle, signal) = shutdown_handle(); - serve_with_shutdown(config, rt_config, pipeline, signal) + serve_with_shutdown(config, rt_config, pipeline, app_children, signal) } /// Defaults all round: default [`Config`], default smarm runtime (one -/// scheduler thread per CPU), serve until killed. +/// scheduler thread per CPU), no app children, 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), smarm::Config::default(), pipeline) + serve_with(Config::new(addr), smarm::Config::default(), pipeline, Vec::new()) } diff --git a/tests/integration.rs b/tests/integration.rs index a196cd5..64dd639 100644 --- a/tests/integration.rs +++ b/tests/integration.rs @@ -26,6 +26,11 @@ fn free_port() -> u16 { } fn spawn_server(pipeline: Pipeline) -> u16 { + spawn_server_children(pipeline, Vec::new()) +} + +/// [`spawn_server`] with app children started ahead of the endpoint. +fn spawn_server_children(pipeline: Pipeline, app_children: Vec) -> u16 { let port = free_port(); let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); std::thread::spawn(move || { @@ -33,7 +38,7 @@ fn spawn_server(pipeline: Pipeline) -> u16 { listener_pool: 2, ..Config::new(addr) }; - serve_with(cfg, smarm::Config::exact(2), pipeline).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipeline, app_children).unwrap(); }); // Wait for the server to actually be listening. for _ in 0..50 { @@ -253,7 +258,7 @@ fn panicking_listener_restarts() { listener_pool: 1, ..Config::new(addr) }; - serve_with(cfg, smarm::Config::exact(2), pipe).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipe, Vec::new()).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { @@ -296,6 +301,16 @@ fn panicking_listener_restarts() { fn spawn_server_with_handle( pipeline: Pipeline, drain: Duration, +) -> (u16, urus::Handle, std::sync::mpsc::Receiver<()>) { + spawn_server_with_children(pipeline, drain, Vec::new()) +} + +/// [`spawn_server_with_handle`] with app children started ahead of the +/// endpoint — the pubsub table, a channel hub's actors, and so on. +fn spawn_server_with_children( + pipeline: Pipeline, + drain: Duration, + app_children: Vec, ) -> (u16, urus::Handle, std::sync::mpsc::Receiver<()>) { let port = free_port(); let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); @@ -307,7 +322,8 @@ fn spawn_server_with_handle( drain_timeout: drain, ..Config::new(addr) }; - urus::serve_with_shutdown(cfg, smarm::Config::exact(2), pipeline, signal).unwrap(); + urus::serve_with_shutdown(cfg, smarm::Config::exact(2), pipeline, app_children, signal) + .unwrap(); let _ = done_tx.send(()); }); for _ in 0..50 { @@ -447,7 +463,7 @@ fn spawn_server_with_timeouts( body_timeout: body, ..Config::new(addr) }; - serve_with(cfg, smarm::Config::exact(2), pipeline).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipeline, Vec::new()).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { @@ -479,7 +495,7 @@ fn spawn_server_with_body_gate( body_stall_timeout: stall, ..Config::new(addr) }; - serve_with(cfg, smarm::Config::exact(2), pipeline).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipeline, Vec::new()).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { @@ -811,7 +827,7 @@ fn stalled_reader_killed_at_write_timeout() { write_timeout: Duration::from_millis(300), ..Config::new(addr) }; - serve_with(cfg, smarm::Config::exact(2), pipe).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipe, Vec::new()).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { @@ -900,7 +916,7 @@ fn chunked_request_over_limit_413() { max_body_bytes: 8, // tiny ..Config::new(addr) }; - serve_with(cfg, smarm::Config::exact(2), pipe).unwrap(); + serve_with(cfg, smarm::Config::exact(2), pipe, Vec::new()).unwrap(); }); for _ in 0..50 { if TcpStream::connect(addr).is_ok() { break; } @@ -1438,15 +1454,14 @@ fn shutdown_force_stops_open_ws() { /// Minimal chat handler mirroring the example: on_open subscribes (conn /// actor pid) + spawns a relay holding ONLY the Receiver and a WsSender /// clone; on_message broadcast_from's, skipping the sender's own relay. -struct Chat { - bus: std::sync::Arc>>, -} +const CHAT_BUS: urus::PubSub = urus::PubSub::new("chat-test"); + +struct Chat; impl urus::WsHandler for Chat { fn on_open(&mut self, sender: &urus::WsSender) { - let bus = self.bus.get_or_init(urus::PubSub::new); - let rx = bus.subscribe("room").unwrap(); - let _ = bus.broadcast_from(smarm::self_pid(), "room", "joined".into()); + let rx = CHAT_BUS.subscribe("room").unwrap(); + let _ = CHAT_BUS.broadcast_from(smarm::self_pid(), "room", "joined".into()); let out = sender.clone(); smarm::spawn(move || { while let Ok(msg) = rx.recv() { @@ -1468,19 +1483,14 @@ impl urus::WsHandler for Chat { let _ = sender.send(urus::Message::Text("synced".into())); return; } - let bus = self.bus.get_or_init(urus::PubSub::new); - let _ = bus.broadcast_from(smarm::self_pid(), "room", t); + let _ = CHAT_BUS.broadcast_from(smarm::self_pid(), "room", t); } } } fn chat_pipeline() -> Pipeline { - let bus: std::sync::Arc>> = - std::sync::Arc::new(std::sync::OnceLock::new()); - Pipeline::new().plug(Router::new().get("/ws", move |c: Conn, _n: Next| { - let bus = bus.clone(); - c.upgrade(Chat { bus }) - })) + Pipeline::new() + .plug(Router::new().get("/ws", |c: Conn, _n: Next| c.upgrade(Chat))) } /// Two clients in a room: a broadcast reaches the other client and (via @@ -1489,7 +1499,7 @@ fn chat_pipeline() -> Pipeline { /// having sent a single frame when A's first message arrives. #[test] fn ws_chat_broadcast_reaches_other_client_not_sender() { - let port = spawn_server(chat_pipeline()); + let port = spawn_server_children(chat_pipeline(), vec![CHAT_BUS.child()]); let mut a = ws_connect(port); let mut a_buf = Vec::new(); @@ -1532,7 +1542,11 @@ fn ws_chat_broadcast_reaches_other_client_not_sender() { #[test] fn shutdown_with_open_chat_terminates() { let (port, handle, done_rx) = - spawn_server_with_handle(chat_pipeline(), Duration::from_millis(300)); + spawn_server_with_children( + chat_pipeline(), + Duration::from_millis(300), + vec![CHAT_BUS.child()], + ); let mut a = ws_connect(port); let mut a_buf = Vec::new(); @@ -1592,21 +1606,19 @@ mod channels_wire { } } - fn channels_pipeline() -> Pipeline { - let hub: std::sync::Arc>> = - std::sync::Arc::new(std::sync::OnceLock::new()); - Pipeline::new().plug(Router::new().get("/socket", move |c: Conn, _n: Next| { - // Hub construction is in-runtime only (it spawns the pubsub - // table); the NON-static OnceLock is what lets shutdown - // drain the table (the v0.5 pattern). - let hub = hub.clone(); - let hub = hub.get_or_init(|| { - ChannelHub::new(PrefixRouter::new().channel("room:*", |_: &str| { - Box::new(Lobby) as Box> - })) - }); - hub.upgrade(c) - })) + /// The hub is a description: no actors, so it is built once here and + /// its `children()` handed to the supervisor separately. + fn channels_hub() -> (ChannelHub

, Vec) { + ChannelHub::new( + "chan-test", + PrefixRouter::new() + .channel("room:*", |_: &str| Box::new(Lobby) as Box>), + ) + } + + fn channels_pipeline(hub: ChannelHub

) -> Pipeline { + Pipeline::new() + .plug(Router::new().get("/socket", move |c: Conn, _n: Next| hub.upgrade(c))) } const CHAN_HANDSHAKE: &[u8] = @@ -1634,7 +1646,8 @@ mod channels_wire { #[test] fn channels_join_heartbeat_event_broadcast_leave() { - let port = spawn_server(channels_pipeline()); + let (hub, children) = channels_hub(); + let port = spawn_server_children(channels_pipeline(hub), children); let mut a = chan_connect(port); let mut a_buf = Vec::new(); @@ -1720,7 +1733,8 @@ mod channels_wire { #[test] fn channels_codec_garbage_closes_1002() { - let port = spawn_server(channels_pipeline()); + let (hub, children) = channels_hub(); + let port = spawn_server_children(channels_pipeline(hub), children); let mut s = chan_connect(port); let mut buf = Vec::new(); ws_send(&mut s, &Frame::new(Opcode::Text, "not a v2 frame")); @@ -1737,7 +1751,14 @@ mod channels_wire { // is monitor-pruned -> the non-static hub's last handle drops // with the drained pipeline -> pubsub table exits -> AllDone. let (port, handle, done_rx) = - spawn_server_with_handle(channels_pipeline(), Duration::from_millis(300)); + { + let (hub, children) = channels_hub(); + spawn_server_with_children( + channels_pipeline(hub), + Duration::from_millis(300), + children, + ) + }; let mut a = chan_connect(port); let mut a_buf = Vec::new(); @@ -1765,19 +1786,20 @@ mod channels_wire { } } - fn session_pipeline() -> Pipeline { - let hub: std::sync::Arc>> = - std::sync::Arc::new(std::sync::OnceLock::new()); - Pipeline::new().plug(Router::new().get("/socket", move |c: Conn, _n: Next| { - let hub = hub.clone(); - let hub = hub.get_or_init(|| { - ChannelHub::new(PrefixRouter::new().channel_session::( - "room:*", - |_: &str| Box::new(Lobby) as Box>, - )) - }); - hub.upgrade(c) - })) + fn session_hub() -> (ChannelHub

, Vec) { + ChannelHub::new( + "sess-test-bus", + PrefixRouter::new().channel_session::( + "room:*", + "sess-test-registry", + |_: &str| Box::new(Lobby) as Box>, + ), + ) + } + + fn session_pipeline(hub: ChannelHub

) -> Pipeline { + Pipeline::new() + .plug(Router::new().get("/socket", move |c: Conn, _n: Next| hub.upgrade(c))) } #[test] @@ -1790,7 +1812,14 @@ mod channels_wire { // terminates, and exits -> AllDone. No links anywhere in that // chain; this test is the proof it composes. let (port, handle, done_rx) = - spawn_server_with_handle(session_pipeline(), Duration::from_millis(300)); + { + let (hub, children) = session_hub(); + spawn_server_with_children( + session_pipeline(hub), + Duration::from_millis(300), + children, + ) + }; let mut a = chan_connect(port); let mut a_buf = Vec::new();