From bfa513cd6da32915b76b9dedac286eade7586823 Mon Sep 17 00:00:00 2001 From: smarm Date: Sat, 20 Jun 2026 13:34:17 +0200 Subject: [PATCH] doc: rework gen_server module docs completely --- src/gen_server.rs | 521 ++++++++++++++++++++++++++++------------------ 1 file changed, 317 insertions(+), 204 deletions(-) diff --git a/src/gen_server.rs b/src/gen_server.rs index 58c04fa..22954d1 100644 --- a/src/gen_server.rs +++ b/src/gen_server.rs @@ -1,71 +1,182 @@ -//! gen_server — synchronous call / asynchronous cast over a single actor. +//! A stateful actor with typed request-reply messaging. //! -//! A thin request-reply layer on top of [`channel`](crate::channel()), modelled -//! on Erlang's `gen_server`. A *server* is an actor owning a state value that -//! implements [`GenServer`]; clients hold a clonable [`GenServerRef`] and issue -//! [`call`](GenServerRef::call) (synchronous, returns a reply) or -//! [`cast`](GenServerRef::cast) (fire-and-forget). +//! ## What is a gen_server? //! -//! ## One inbox, many arms +//! Whenever you need mutable state shared between threads, the obvious tool is +//! `Arc>`. That works, but it scatters lock guards and error +//! handling everywhere, and offers no natural place to put the logic that +//! operates on the state. //! -//! Every call and cast travels the *same* inbox channel as an [`Envelope`] -//! (calls and casts are ordered relative to each other, like Erlang); the -//! server loop dispatches by variant. A `call` carries a freshly made -//! one-shot reply channel; the server sends the reply straight back down it. +//! A gen_server is a cleaner alternative: a dedicated thread that *owns* the +//! state and handles one message at a time. Other threads talk to it by sending +//! typed messages and optionally waiting for a typed reply. Because messages are +//! serialised through a single inbox. The server is the only thing that ever +//! touches the state, so there are no lock guards at all. //! -//! Out-of-band messages ride *separate* channels composed at the wait via -//! [`select`](channel::select): info channels handed over at start -//! ([`GenServerBuilder::with_info`]) are dispatched to -//! [`handle_info`](GenServer::handle_info). Arm priority is -//! **infos before inbox**, in declaration order — a hot inbox cannot starve -//! an out-of-band message; conversely a hot info channel CAN starve the -//! inbox, deliberately (system-message semantics). An info channel whose -//! senders are all gone is silently dropped from the arm set (per the -//! closed-arm-is-ready-forever rule on `select`); a closed *inbox* still -//! means graceful shutdown. +//! You write the logic by implementing [`GenServer`]. smarm runs the receive +//! loop, the reply plumbing, and the lifecycle. //! -//! ## Server death +//! ## A first server //! -//! Detection falls out of channel closure, so no monitor is required: -//! - if the server is already gone, its inbox is closed and the `send` in -//! `call`/`cast` fails → [`CallError::ServerDown`] / [`CastError::ServerDown`]; -//! - if the server dies *after* a call is enqueued but before it replies -//! (a handler panic, or a cooperative `request_stop`), the reply sender is -//! dropped as the server's stack unwinds, closing the reply channel; the -//! parked caller wakes and its `recv` returns `Err` → `ServerDown`. +//! Let's build a counter that can be incremented from any thread and queried for +//! its current value. //! -//! ## Lifecycle / callbacks +//! ```ignore +//! use smarm::gen_server::{self, GenServer, GenServerRef}; //! -//! [`GenServer::init`] runs once inside the server actor before the first -//! message; [`GenServer::terminate`] runs on the way out. terminate is wired -//! through a drop guard, so it fires on *every* exit path — graceful inbox -//! close, a handler panic, or a cooperative `request_stop` — not only the clean -//! one. Keep it cheap and non-blocking: it may run mid-unwind, and a panic -//! inside it during an unwind aborts the process (a double panic). +//! // State //! -//! ## Time (timers and idle) +//! struct Counter { +//! count: u64, +//! } //! -//! A server arms timers through a [`TimerHandle`] cloned from the [`GenServerCtx`] -//! in `init` and stored on the state — the same shape as [`Watcher`] for -//! monitors. [`arm_after`](TimerHandle::arm_after) is a one-shot, -//! [`tick_every`](TimerHandle::tick_every) a periodic; both fire into -//! [`handle_timer`](GenServer::handle_timer) at *system priority* (above infos -//! and the inbox, by arm position), and [`cancel`](TimerHandle::cancel) carries -//! the substrate's race signal. Separately, [`GenServerCtx::idle_after`] sets a -//! receive-timeout window: quiet for the whole window fires -//! [`handle_idle`](GenServer::handle_idle). +//! // Message types //! -//! Two unrelated things share the word *timeout*: the server-side **idle / -//! receive timeout** above, and the client-side **call deadline** -//! ([`GenServerRef::call_timeout`]) — how long a caller waits for a reply. They sit -//! on different axes and never interact (RFC 015 §7). +//! // Calls expect a reply; casts do not. +//! enum CounterCall { Get } +//! enum CounterCast { Increment } //! -//! ## Not here (yet) +//! // Behaviour //! -//! No dynamic info subscription: the info-channel set is fixed at start. -//! Revisit against a real consumer. (Monitor `Down` forwarding is dynamic — -//! see `handle_down` — because monitors are inherently created at runtime.) -//! The idle window is likewise set once, in `init` (RFC 015 §4.4). +//! impl GenServer for Counter { +//! type Call = CounterCall; +//! type Reply = u64; +//! type Cast = CounterCast; +//! type Info = (); // no out-of-band messages +//! type Timer = (); // no timers +//! +//! fn handle_call(&mut self, request: CounterCall) -> u64 { +//! match request { +//! CounterCall::Get => self.count, +//! } +//! } +//! +//! fn handle_cast(&mut self, request: CounterCast) { +//! match request { +//! CounterCast::Increment => self.count += 1, +//! } +//! } +//! } +//! ``` +//! +//! Start the server and talk to it: +//! +//! ```ignore +//! let server: GenServerRef = gen_server::start(Counter { count: 0 }); +//! +//! server.cast(CounterCast::Increment).unwrap(); +//! server.cast(CounterCast::Increment).unwrap(); +//! +//! let n = server.call(CounterCall::Get).unwrap(); +//! // n == 2 +//! ``` +//! +//! [`GenServerRef`] is cheaply clonable; hand copies to as many threads as you +//! like. They all share the same inbox; the server handles messages one at a +//! time, in arrival order. +//! +//! ## Client / server APIs +//! +//! In practice, callers should not have to know about `CounterCall` or +//! `CounterCast`. The recommended pattern is to wrap the message types in plain +//! functions on the same module: +//! +//! ```ignore +//! // Client API +//! +//! impl Counter { +//! pub fn start() -> GenServerRef { +//! gen_server::start(Counter { count: 0 }) +//! } +//! +//! pub fn increment(server: &GenServerRef) { +//! server.cast(CounterCast::Increment).unwrap(); +//! } +//! +//! pub fn get(server: &GenServerRef) -> u64 { +//! server.call(CounterCall::Get).unwrap() +//! } +//! } +//! +//! // Callers now just do: +//! let counter = Counter::start(); +//! Counter::increment(&counter); +//! Counter::increment(&counter); +//! let n = Counter::get(&counter); // n == 2 +//! ``` +//! +//! ## call vs cast +//! +//! `call` sends a request and blocks the calling thread until the server +//! replies. Use it when you need a return value, or when you need to know that +//! the server has processed the message before continuing. +//! +//! `cast` enqueues a message and returns immediately, without waiting for the +//! server to handle it. Use it for fire-and-forget updates where you don't need +//! confirmation. +//! +//! Both return `Err(ServerDown)` if the inbox is closed (the server is gone). +//! If you need to cap how long you wait for a reply, use +//! [`GenServerRef::call_timeout`], which additionally returns `Err(Timeout)` if +//! the deadline passes before a reply arrives. +//! +//! ## Lifecycle +//! +//! Two optional callbacks bracket the server's life: +//! +//! - [`GenServer::init`] runs once before the first message. Use it to start +//! timers or set up monitors; see the [`GenServerCtx`] it receives. +//! - [`GenServer::terminate`] runs when the server is about to exit. It fires +//! on every exit path (all `GenServerRef`s dropped, a handler panic, or an +//! explicit [`GenServerRef::shutdown`]), not only on clean shutdown. Keep it +//! short and non-blocking: if `terminate` panics while the server is already +//! unwinding from a handler panic, the process aborts. +//! +//! ## When the server stops +//! +//! The server runs as long as at least one [`GenServerRef`] exists. When the last +//! one is dropped, the inbox closes and the loop exits gracefully. To stop a +//! server explicitly and wait for it to finish, call [`GenServerRef::shutdown`]. +//! +//! If the server panics inside a handler, the panic unwinds the server thread. +//! Any caller currently waiting in `call` sees `Err(ServerDown)`: the reply +//! channel closes as the server unwinds, which wakes the caller. +//! +//! ## Going further +//! +//! Out-of-band messages: a server can receive messages from sources other +//! than `call`/`cast`, for example notifications from a background task. Pass +//! extra [`Receiver`] channels to [`GenServerBuilder::with_info`]; they are +//! dispatched to [`GenServer::handle_info`], always before inbox messages. +//! +//! Monitors: to watch another actor and be notified when it exits, clone a +//! [`Watcher`] from [`GenServerCtx::watcher`] during `init` and call +//! [`Watcher::watch`] with a [`Monitor`] from any handler. The loop dispatches +//! the resulting [`Down`] to [`GenServer::handle_down`]. +//! +//! Timers: clone a [`TimerHandle`] from [`GenServerCtx::timer`] during `init`. +//! Use [`TimerHandle::arm_after`] for a one-shot and [`TimerHandle::tick_every`] +//! for a periodic; both fire into [`GenServer::handle_timer`]. +//! +//! Idle detection: call [`GenServerCtx::idle_after`] during `init` to set a +//! quiet window: if no message of any kind is dispatched for that duration, +//! the loop calls [`GenServer::handle_idle`] and resets the window. This differs +//! from the client-side `call_timeout`. Idle measures silence on the *server*, +//! while `call_timeout` caps how long *one caller* waits. +//! +//! Names and registration: a server can be given a static name so other +//! actors can reach it without holding a `GenServerRef`. Use +//! [`GenServerBuilder::named`] to register on start, and the free functions +//! [`call`], [`cast`], and [`whereis_server`] to address it by name. Registered +//! servers are a natural fit for supervision; see `supervisor` for how to +//! build a tree that restarts servers on failure. +//! +//! ## Limitations +//! +//! The info-channel set is fixed at start; channels cannot be added or removed +//! while the server is running. Monitors are dynamic (they can be registered +//! from any handler via [`Watcher::watch`]) because monitors are inherently +//! created at runtime. The idle window is set once, in `init`. use crate::channel::{channel, select, select_timeout, Receiver, RecvTimeoutError, Selectable, Sender}; use crate::monitor::{demonitor, monitor, Down, Monitor}; @@ -79,37 +190,33 @@ use std::marker::PhantomData; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; -/// Behaviour for a gen_server: a state value plus call/cast handlers. +/// The behaviour you implement to make a type into a gen_server. /// -/// `handle_call` and `handle_cast` are required; [`init`](Self::init) and -/// [`terminate`](Self::terminate) are optional lifecycle hooks with no-op -/// defaults. +/// Implement this on your state struct. `handle_call` and `handle_cast` are +/// required; all other methods have no-op defaults and can be added as needed. pub trait GenServer: Send + 'static { - /// Synchronous request type (carried by [`GenServerRef::call`]). + /// The request type sent by [`GenServerRef::call`]. Must produce a [`Reply`](Self::Reply). type Call: Send + 'static; - /// Reply type returned for a `Call`. + /// The value returned to the caller by [`handle_call`](Self::handle_call). type Reply: Send + 'static; - /// Asynchronous request type (carried by [`GenServerRef::cast`]). + /// The request type sent by [`GenServerRef::cast`]. No reply is produced. type Cast: Send + 'static; - /// Out-of-band message type, delivered to [`handle_info`](Self::handle_info) - /// from the info channels registered at start - /// ([`GenServerBuilder::with_info`]). Servers with several out-of-band - /// sources enum them up into one `Info`. Use `()` if unused. + /// Out-of-band message type delivered to [`handle_info`](Self::handle_info) + /// from channels registered via [`GenServerBuilder::with_info`]. If you don't + /// use info channels, set this to `()`. type Info: Send + 'static; - /// The server's own scheduled-timer payload, delivered to - /// [`handle_timer`](Self::handle_timer) when a timer armed via the loop's - /// [`TimerHandle`] fires (RFC 015). Kept distinct from - /// [`Info`](Self::Info) — `Info` is *external* out-of-band traffic, `Timer` - /// is the server's *own* fires — so the system/userspace split stays - /// legible and a tagged timer round-trips - /// (`arm_after(d, Tk::Retry(n))` → `handle_timer(Tk::Retry(n))`). Use `()` - /// if unused, exactly as `Info` does. + /// Payload type for timers armed via [`TimerHandle`], delivered to + /// [`handle_timer`](Self::handle_timer). Kept separate from [`Info`](Self::Info) + /// so that timer fires (which your server schedules itself) stay distinct + /// from external messages (which arrive from outside). Use `()` if unused. type Timer: Send + 'static; /// Runs once inside the server actor before any message is handled. The - /// [`GenServerCtx`] is the loop's one runtime hook: clone its [`Watcher`] - /// into the state here to be able to [`watch`](Watcher::watch) monitors - /// from any later handler. + /// [`GenServerCtx`] is the loop's one runtime hook: clone whichever of + /// [`ctx.watcher()`](GenServerCtx::watcher) (monitors) and + /// [`ctx.timer()`](GenServerCtx::timer) (timers) you will need from later + /// handlers into the state here. If you need neither, the unused ctx drops + /// and the loop's system arm auto-closes (see module docs). fn init(&mut self, _ctx: &GenServerCtx) where Self: Sized, @@ -130,19 +237,19 @@ pub trait GenServer: Send + 'static { /// [`Watcher::watch`]. Default: drop it. fn handle_down(&mut self, _down: Down) {} - /// Handle a fired timer armed through the loop's [`TimerHandle`] + /// Handle a fired timer armed through [`TimerHandle`] /// ([`arm_after`](TimerHandle::arm_after) / - /// [`tick_every`](TimerHandle::tick_every)). Timer fires re-enter at system - /// priority — above infos and the inbox — so a heartbeat cannot be starved - /// by userspace traffic. Default: drop it (RFC 015 §6). + /// [`tick_every`](TimerHandle::tick_every)). Timer fires are delivered + /// before info and inbox messages, so a heartbeat cannot be starved by + /// application traffic. Default: drop the message. fn handle_timer(&mut self, _msg: Self::Timer) {} - /// Handle a receive/idle timeout: fired when the loop has waited a full - /// idle window (set once via [`GenServerCtx::idle_after`]) with no message of - /// any kind dispatched. The window resets on every dispatched message and - /// re-arms after this fires (a steady idle detector); a server wanting - /// one-shot idle-shutdown simply requests its own stop here. Default: no-op - /// (RFC 015 §4.4, §6). + /// Handle a quiet-period notification: called when the server has gone the + /// full idle window (set via [`GenServerCtx::idle_after`] in `init`) without + /// dispatching any message. The window resets automatically after this + /// fires, so it acts as a steady idle detector. To shut down after one idle + /// period, call [`request_stop`](crate::scheduler::request_stop) here. + /// Default: no-op. fn handle_idle(&mut self) {} /// Runs as the server actor exits, on any exit path (see module docs). @@ -150,7 +257,7 @@ pub trait GenServer: Send + 'static { } /// What travels the server's single inbox channel: a synchronous call (with a -/// reply sender) or an asynchronous cast. +/// reply sender) or an asynchronous cast. Private — callers use [`GenServerRef`]. enum Envelope { Call(G::Call, Sender), Cast(G::Cast), @@ -182,11 +289,10 @@ pub enum CallError { pub enum CallTimeoutError { /// The server was already gone, or died before replying. ServerDown, - /// The deadline passed before a reply arrived. The request stays in the - /// server's inbox: it will still be *handled*, but the reply is discarded - /// (the abandoned reply channel's receiver is dropped, so the server's - /// reply send fails harmlessly). Erlang behaves the same way; design - /// idempotent calls accordingly. + /// The deadline passed before a reply arrived. The request is still in the + /// server's inbox and will be handled — the reply is simply discarded + /// because the reply channel was dropped when the timeout fired. Design + /// calls that may time out to be idempotent, so a late reply causes no harm. Timeout, } @@ -214,25 +320,17 @@ impl GenServerRef { reply_rx.recv().map_err(|_| CallError::ServerDown) } - /// Bounded synchronous request-reply: like [`call`](Self::call), but - /// gives up after `timeout`, returning [`CallTimeoutError::Timeout`]. + /// Like [`call`](Self::call), but gives up after `timeout` and returns + /// [`CallTimeoutError::Timeout`] if no reply arrives in time. /// - /// This `timeout` is a **client-side call deadline** — how long *this caller* - /// waits for a reply — and is wholly separate from the server-side idle / - /// receive timeout ([`GenServerCtx::idle_after`] → - /// [`GenServer::handle_idle`]), which measures quiet on the *server's* inbox. - /// Same word ("timeout"), two different axes (RFC 015 §7); neither touches - /// the other. + /// Note that "timeout" means two different things here and in + /// [`GenServerCtx::idle_after`]: this one is a *client-side call deadline* — + /// how long this particular caller waits. The idle timeout measures silence + /// on the *server's* inbox. They are completely independent. /// - /// The roadmap sketched this as monitor + wait-reply-or-Down + demonitor; - /// that machinery is unnecessary here because server death is already - /// observable on the reply channel itself — the reply sender is dropped - /// as the server's stack unwinds, closing the channel and waking the - /// parked caller (see the module docs). So a bounded call is exactly - /// [`Receiver::recv_timeout`] on the reply channel: `Ok` on a reply, - /// `Disconnected` → [`CallTimeoutError::ServerDown`], `Timeout` → - /// [`CallTimeoutError::Timeout`]. Nothing is registered, so nothing can - /// leak on the timeout path by construction. + /// On timeout, the server still handles the request; only the reply is + /// discarded (the reply channel is dropped, so the server's send fails + /// silently). Design timed-out calls to be idempotent. pub fn call_timeout( &self, request: G::Call, @@ -257,18 +355,17 @@ impl GenServerRef { .map_err(|_| CastError::ServerDown) } - /// Ask the server to terminate and block until it has — the `sys`-style - /// stop (cf. Erlang's `gen_server:stop/1`). Sends a cooperative stop to the - /// server actor and waits for its `Down`, so [`GenServer::terminate`] has - /// run by the time this returns (it fires from the loop's drop guard on the - /// stop unwind). Returns immediately if the server was already gone. + /// Stop the server and block until it has fully exited. /// - /// This is the explicit teardown for a server pinned alive by a registered - /// [`GenServerName`] (whose stored sender means dropping every external - /// [`GenServerRef`] no longer closes the inbox). Best-effort like all - /// cooperative cancellation: a server wedged in a tight loop with no - /// observation point cannot be stopped. Panics if called outside - /// `Runtime::run()`. + /// Sends a cooperative stop signal to the server actor and waits for it to + /// exit, so [`GenServer::terminate`] has run by the time this returns. + /// Returns immediately if the server is already gone. + /// + /// This is the right teardown for a server kept alive by a registered + /// [`GenServerName`], where dropping every external `GenServerRef` is not enough + /// to close the inbox. Like all cooperative cancellation, it is best-effort: + /// a server wedged in a tight loop with no observation point cannot be + /// stopped this way. Panics if called outside `Runtime::run()`. pub fn shutdown(&self) { let mon = monitor(self.pid); request_stop(self.pid); @@ -279,29 +376,28 @@ impl GenServerRef { } } -/// The server loop's internal *system intake* channel (RFC 015 §4.5). Control -/// (monitor handoff) and armed-timer fires are structurally the same — both are -/// loop-internal sources selected above the inbox — so they fold into one -/// channel, dispatched by variant. The arm stays open while any [`Watcher`] or -/// [`TimerHandle`] clone (or an in-flight fire thunk) still holds a sender. +/// Internal channel carrying control messages to the server loop. Monitors +/// and timer fires both need to land above the inbox in priority, so they +/// share one channel dispatched by variant. The arm stays open while any +/// [`Watcher`] or [`TimerHandle`] clone (or an in-flight timer fire) holds a sender. enum Sys { /// A monitor handed in via [`Watcher::watch`]; pushed onto the loop's /// monitor set. Watch(Monitor), - /// A fired one-shot timer carrying its loop-local id (so the loop can retire - /// it from the live set) and the server's payload, dispatched to - /// [`GenServer::handle_timer`]. + /// A fired one-shot timer: loop-local id (so the loop can retire the + /// registry entry) plus the server's payload. Timer(crate::timer::TimerId, G::Timer), - /// A fired periodic tick carrying its stable local id. The loop resolves the - /// payload from the registry's factory, dispatches it to - /// [`GenServer::handle_timer`], and re-arms the next tick (RFC 015 §4.3). + /// A fired periodic tick carrying its stable local id. The loop looks up + /// the payload factory, dispatches it to [`GenServer::handle_timer`], and + /// re-arms the next tick before returning. Tick(crate::timer::TimerId), } /// The server loop's runtime hook, passed to [`GenServer::init`]. Hands out the /// loop's two clonable intake handles — the [`Watcher`] (monitors) and the -/// [`TimerHandle`] (timers) — plus the one-shot idle-window setter. Opaque, so -/// fields can grow without breaking. +/// [`TimerHandle`] (timers) — plus the one-shot idle-window setter. All fields +/// are private so the struct can gain new hooks in future without breaking +/// existing implementations. pub struct GenServerCtx { sys_tx: Sender>, reg: Arc>>, @@ -333,27 +429,25 @@ impl GenServerCtx { TimerHandle { sys_tx: self.sys_tx.clone(), reg: self.reg.clone() } } - /// Set the idle / receive-timeout window: if no message of any kind is - /// dispatched for `after`, the loop fires [`GenServer::handle_idle`] - /// (RFC 015 §4.4). Set once, in `init`; the window is loop-owned and fixed - /// (dynamic per-message reconfiguration is a non-goal). The window resets on - /// every dispatched message and re-arms after `handle_idle` (a steady idle - /// detector); a server wanting one-shot idle-shutdown requests its own stop - /// in `handle_idle`. + /// Set a quiet-period window: if the loop goes `after` without dispatching + /// any message, it calls [`GenServer::handle_idle`] and resets the window. + /// Call this once during `init`; the window is fixed for the server's + /// lifetime. /// - /// Distinct from [`GenServerRef::call_timeout`], which is a *client-side* call - /// deadline — same word, different axis (§7). + /// This is distinct from [`GenServerRef::call_timeout`], which caps how long a + /// single caller waits for a reply. This one measures silence on the whole + /// inbox. pub fn idle_after(&self, after: Duration) { self.idle.set(Some(after)); } } /// Per-server timer bookkeeping, shared between the loop and every -/// [`TimerHandle`] clone. A gen_server actor is single-threaded — handle calls -/// (made from inside handlers) and the loop never run concurrently — so this -/// `Mutex` is always uncontended; it is `Arc` rather than `Rc` -/// only because the handle is stored on `self` and `GenServer: Send` forces the -/// handle (hence its shared state) to be `Send`. +/// [`TimerHandle`] clone. A gen_server actor is single-threaded — handlers +/// and the loop never run concurrently — so this `Mutex` is always +/// uncontended at runtime; `Arc` (rather than `Rc`) is used +/// only because `GenServer: Send` forces the handle (hence this shared state) +/// to be `Send`. struct TimerReg { /// Monotonic minter for loop-local [`TimerId`]s — the ids handed to users, /// kept distinct from the substrate `seq` (which changes on every periodic @@ -379,15 +473,14 @@ struct TimerReg { /// One periodic timer's loop-side bookkeeping. struct Periodic { - /// Re-arm interval; each fire schedules the next at `now + every` (§4.3). + /// Re-arm interval; each fire schedules the next at `now + every`. every: Duration, - /// The currently-armed substrate timer for this periodic (changes on every + /// The currently-armed underlying timer for this periodic (changes on every /// re-arm); cancelled by [`cancel`](TimerHandle::cancel). live: crate::timer::TimerId, /// Produces a fresh payload for each tick. Built in `tick_every` as - /// `move || msg.clone()` — so the `Clone` bound lives there, on the one - /// method that needs it, and never leaks onto `type Timer` or the loop - /// (which is monomorphised per server and only *calls* this box). + /// `move || msg.clone()` — so the `Clone` bound lives on that one method + /// and never appears on `type Timer` or anywhere else in the loop. make: Box G::Timer + Send>, } @@ -413,11 +506,10 @@ impl TimerReg { } } -/// Arms, cancels, and (chunk 4) periodically ticks server timers, the time-side -/// twin of [`Watcher`] (RFC 015 §4.3). Handed out by [`GenServerCtx::timer`] in -/// `init` and stored on the state; later handlers arm timers without any change -/// to their signatures. Clonable; the arm stays open while any clone (or an -/// in-flight fire) lives. +/// Arms and cancels timers for a server loop, the time-side twin of [`Watcher`]. +/// Handed out by [`GenServerCtx::timer`] in `init`; store it on the state so +/// handlers can arm timers without changing their signatures. Clonable; the +/// loop's timer arm stays open while any clone (or an in-flight fire) lives. pub struct TimerHandle { sys_tx: Sender>, reg: Arc>>, @@ -431,11 +523,10 @@ impl Clone for TimerHandle { } impl TimerHandle { - /// Arm a one-shot timer: deliver `msg` to [`GenServer::handle_timer`] after - /// `after`, unless [`cancel`](Self::cancel)led first. Returns a loop-local - /// [`TimerId`] for cancellation. A thin wrapper over the `send_after` - /// substrate that lands the fire on the loop's own system arm (above infos - /// and the inbox), not in the inbox. + /// Arm a one-shot timer: deliver `msg` to [`GenServer::handle_timer`] once, + /// after `after`, unless [`cancel`](Self::cancel)led first. Returns a + /// [`TimerId`] you can pass to `cancel`. Timer fires arrive before info and + /// inbox messages. pub fn arm_after(&self, after: Duration, msg: G::Timer) -> TimerId { let mut reg = self.reg.lock().unwrap(); let local = reg.mint(); @@ -448,15 +539,15 @@ impl TimerHandle { /// Arm a periodic timer: deliver a fresh `msg` to /// [`GenServer::handle_timer`] every `every`, until - /// [`cancel`](Self::cancel)led. Loop-managed sugar over the one-shot - /// substrate (RFC 015 §4.3): the substrate stays one-shot; the loop re-arms - /// at `now + every` after each fire and exposes one stable [`TimerId`], so a - /// single `cancel` stops the re-arm *and* the pending instance. + /// [`cancel`](Self::cancel)led. The underlying timer is one-shot; the loop + /// re-arms it after each fire to produce the repeating cadence, and exposes + /// a single stable [`TimerId`] so one `cancel` stops both the pending + /// instance and all future ones. /// - /// Requires `Self::Timer: Clone` — a periodic re-delivers the same logical - /// message each tick, so the loop needs a fresh copy per period. The bound - /// sits on this method alone; one-shot [`arm_after`](Self::arm_after) and - /// the `type Timer` declaration stay unconstrained. + /// Requires `G::Timer: Clone` because the same logical message is cloned + /// fresh for each tick. This bound is on this method only; one-shot + /// [`arm_after`](Self::arm_after) and the `type Timer` declaration are + /// unconstrained. pub fn tick_every(&self, every: Duration, msg: G::Timer) -> TimerId where G::Timer: Clone, @@ -473,15 +564,18 @@ impl TimerHandle { // from `make` on fire, so the tick carries only the stable id. let sub = send_after_to(every, self.sys_tx.clone(), Sys::Tick(local)); reg.periodics.insert(local, Periodic { every, live: sub, make }); + debug_assert!(reg.rearm_tx.is_some(), "rearm_tx must be Some while periodics is non-empty"); local } /// Cancel an armed timer (one-shot or periodic). Returns the substrate's - /// race signal — `true` if the cancel beat the (pending) fire, `false` if - /// it had already fired / been cancelled / is unknown. For a periodic this - /// also stops the re-arm: the entry is dropped, so the loop will not - /// schedule another tick even if the pending one already escaped onto the - /// channel (that in-flight tick is discarded on dispatch). + /// race signal — `true` if the cancel beat the fire, `false` if the timer + /// had already fired, been cancelled, or is otherwise unknown. Most callers + /// can ignore the return value: a fired-but-not-yet-dispatched one-shot is + /// discarded by the loop on delivery (see the `Tick` arm in `server_loop`), + /// and a cancelled periodic will not re-arm. For a periodic, `cancel` also + /// stops re-arming: the entry is removed so the loop will not schedule + /// another tick even if the pending one already escaped onto the channel. pub fn cancel(&self, id: TimerId) -> bool { let mut reg = self.reg.lock().unwrap(); if let Some(sub) = reg.oneshots.remove(&id) { @@ -492,6 +586,10 @@ impl TimerHandle { if reg.periodics.is_empty() { reg.rearm_tx = None; } + debug_assert!( + reg.rearm_tx.is_some() == !reg.periodics.is_empty(), + "rearm_tx must be Some iff periodics is non-empty" + ); return beat; } false @@ -573,10 +671,10 @@ impl GenServerBuilder { NamedGenServerBuilder { builder: self, name: name.as_str() } } - /// The shared spawn body behind [`start`](Self::start) and - /// [`NamedGenServerBuilder::start`]: make the inbox, spawn the loop, return the - /// ref. (The named path additionally publishes the inbox sender under the - /// name before returning.) + /// Private shared body behind [`start`](Self::start) and + /// [`NamedGenServerBuilder::start`]: allocate the inbox, spawn the loop, + /// return the ref. The named path additionally publishes the inbox sender + /// under the name before returning. fn spawn_server(self) -> GenServerRef { let (tx, rx) = channel::>(); let GenServerBuilder { state, infos, supervisor } = self; @@ -588,19 +686,20 @@ impl GenServerBuilder { } } -/// A durable name typed by the *server* (RFC 014): a gen_server is multi-message -/// (call / cast over one inbox), so it is addressed by `G` rather than by a -/// single message type. By-name [`call`] / [`cast`] check the request against -/// `G`'s `Call` / `Cast` / `Reply`. Declared as a constant and shared freely: +/// A typed, static name for a gen_server, used to address it through the +/// registry without holding a [`GenServerRef`]. Because a gen_server handles +/// multiple message types, the name is typed by the whole server (`G`) rather +/// than by a single message type. Declare it as a constant: /// /// ```ignore /// const COUNTER: GenServerName = GenServerName::new("counter"); /// ``` /// -/// Under the hood the server's inbox is published into the registry as a -/// `Sender>` keyed by its message `TypeId` — the same channel store -/// every other name uses — so naming needs no separate directory. `Envelope` -/// stays private: only `GenServerName` opens that door. +/// Use [`GenServerBuilder::named`] to bind the name on start, and the free +/// functions [`call`], [`cast`], and [`whereis_server`] to address the server +/// by name. Under the hood the server's inbox sender is stored in the registry +/// keyed by name and message `TypeId`; `Envelope` is private, so only +/// `GenServerName` can open that slot. pub struct GenServerName { name: &'static str, _marker: PhantomData G>, @@ -732,12 +831,17 @@ fn server_loop( state: G, mut infos: Vec>, ) { - // terminate() must run on every exit path (clean close, panic, stop), so it - // lives in this guard's Drop rather than after the loop. The guard also owns - // a registry handle so the *same* drop drains every live timer (RFC 015 - // §4.7): once the loop is gone no periodic can re-arm anyway, but cancelling - // here keeps no armed substrate entry lingering past the server, and the - // assert pins the invariant. + // Drop guard — owns the server state and the timer registry. + // + // Why a guard rather than code after the loop: + // - `terminate()` must fire on *every* exit path: clean inbox close, + // a handler panic, and cooperative `request_stop`. A guard's `Drop` + // covers all three; code after the loop only covers the clean path. + // - The timer drain must run *before* `terminate()`: terminate may + // inspect state but must not arm new timers into a dead loop. By + // owning the registry here the drain and the terminate call are + // sequenced correctly and can never be reordered. + // - A debug_assert pins the post-drain invariant cheaply. struct Terminate(G, Arc>>); impl Drop for Terminate { fn drop(&mut self) { @@ -761,6 +865,9 @@ fn server_loop( } } + // Unwrap an envelope and hand it to the right handler. For calls, a failed + // reply send means the caller went away (e.g. cancelled while parked) — + // that is the caller's problem, not the server's. fn dispatch(state: &mut G, env: Envelope) { match env { Envelope::Call(request, reply_tx) => { @@ -807,8 +914,9 @@ fn server_loop( loop { if monitors.is_empty() && !sys_open && infos.is_empty() { - // Nothing to select over: park on the inbox alone (the pre-v0.8 - // loop), with the idle window as the recv timeout when one is set. + // Fast path: no extra arms, no select overhead — park directly on + // the inbox. Mirrors the inbox arm of the select path below; any + // change there must be applied here too. match idle_deadline { Some(dl) => { let wait = dl.saturating_duration_since(Instant::now()); @@ -832,12 +940,15 @@ fn server_loop( }, } } else { - // Arm priority: downs, then the system arm, then infos (each in - // declaration order), then the inbox. The arm slice is rebuilt - // per iteration because every set but the inbox shrinks or grows. - // The whole wait carries the idle window as its timeout. - let nd = monitors.len(); - let nw = sys_open as usize; + // Slow path: one or more extra arms live — build the arm slice and + // select. Arm order encodes priority: downs → system → infos → + // inbox. The slice is rebuilt each iteration because the monitor + // and info sets shrink/grow. Mirrors the fast-path inbox park + // above; keep them in sync. + let nd = monitors.len(); // monitor band: [0, nd) + let nw = sys_open as usize; // system arm: [nd, nd+nw) + // info band: [nd+nw, nd+nw+ni) + // inbox arm: [nd+nw+ni] let sel = { let mut arms: Vec<&dyn Selectable> = Vec::with_capacity(nd + nw + infos.len() + 1); @@ -869,8 +980,8 @@ fn server_loop( } }; if i < nd { - // A Down retires its arm either way: delivered (one-shot) or - // closed without delivering (defensive; shouldn't happen). + // Monitor band: a Down retires its arm either way (one-shot) + // or closes without delivering (defensive; shouldn't happen). let m = monitors.remove(i); if let Ok(Some(down)) = m.rx.try_recv() { guard.0.handle_down(down); @@ -924,6 +1035,7 @@ fn server_loop( Err(_) => sys_open = false, } } else if i < nd + nw + infos.len() { + // Info band. let j = i - nd - nw; match infos[j].try_recv() { Ok(Some(info)) => { @@ -938,6 +1050,7 @@ fn server_loop( } } } else { + // Inbox arm (mirrors the fast-path park above). match rx.try_recv() { Ok(Some(env)) => { dispatch(&mut guard.0, env);