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();