- crud is now the demonstrator: app owns smarm runtime + root supervisor,
store actor and urus::endpoint as ordered siblings (store first, so
reverse-order shutdown drains HTTP before stopping the store),
Shutdown::Infinity on the endpoint, shutdown via
rt.handle().request_shutdown(root_sup) from the stdin thread.
Deletes two kludges the old shape forced:
* static OnceLock<Sender> spawn-on-first-use -> a supervised child
that self-registers a typed Name; handlers use smarm::send per
request and turn 'between incarnations' into a 503 instead of
panicking on a dropped store.
* static SHUTTING_DOWN AtomicBool + 250ms recv_timeout poll in the
store loop and in the SSE ticker -> a plain park; the tree stops
both. Smoke-tested live: CRUD round-trips, SSE stream, clean drain
with the stream open, port closed after.
- Other examples stay short and on serve*, updated for the split config
(serve_with(cfg, smarm::Config, pipe) / serve_with_shutdown(..., signal)).
plain_serve's URUS_SCHED_THREADS now builds a smarm::Config.
ws_chat's doc block explains the OnceLock is a serve*-only workaround
and points at crud for the clean shape.
- README: new 'Your Own Supervision Tree' section (endpoint as the real
API), graceful shutdown reframed as the serve*-only path, Config table
loses scheduler_threads and gains name, PubSub rule 1 notes the
supervised-sibling alternative.
111 lib + 50 integration + 2 doc tests green; clippy clean.
137 lines
4.9 KiB
Rust
137 lines
4.9 KiB
Rust
//! WebSocket chat rooms (v0.5): `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)
|
|
//!
|
|
//! The shapes this demonstrates:
|
|
//!
|
|
//! - **One `PubSub<String>` for the whole app**, topics are rooms
|
|
//! (`room:{name}`). Built lazily on the first connection via a
|
|
//! NON-static `Arc<OnceLock<PubSub>>` captured by the route closure,
|
|
//! because `PubSub::new()` spawns the table actor and smarm only
|
|
//! allows `spawn` in-runtime — while `serve_with_shutdown` owns the
|
|
//! runtime, there is no in-runtime moment before the first request.
|
|
//! Keep the cell non-static so the table is dropped with the drained
|
|
//! pipeline rather than pinned for the life of the process.
|
|
//!
|
|
//! If your app owns its own runtime and tree (see `crud`), skip this
|
|
//! entirely: start the table as a supervised sibling of
|
|
//! `urus::endpoint(...)` and address it by name.
|
|
//!
|
|
//! - **`on_open` subscribes and spawns the relay** — a listen-only
|
|
//! client receives the room without ever sending. The subscription is
|
|
//! 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};
|
|
|
|
use urus::{
|
|
serve_with_shutdown, shutdown_handle, Config, Conn, Message, Next, Pipeline, PubSub, Router,
|
|
WsHandler, WsSender,
|
|
};
|
|
|
|
type Bus = PubSub<String>;
|
|
|
|
struct ChatHandler {
|
|
bus: Arc<OnceLock<Bus>>,
|
|
room: String,
|
|
}
|
|
|
|
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) {
|
|
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}"));
|
|
|
|
// 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.
|
|
let out = sender.clone();
|
|
smarm::spawn(move || {
|
|
while let Ok(msg) = rx.recv() {
|
|
if out.send(Message::Text((*msg).clone())).is_err() {
|
|
break; // client gone or server draining
|
|
}
|
|
}
|
|
});
|
|
}
|
|
|
|
fn on_message(&mut self, msg: Message, sender: &WsSender) {
|
|
let Message::Text(text) = msg else {
|
|
let _ = sender.close(1003, "text only");
|
|
return;
|
|
};
|
|
// 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);
|
|
}
|
|
|
|
fn on_close(&mut self, _code: Option<u16>, _reason: &str) {
|
|
let topic = self.topic();
|
|
let _ = self
|
|
.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).
|
|
}
|
|
}
|
|
|
|
fn main() {
|
|
let bus: Arc<OnceLock<Bus>> = Arc::new(OnceLock::new());
|
|
|
|
let pipeline = Pipeline::new().plug(Router::new().get("/chat/:room", move |c: Conn, _n: Next| {
|
|
let room = c.params.get("room").unwrap_or("lobby").to_string();
|
|
let bus = bus.clone();
|
|
c.upgrade(ChatHandler { bus, room })
|
|
}));
|
|
|
|
let (handle, signal) = shutdown_handle();
|
|
std::thread::spawn(move || {
|
|
println!("ws_chat on ws://127.0.0.1:8080/chat/:room — press Enter to shut down");
|
|
let mut line = String::new();
|
|
let _ = std::io::stdin().read_line(&mut line);
|
|
handle.shutdown();
|
|
});
|
|
|
|
serve_with_shutdown(
|
|
Config::new("0.0.0.0:8080".parse().unwrap()),
|
|
smarm::Config::default(),
|
|
pipeline,
|
|
signal,
|
|
)
|
|
.unwrap();
|
|
println!("ws_chat: drained, bye");
|
|
}
|