From 415effb2e9ab16f33273edac0f04c50e43fe812c Mon Sep 17 00:00:00 2001 From: "Claude (sandbox)" Date: Wed, 19 Aug 2026 17:47:31 +0000 Subject: [PATCH] =?UTF-8?q?feat(gen=5Fserver,gen=5Fstatem):=20lifetime=20i?= =?UTF-8?q?s=20the=20actor's=20=E2=80=94=20refs=20are=20addresses;=20inlin?= =?UTF-8?q?e=20named=20run?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Root cause behind the "pin the endpoint" gotcha and the trapping-wrapper pattern: a gen_server had two lifetime authorities — its refs (last one dropped → inbox closes → exit) and, when supervised, its supervisor. OTP has one: a process lives until it stops, is shut down, or is killed; a pid is an address. Root exit now shutting down every forest root removes the reason the ref-governed idiom existed (a forgotten server no longer hangs the run), so adopt the one rule: - The server/machine loop holds one inbox sender for its life; the inbox never closes. GenServerRef / GenStatemRef are addresses. Explicit close is `shutdown()`; a forgotten one is swept at root exit. - `NamedGenServerBuilder::run()` / `gen_statem::run_named(name, m)` run the loop inline as the current actor: a server is a direct ChildSpec child, gets the supervisor's shutdown as handle_shutdown / a shutdown row, re-binds its name on restart, and is addressed by name. The wrapper in examples/graceful_shutdown.rs is gone. - gen_statem gains GenStatemName + whereis_machine/send/call/shutdown by name (parity with gen_server); the macro gets `Sm::new`. - Root-exit sweep records `Event::RootSweep { target, trapping }` under smarm-trace ("root_sweep shutdown|stopped"): unsupervised leftovers are visible rather than silently owned-by-refs. - Named start() name-clash path stops the spawned actor instead of relying on ref drop. Tests: tests/gen_server_lifetime.rs, tests/gen_statem_lifetime.rs, tests/root_sweep_trace.rs (feature-gated); three existing tests that used drop-closes-inbox now use shutdown(). Docs/README/ROADMAP/Deep Dive updated. --- README.md | 8 ++ ROADMAP.md | 13 +- docs/smarm - Deep Dive.html | 7 +- examples/graceful_shutdown.rs | 47 +++---- src/gen_server.rs | 108 +++++++++++---- src/gen_statem.rs | 136 ++++++++++++++++--- src/lib.rs | 2 +- src/runtime.rs | 13 +- src/scheduler.rs | 19 ++- src/trace.rs | 26 +++- tests/gen_server.rs | 12 +- tests/gen_server_lifetime.rs | 243 ++++++++++++++++++++++++++++++++++ tests/gen_server_shutdown.rs | 14 +- tests/gen_statem_lifetime.rs | 128 ++++++++++++++++++ tests/root_sweep_trace.rs | 36 +++++ 15 files changed, 715 insertions(+), 97 deletions(-) create mode 100644 tests/gen_server_lifetime.rs create mode 100644 tests/gen_statem_lifetime.rs create mode 100644 tests/root_sweep_trace.rs diff --git a/README.md b/README.md index 75239be..d75e3fc 100644 --- a/README.md +++ b/README.md @@ -69,6 +69,14 @@ root actor returning means "the program is done": every top-level actor gets a runtime (a signal thread), `Runtime::handle().request_shutdown(pid)` does the same. `examples/graceful_shutdown.rs` shows all of it. +A gen_server or gen_statem lives until it stops, is shut down, or is killed; its +refs are addresses — dropping them never ends it (a forgotten one is swept at +root exit; with `--features smarm-trace` each such sweep is a `root_sweep` +trace line). The supervised shape is `GenServerBuilder::named(N).run()` / +`gen_statem::run_named(N, m)`: the server runs inline as the `ChildSpec` child +itself, so the supervisor's shutdown reaches it directly, a restart re-binds the +name, and the program addresses it by name. + ## Layout ``` diff --git a/ROADMAP.md b/ROADMAP.md index f05ecaa..ebc36b5 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -273,14 +273,11 @@ outright — `join` what you need finished. No forcing sweep follows. --- ### Open items from the graceful-shutdown work (not scheduled) -- **gen_server / gen_statem as a direct supervised child.** A `ChildSpec` start - fn is `Fn()`, and `GenServerBuilder::start` spawns a *new* actor and hands - back a ref, so a supervised server today needs a trapping wrapper actor that - starts it `under(self_pid())`, forwards the shutdown and waits - (`examples/graceful_shutdown.rs::drainer_child`). An inline - `GenServerBuilder::run()` (loop as the current actor; ref handed out via a - name or a start callback) would make the primary OTP use case direct. Needed - by urus's endpoint child. +- ~~gen_server / gen_statem as a direct supervised child.~~ Done: lifetime is + the actor's (refs are addresses, the loop holds an inbox sender); + `NamedGenServerBuilder::run` / `gen_statem::run_named` run the loop inline as + the `ChildSpec` child; the root-exit sweep traces each leftover as + `root_sweep` under `smarm-trace`. - Supervisor `Live` drop-guard sweep is `request_stop` (kill propagates as kill); OTP would deliver a trappable `killed`. Chosen for boundedness. - A root-exit shutdown reaches only actors live *at that instant*; a diff --git a/docs/smarm - Deep Dive.html b/docs/smarm - Deep Dive.html index f772877..e983962 100644 --- a/docs/smarm - Deep Dive.html +++ b/docs/smarm - Deep Dive.html @@ -1620,8 +1620,9 @@ wait: select priority is Down arms › Watcher arm › info channels (declaration order) › inbox, rebuilt each turn. A hot inbox can't starve a death notice or a system message; conversely a hot info channel can starve the inbox — deliberately. A closed info arm is - silently dropped from the set; a closed inbox (every ServerRef gone) is graceful - shutdown.

+ silently dropped from the set. The inbox never closes — the loop holds one sender for its whole + life, so a GenServerRef is an address, not an owner: the server ends only by + StopHandle::stop, a shutdown, a hard stop, or a panic.

Death needs no monitor

Server death detection falls out of channel closure. Already dead → the inbox is closed and @@ -1700,7 +1701,7 @@

Panics in terminate()
-

gen_server's terminate() runs from a drop guard, possibly mid-unwind. A panic inside it during an unwind is a double panic → process abort, no supervision tree to save you. On the panic and hard-stop paths keep it cheap, non-blocking, non-panicking. Only the graceful path (handle_shutdown → Exit, StopHandle::stop, inbox close) runs it outside an unwind, where it may do real work.

+

gen_server's terminate() runs from a drop guard, possibly mid-unwind. A panic inside it during an unwind is a double panic → process abort, no supervision tree to save you. On the panic and hard-stop paths keep it cheap, non-blocking, non-panicking. Only the graceful path (handle_shutdown → Exit, StopHandle::stop) runs it outside an unwind, where it may do real work.

Cold locks are leaf locks
diff --git a/examples/graceful_shutdown.rs b/examples/graceful_shutdown.rs index 9a9031d..e5ec8da 100644 --- a/examples/graceful_shutdown.rs +++ b/examples/graceful_shutdown.rs @@ -24,10 +24,11 @@ //! just waits for the tree to come down. use smarm::gen_server::{ - GenServer, GenServerBuilder, GenServerCtx, ShutdownAction, StopHandle, TimerHandle, + GenServer, GenServerBuilder, GenServerCtx, GenServerName, ShutdownAction, StopHandle, + TimerHandle, }; use smarm::supervisor::{ChildSpec, OneForOne, Restart, Shutdown}; -use smarm::{monitor, request_shutdown, self_pid, sleep, spawn, trap_exit, DownReason}; +use smarm::{sleep, spawn}; use std::thread; use std::time::Duration; @@ -77,28 +78,9 @@ impl GenServer for Drainer { } } -/// A supervised child wrapping the server. (A gen_server is not yet directly -/// usable as a `ChildSpec` start fn; the wrapper traps, forwards the shutdown, -/// and waits for the server to finish. See ROADMAP "open items".) -fn drainer_child() { - let inbox = trap_exit(); - let srv = GenServerBuilder::new(Drainer { - pending: 3, - stop: None, - timer: None, - }) - .under(self_pid()) - .start(); - let mon = monitor(srv.pid()); - // Wait for our shutdown, forward it, wait for the server. - while let Ok(sig) = inbox.recv() { - if sig.reason == DownReason::Shutdown { - request_shutdown(srv.pid()); - let _ = mon.rx.recv(); - return; - } - } -} +/// The server's name: how the rest of the app reaches it (and the only handle +/// that survives a restart). +const DRAINER: GenServerName = GenServerName::new("drainer"); fn app_tree() -> OneForOne { OneForOne::new() @@ -111,7 +93,22 @@ fn app_tree() -> OneForOne { }) .shutdown(Shutdown::Timeout(Duration::from_millis(100))), ) - .child(ChildSpec::new(Restart::Permanent, drainer_child).shutdown(Shutdown::Infinity)) + // A gen_server is a direct child: `named(N).run()` runs the loop as + // the child actor itself, so the supervisor's shutdown arrives as + // `handle_shutdown` and a restart re-binds the name. + .child( + ChildSpec::new(Restart::Permanent, || { + GenServerBuilder::new(Drainer { + pending: 3, + stop: None, + timer: None, + }) + .named(DRAINER) + .run() + .expect("drainer name is free"); + }) + .shutdown(Shutdown::Infinity), + ) } fn main() { diff --git a/src/gen_server.rs b/src/gen_server.rs index 54afb3d..73dd2fb 100644 --- a/src/gen_server.rs +++ b/src/gen_server.rs @@ -127,8 +127,8 @@ //! - [`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, a -//! cooperative stop, or a graceful shutdown), not only on clean shutdown. +//! on every exit path (a self-stop, a graceful shutdown, a cooperative hard +//! stop, or a handler panic), not only on clean shutdown. //! Keep it non-panicking: on the panic and hard-stop paths it runs //! mid-unwind, where a second panic aborts the process and where it must //! not block (any park re-observes the stop). Only on the graceful path @@ -136,13 +136,22 @@ //! //! ## 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 normally. It can also -//! end itself: clone a [`StopHandle`] from [`GenServerCtx::stop_handle`] in -//! `init` and call [`StopHandle::stop`] from any handler — the loop breaks -//! after the current message and exits *normally* (OTP's `{stop, normal}`). -//! This is distinct from `request_stop(self_pid())`, which is an abnormal -//! `Stopped` and gets a `Transient` child restarted. +//! A server lives until it stops, is shut down, or is killed — as an OTP +//! process does. A [`GenServerRef`] is an *address*: cloning and dropping it +//! never changes the server's lifetime, and a ref that nobody holds is not a +//! leak — a forgotten server idles until the run ends, when the root-exit +//! shutdown (see [`Runtime::run`](crate::Runtime::run)) takes it down with +//! every other unsupervised actor. Anything meant to live long should be +//! supervised (see *Supervised servers* below); [`start`] / [`start_under`] +//! are for scripts, tests and short-lived helpers, and the explicit close is +//! [`GenServerRef::shutdown`]. +//! +//! A server can end itself: clone a [`StopHandle`] from +//! [`GenServerCtx::stop_handle`] in `init` and call [`StopHandle::stop`] from +//! any handler — the loop breaks after the current message and exits +//! *normally* (OTP's `{stop, normal}`). This is distinct from +//! `request_stop(self_pid())`, which is an abnormal `Stopped` and gets a +//! `Transient` child restarted. //! //! ## Graceful shutdown //! @@ -193,9 +202,23 @@ //! 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. +//! [`call`], [`cast`], and [`whereis_server`] to address it by name. +//! +//! ## Supervised servers +//! +//! The supervised shape is [`NamedGenServerBuilder::run`]: it runs the loop +//! **inline, as the current actor**, so the closure of a +//! [`ChildSpec`](crate::supervisor::ChildSpec) *is* the server — the +//! supervisor's shutdown arrives as [`GenServer::handle_shutdown`], a restart +//! runs the factory again and re-binds the name, and the rest of the program +//! addresses it by name (a held ref would go stale on restart anyway). +//! +//! ```ignore +//! const COUNTER: GenServerName = GenServerName::new("counter"); +//! OneForOne::new().child(ChildSpec::new(Restart::Permanent, || { +//! GenServerBuilder::new(Counter::default()).named(COUNTER).run().unwrap(); +//! })); +//! ``` //! //! ## Limitations //! @@ -317,9 +340,9 @@ enum Envelope { Cast(G::Cast), } -/// A clonable handle to a running server. Cloning yields another sender to the -/// same inbox; the server lives until the last `GenServerRef` is dropped, at which -/// point its inbox closes and the loop exits normally. +/// A clonable handle to a running server: an *address*, not an owner. Cloning +/// yields another sender to the same inbox; dropping refs never ends the +/// server (see the module docs, *When the server stops*). pub struct GenServerRef { tx: Sender>, pid: Pid, @@ -423,9 +446,8 @@ impl GenServerRef { /// a supervisor uses its child's [`Shutdown`](crate::supervisor::Shutdown) /// policy to bound it. /// - /// 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: + /// This is the explicit close: dropping refs never ends a server. 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) { @@ -807,9 +829,10 @@ impl GenServerBuilder { self } - /// Spawn the server actor and hand back its [`GenServerRef`]. The server's - /// lifetime is governed by its refs, not by joining, so the backing join - /// handle is dropped. + /// Spawn the server actor and hand back its [`GenServerRef`] (an address; + /// the server's lifetime is its own, see the module docs). The backing + /// join handle is dropped. For a supervised server use + /// [`named`](Self::named) + [`NamedGenServerBuilder::run`] instead. pub fn start(self) -> GenServerRef { self.spawn_server() } @@ -837,13 +860,14 @@ impl GenServerBuilder { supervisor, stack_opts, } = self; + let keep = tx.clone(); let handle = match supervisor { Some(sup) => crate::scheduler::spawn_under_with(sup, stack_opts, move || { - server_loop::(rx, state, infos) + server_loop::(keep, rx, state, infos) + }), + None => crate::scheduler::spawn_with(stack_opts, move || { + server_loop::(keep, rx, state, infos) }), - None => { - crate::scheduler::spawn_with(stack_opts, move || server_loop::(rx, state, infos)) - } }; GenServerRef { tx, @@ -931,19 +955,38 @@ impl NamedGenServerBuilder { /// The inbox sender is published under the name **from the parent side, /// before this returns**, so a by-name `call` / `cast` resolves the instant /// `start()` returns — no race with the server body. On a name clash the - /// just-spawned server is wound down (its only ref is dropped, closing the - /// inbox), so a failed bind leaks no actor. + /// just-spawned server is stopped, so a failed bind leaks no actor. pub fn start(self) -> Result, RegisterError> { let NamedGenServerBuilder { builder, name } = self; let server = builder.spawn_server(); match register_with::>(server.pid, name, server.tx.clone()) { Ok(()) => Ok(server), Err(e) => { - drop(server); // inbox closes → loop exits gracefully + crate::scheduler::request_stop(server.pid); // never ran init Err(e) } } } + + /// Run the server **inline, as the current actor**, bound to its name. + /// This is the supervised shape: the closure of a + /// [`ChildSpec`](crate::supervisor::ChildSpec) *is* the server, so the + /// supervisor's shutdown reaches it as [`GenServer::handle_shutdown`], a + /// restart runs the factory again and re-binds the name, and clients + /// address it by name ([`call`], [`cast`], [`whereis_server`]). Returns + /// when the server exits; [`RegisterError::NameTaken`] (before `init`) if + /// the name is held by a different live actor. + /// + /// `under` / `stack_opts` are spawn options and do not apply here — the + /// actor already exists. + pub fn run(self) -> Result<(), RegisterError> { + let NamedGenServerBuilder { builder, name } = self; + let GenServerBuilder { state, infos, .. } = builder; + let (tx, rx) = channel::>(); + register_with::>(crate::scheduler::self_pid(), name, tx.clone())?; + server_loop::(tx, rx, state, infos); + Ok(()) + } } /// Resolve a [`GenServerName`] to a [`GenServerRef`] when you want a handle to hold or @@ -1002,10 +1045,17 @@ pub fn start_under(supervisor: Pid, state: G) -> GenServerRef { } fn server_loop( + keep: Sender>, rx: Receiver>, state: G, mut infos: Vec>, ) { + // The loop holds one inbox sender for its whole life: the inbox never + // closes, so refs are addresses and the server's lifetime is the actor's + // (stop handle, shutdown, stop, panic). The `Disconnected` arms below are + // defensive only. + let _keep = keep; + // Drop guard — owns the server state and the timer registry. // // Why a guard rather than code after the loop: @@ -1117,7 +1167,7 @@ fn server_loop( guard.0.handle_idle(); reset_idle(&mut idle_deadline); } - // All ServerRefs dropped → inbox closed → shutdown. + // Defensive: the loop holds a sender, so unreachable. Err(RecvTimeoutError::Disconnected) => break, } } diff --git a/src/gen_statem.rs b/src/gen_statem.rs index 7fdf96a..357a940 100644 --- a/src/gen_statem.rs +++ b/src/gen_statem.rs @@ -71,8 +71,12 @@ //! //! ## Stopping, shutdown, and exits //! -//! A machine runs while a [`GenStatemRef`] exists; dropping the last one -//! closes the inbox and ends it normally. It can also end itself: any row +//! A machine lives until it stops, is shut down, or is killed; a +//! [`GenStatemRef`] is an address, and dropping refs never ends it (the +//! gen_server rule — see its *When the server stops*). The supervised shape +//! is [`run_named`], which runs the machine inline as the current actor so it +//! is a direct `ChildSpec` child addressed by [`GenStatemName`]. A machine +//! can end itself: any row //! body may call [`cx.stop()`](Cx::stop) (or use the `stop` tail keyword, //! sugar for `{ cx.stop(); prev }`) — the loop breaks after that event and the //! actor exits *normally* (OTP's `{stop, normal}`; a `Transient` child is not @@ -94,6 +98,7 @@ use crate::channel::{channel, select, Receiver, Selectable, Sender}; use crate::link::ExitSignal; use crate::monitor::{monitor, DownReason}; use crate::pid::Pid; +use crate::registry::{register_with, resolve_named_sender, RegisterError}; use crate::scheduler::{cancel_timer, request_shutdown, send_after_to}; use crate::timer::TimerId; use std::collections::{HashMap, VecDeque}; @@ -162,11 +167,12 @@ pub trait Machine: Send + 'static { None } - /// Runs as the machine actor exits, on any exit path (inbox closed, a - /// `stop`, a handler panic, a hard stop). Like `gen_server::terminate`: - /// on the panic and hard-stop paths it runs mid-unwind — do not panic or - /// park there; only on the normal path (`stop`, inbox close) may it do real - /// work. The macro's optional `terminate { … }` block generates it. + /// Runs as the machine actor exits, on any exit path (a `stop`, a + /// graceful shutdown, a handler panic, a hard stop). Like + /// `gen_server::terminate`: on the panic and hard-stop paths it runs + /// mid-unwind — do not panic or park there; only on the normal path + /// (`stop`, shutdown rows) may it do real work. The macro's optional + /// `terminate { … }` block generates it. fn terminate(&mut self) {} } @@ -460,8 +466,7 @@ pub enum SendError { } /// A clonable handle to a running machine. Cloning yields another sender to the -/// same inbox; the machine lives until the last `GenStatemRef` is dropped, at which -/// point its inbox closes and the loop exits. +/// same inbox. An address, not an owner: dropping refs never ends the machine. pub struct GenStatemRef { tx: Sender, pid: Pid, @@ -528,7 +533,8 @@ impl GenStatemRef { /// Spawn `machine` as an actor and hand back its [`GenStatemRef`]. Shape mirrors /// `gen_server::start`: make the inbox, spawn the loop, return the ref; the -/// backing join handle is dropped (lifetime is governed by refs, not joining). +/// backing join handle is dropped (the machine's lifetime is its own). For a +/// supervised machine use [`run_named`]. /// /// Panics if called outside `Runtime::run()`. pub fn spawn(machine: M) -> GenStatemRef { @@ -543,16 +549,104 @@ pub fn spawn(machine: M) -> GenStatemRef { /// Panics if called outside `Runtime::run()`. pub fn spawn_with(opts: crate::scheduler::SpawnOpts, machine: M) -> GenStatemRef { let (tx, rx) = channel::(); - let handle = crate::scheduler::spawn_with(opts, move || statem_loop(rx, machine)); + let keep = tx.clone(); + let handle = crate::scheduler::spawn_with(opts, move || statem_loop(keep, rx, machine)); GenStatemRef { tx, pid: handle.pid(), } } +/// A typed, static name for a gen_statem, used to address a machine through +/// the registry without holding a [`GenStatemRef`]. Mirrors +/// [`GenServerName`](crate::gen_server::GenServerName): declare it as a +/// constant and bind it with [`run_named`]. +pub struct GenStatemName { + name: &'static str, + _marker: PhantomData M>, +} + +impl GenStatemName { + /// Bind a static string as a machine name. + #[inline] + pub const fn new(name: &'static str) -> Self { + Self { + name, + _marker: PhantomData, + } + } + + /// The underlying registry key. + #[inline] + pub const fn as_str(self) -> &'static str { + self.name + } +} + +impl Copy for GenStatemName {} +impl Clone for GenStatemName { + fn clone(&self) -> Self { + *self + } +} + +/// Run `machine` **inline, as the current actor**, bound to `name`. The +/// supervised shape: the closure of a +/// [`ChildSpec`](crate::supervisor::ChildSpec) *is* the machine, so the +/// supervisor's shutdown reaches it as a `shutdown` row, a restart runs the +/// factory again and re-binds the name, and clients address it by name +/// ([`send`], [`call`], [`whereis_machine`]). Returns when the machine exits; +/// [`RegisterError::NameTaken`] (before `on_start`) if the name is held by a +/// different live actor. Mirrors +/// [`NamedGenServerBuilder::run`](crate::gen_server::NamedGenServerBuilder::run). +pub fn run_named(name: GenStatemName, machine: M) -> Result<(), RegisterError> { + let (tx, rx) = channel::(); + register_with::(crate::scheduler::self_pid(), name.as_str(), tx.clone())?; + statem_loop(tx, rx, machine); + Ok(()) +} + +/// Resolve a [`GenStatemName`] to a [`GenStatemRef`]; `None` if no live +/// machine holds the name. Panics if called outside `Runtime::run()`. +pub fn whereis_machine(name: GenStatemName) -> Option> { + resolve_named_sender::(name.as_str()).map(|(pid, tx)| GenStatemRef { tx, pid }) +} + +/// Push an event to the machine registered under `name`, resolving per send. +/// [`SendError::Down`] if no live machine holds the name. +pub fn send(name: GenStatemName, ev: M::Ev) -> Result<(), SendError> { + match whereis_machine(name) { + Some(m) => m.send(ev), + None => Err(SendError::Down), + } +} + +/// Synchronous request-reply to the machine registered under `name`, +/// resolving per call (a machine restarted under the same name is reached +/// transparently). [`CallError::Down`] if no live machine holds the name. +pub fn call(name: GenStatemName, make: F) -> Result +where + M: Machine, + T: Send + 'static, + F: FnOnce(Reply) -> M::Ev, +{ + match whereis_machine(name) { + Some(m) => m.call(make), + None => Err(CallError::Down), + } +} + +/// Shut down the machine registered under `name` and wait for it (see +/// [`GenStatemRef::shutdown`]). A no-op if no live machine holds the name. +pub fn shutdown(name: GenStatemName) { + if let Some(m) = whereis_machine(name) { + m.shutdown(); + } +} + /// The machine actor body: `on_start`, then one `handle` per event until the -/// inbox closes (all refs dropped), a row resolves to `stop`, or the actor is -/// stopped from outside. +/// row resolves to `stop`, a shutdown row stops it, or the actor is stopped +/// from outside. /// /// Intake arms are selected each iteration in priority order — **exits** /// (only when trapping) above **timers** above the **inbox** — so a shutdown @@ -566,7 +660,11 @@ pub fn spawn_with(opts: crate::scheduler::SpawnOpts, machine: M) -> /// it back ([`Step::Postponed`]) for the queue; a `handle` that transitions /// ([`Step::Transitioned`]) triggers a [`replay`] of the queue in the new /// state, ahead of the next intake. -fn statem_loop(rx: Receiver, machine: M) { +fn statem_loop(keep: Sender, rx: Receiver, machine: M) { + // One inbox sender lives with the loop: the inbox never closes, refs are + // addresses, the machine's lifetime is the actor's (stop row, shutdown, + // stop, panic). The `Disconnected` inbox arm below is defensive only. + let _keep = keep; // Drop guard — owns the machine and the timer registry, so `terminate` // fires on every exit path (clean close, `stop`, a handler panic, a hard // stop) and the timer drain is sequenced before it. Same shape as @@ -685,7 +783,7 @@ fn statem_loop(rx: Receiver, machine: M) { match rx.try_recv() { Ok(Some(ev)) => dispatch(&mut guard.0, &mut cx, &mut postpone, ev), Ok(None) => debug_assert!(false, "ready inbox was empty"), - // All GenStatemRefs dropped → inbox closed → normal exit. + // Defensive: the loop holds a sender, so unreachable. Err(_) => break, } } @@ -991,8 +1089,14 @@ macro_rules! gen_statem { } impl $sm { + /// The machine value, for [`gen_statem::run_named`] + /// (`$crate::gen_statem::run_named`) or `spawn`. + fn new(init: $State, data: $Data) -> $sm { + $sm { state: init, data } + } + fn start(init: $State, data: $Data) -> $crate::gen_statem::GenStatemRef<$sm> { - $crate::gen_statem::spawn($sm { state: init, data }) + $crate::gen_statem::spawn($sm::new(init, data)) } #[allow(unused_variables)] diff --git a/src/lib.rs b/src/lib.rs index c27d0b5..c7c9b90 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -63,7 +63,7 @@ pub use gen_server::{ ShutdownAction, StopHandle, TimerHandle, Watcher, }; pub use gen_statem::{ - CallError as GenStatemCallError, Cx, GenStatemRef, Machine, Reply, Resolution, + CallError as GenStatemCallError, Cx, GenStatemName, GenStatemRef, Machine, Reply, Resolution, SendError as GenStatemSendError, }; pub use introspect::{ diff --git a/src/runtime.rs b/src/runtime.rs index 16f45a3..da033d4 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -1965,7 +1965,18 @@ fn shutdown_forest_roots(inner: &Arc, root: Pid) { .slot_at(parent) .is_some_and(|ps| ps.is_live_for(parent)); if !parent_live { - crate::scheduler::request_shutdown_inner(inner, pid, root); + // `_probe` so the trace can say what each leftover was: under + // `smarm-trace` every swept actor is a `root_sweep` line — the + // visibility that makes "a forgotten actor costs only a slot, never + // a hung run" a checkable claim rather than a hope. + let found = crate::scheduler::request_shutdown_inner_probe(inner, pid, root); + #[cfg_attr(not(feature = "smarm-trace"), allow(unused_variables))] + if let Some(trapping) = found { + crate::te!(crate::trace::Event::RootSweep { + target: pid, + trapping + }); + } } } } diff --git a/src/scheduler.rs b/src/scheduler.rs index 00d9125..ecd10ae 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -662,6 +662,17 @@ pub fn request_shutdown(pid: Pid) { // cold lock (generation-verified), then acts outside the lock: a trap send // may unpark the receiver, and `request_stop_inner` re-takes the lock. pub(crate) fn request_shutdown_inner(inner: &RuntimeInner, pid: Pid, from: Pid) { + request_shutdown_inner_probe(inner, pid, from); +} + +/// [`request_shutdown_inner`], reporting what it found: `Some(true)` if the +/// target was trapping (got the signal), `Some(false)` if it was stopped +/// outright, `None` if there was nothing live at `pid`. +pub(crate) fn request_shutdown_inner_probe( + inner: &RuntimeInner, + pid: Pid, + from: Pid, +) -> Option { let trap = match inner.slot_at(pid) { Some(slot) => { let cold = slot.cold.lock(); @@ -679,9 +690,13 @@ pub(crate) fn request_shutdown_inner(inner: &RuntimeInner, pid: Pid, from: Pid) from, reason: crate::monitor::DownReason::Shutdown, }); + Some(true) } - Some(None) => request_stop_inner(inner, pid), - None => {} + Some(None) => { + request_stop_inner(inner, pid); + Some(false) + } + None => None, } } diff --git a/src/trace.rs b/src/trace.rs index 923cbbd..fbe4b5c 100644 --- a/src/trace.rs +++ b/src/trace.rs @@ -46,17 +46,32 @@ mod inner { #[derive(Clone, Debug)] pub enum Event { // Actor lifecycle - Spawn { parent: Pid, child: Pid }, + Spawn { + parent: Pid, + child: Pid, + }, Resume(Pid), Yield(Pid), Park(Pid), Done(Pid), + /// Root exit found a live forest root (an actor nobody supervises) + /// and delivered `request_shutdown` to it. `trapping` says whether it + /// got the chance to drain (true) or was stopped outright (false). + /// Every such line is an actor whose lifetime was nobody's business + /// but the runtime's — the way to *see* unsupervised leftovers. + RootSweep { + target: Pid, + trapping: bool, + }, // Wakeup paths UnparkDirect(Pid), // unpark() saw Parked -> re-queued immediately UnparkDeferred(Pid), // unpark() saw Runnable -> set pending_unpark flag UnparkFlagConsumed(Pid), // scheduler saw flag on Park -> re-queued instead // Channel - Send { sender: Pid, receiver: Option }, + Send { + sender: Pid, + receiver: Option, + }, RecvPark(Pid), RecvWake(Pid), // Queue @@ -253,6 +268,13 @@ mod inner { Event::Yield(p) => ("yield".into(), p.index()), Event::Park(p) => ("park".into(), p.index()), Event::Done(p) => ("done".into(), p.index()), + Event::RootSweep { target, trapping } => ( + format!( + "root_sweep {}", + if *trapping { "shutdown" } else { "stopped" } + ), + target.index(), + ), Event::UnparkDirect(p) => ("unpark_direct".into(), p.index()), Event::UnparkDeferred(p) => ("unpark_deferred".into(), p.index()), Event::UnparkFlagConsumed(p) => ("unpark_flag_consumed".into(), p.index()), diff --git a/tests/gen_server.rs b/tests/gen_server.rs index 694a2b6..33292c1 100644 --- a/tests/gen_server.rs +++ b/tests/gen_server.rs @@ -88,7 +88,7 @@ impl GenServer for Lifecycle { } } -// init -> handle_call -> (drop last ref closes inbox) -> terminate. +// init -> handle_call -> shutdown -> terminate. #[test] fn init_and_terminate_run() { let log = Arc::new(Mutex::new(Vec::new())); @@ -96,9 +96,9 @@ fn init_and_terminate_run() { run(move || { let server = start(Lifecycle { log: log2 }); server.call(()).unwrap(); - // Dropping the only ref closes the inbox; the server breaks out of its - // recv loop and runs terminate. run() will not return until it has. - drop(server); + // Refs are addresses: dropping one does not end the server. The + // explicit close does, and waits for terminate. + server.shutdown(); }); assert_eq!(*log.lock().unwrap(), vec!["init", "call", "terminate"]); } @@ -702,8 +702,8 @@ fn no_timer_survives_exit() { let _ = server.call(()).unwrap(); // sync: periodic armed smarm::sleep(Duration::from_millis(45)); // a couple of ticks let mon = smarm::monitor(server.pid()); - drop(server); // inbox closes → loop exits → guard drains timers - // Clean Down ⇒ the loop returned without the no-leak assert aborting. + server.shutdown(); // loop exits → guard drains timers + // Clean Down ⇒ the loop returned without the no-leak assert aborting. assert!(mon.rx.recv().is_ok()); let at_exit = f_read.lock().unwrap().len(); smarm::sleep(Duration::from_millis(90)); // would be several more ticks diff --git a/tests/gen_server_lifetime.rs b/tests/gen_server_lifetime.rs new file mode 100644 index 0000000..b5ab393 --- /dev/null +++ b/tests/gen_server_lifetime.rs @@ -0,0 +1,243 @@ +//! gen_server lifetime is the actor's, not its refs' (OTP: a pid is an +//! address, a process lives until it stops, is shut down, or is killed). +//! +//! - Dropping the last `GenServerRef` does NOT end the server. It ends via +//! `StopHandle::stop`, `request_shutdown` / `GenServerRef::shutdown`, +//! `request_stop`, or a handler panic. +//! - `GenServerBuilder::named(N).run()` runs the loop inline as the *current* +//! actor, so a server is a direct `ChildSpec` child: the supervisor's +//! shutdown reaches it as `handle_shutdown`, a restart re-binds the name, +//! and by-name `call`/`cast` reach whichever incarnation is live. + +use smarm::gen_server::{ + self, GenServer, GenServerBuilder, GenServerCtx, GenServerName, ShutdownAction, StopHandle, +}; +use smarm::registry::RegisterError; +use smarm::supervisor::{ChildSpec, OneForOne, Restart, Shutdown}; +use smarm::{monitor, request_shutdown, run, sleep, spawn, DownReason}; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +#[derive(Default, Clone)] +struct Log(Arc>>); +impl Log { + fn push(&self, s: impl Into) { + self.0.lock().unwrap().push(s.into()); + } + fn get(&self) -> Vec { + self.0.lock().unwrap().clone() + } +} + +struct Counter { + log: Log, + n: u64, + trap: bool, + stop: Option>, +} + +enum Call { + Get, +} +enum Cast { + Inc, + Stop, +} + +impl GenServer for Counter { + type Call = Call; + type Reply = u64; + type Cast = Cast; + type Info = (); + type Timer = (); + + fn init(&mut self, ctx: &GenServerCtx) { + if self.trap { + ctx.trap_exit(); + } + self.stop = Some(ctx.stop_handle()); + self.log.push("init"); + } + fn handle_call(&mut self, Call::Get: Call) -> u64 { + self.n + } + fn handle_cast(&mut self, c: Cast) { + match c { + Cast::Inc => self.n += 1, + Cast::Stop => self.stop.as_ref().unwrap().stop(), + } + } + fn handle_shutdown(&mut self) -> ShutdownAction { + self.log.push("handle_shutdown"); + ShutdownAction::Exit + } + fn terminate(&mut self) { + self.log.push("terminate"); + } +} + +fn counter(log: &Log, trap: bool) -> Counter { + Counter { + log: log.clone(), + n: 0, + trap, + stop: None, + } +} + +// --------------------------------------------------------------------------- +// Refs are addresses: dropping the last one does not end the server. +// --------------------------------------------------------------------------- + +#[test] +fn dropping_last_ref_does_not_end_server() { + let log = Log::default(); + let l = log.clone(); + run(move || { + let srv = gen_server::start(counter(&l, true)); + let pid = srv.pid(); + srv.cast(Cast::Inc).unwrap(); + assert_eq!(srv.call(Call::Get).unwrap(), 1); + let mon = monitor(pid); + drop(srv); + sleep(Duration::from_millis(30)); + assert!( + mon.rx.try_recv().unwrap().is_none(), + "server must outlive its last ref" + ); + assert_eq!(l.get(), vec!["init"], "terminate must not have run"); + // Explicit teardown still works, and is what ends it. + request_shutdown(pid); + let down = mon.rx.recv().unwrap(); + assert_eq!(down.reason, DownReason::Exit); + }); + assert_eq!(log.get(), vec!["init", "handle_shutdown", "terminate"]); +} + +#[test] +fn ref_shutdown_is_the_explicit_close() { + let log = Log::default(); + let l = log.clone(); + run(move || { + let srv = gen_server::start(counter(&l, true)); + srv.call(Call::Get).unwrap(); // sync: init (and trap_exit) has run + srv.shutdown(); // graceful, waits + assert_eq!(l.get(), vec!["init", "handle_shutdown", "terminate"]); + }); +} + +#[test] +fn forgotten_server_is_shut_down_at_root_exit() { + // A ref-less server is not a hung run: root exit shuts it down. + let log = Log::default(); + let l = log.clone(); + run(move || { + let srv = gen_server::start(counter(&l, false)); + drop(srv); + }); + assert_eq!(log.get(), vec!["init", "terminate"]); +} + +// --------------------------------------------------------------------------- +// Inline run: a gen_server as a direct ChildSpec child. +// --------------------------------------------------------------------------- + +const COUNTER: GenServerName = GenServerName::new("lifetime-counter"); + +#[test] +fn named_run_is_a_direct_supervised_child_and_gets_shutdown() { + let log = Log::default(); + let l = log.clone(); + run(move || { + let l2 = l.clone(); + let sup = spawn(move || { + let l3 = l2.clone(); + OneForOne::new() + .child( + ChildSpec::new(Restart::Permanent, move || { + GenServerBuilder::new(counter(&l3, true)) + .named(COUNTER) + .run() + .expect("name free"); + }) + .shutdown(Shutdown::Infinity), + ) + .run(); + }); + sleep(Duration::from_millis(10)); + gen_server::cast(COUNTER, Cast::Inc).unwrap(); + assert_eq!(gen_server::call(COUNTER, Call::Get).unwrap(), 1); + request_shutdown(sup.pid()); + sup.join() + .expect("ordered shutdown, supervisor returns normally"); + assert_eq!(l.get(), vec!["init", "handle_shutdown", "terminate"]); + assert!(gen_server::whereis_server(COUNTER).is_none()); + }); +} + +#[test] +fn named_run_child_restarts_and_rebinds_name() { + let log = Log::default(); + let inits = Arc::new(AtomicUsize::new(0)); + let l = log.clone(); + let i = inits.clone(); + run(move || { + let l2 = l.clone(); + let i2 = i.clone(); + let sup = spawn(move || { + let l3 = l2.clone(); + let i3 = i2.clone(); + OneForOne::new() + .child(ChildSpec::new(Restart::Permanent, move || { + i3.fetch_add(1, Ordering::SeqCst); + GenServerBuilder::new(counter(&l3, false)) + .named(COUNTER) + .run() + .expect("name free on (re)start"); + })) + .run(); + }); + sleep(Duration::from_millis(10)); + gen_server::cast(COUNTER, Cast::Inc).unwrap(); + assert_eq!(gen_server::call(COUNTER, Call::Get).unwrap(), 1); + // Normal self-exit → Permanent restarts it, fresh state, same name. + gen_server::cast(COUNTER, Cast::Stop).unwrap(); + sleep(Duration::from_millis(30)); + assert_eq!(i.load(Ordering::SeqCst), 2, "restarted once"); + assert_eq!(gen_server::call(COUNTER, Call::Get).unwrap(), 0); + request_shutdown(sup.pid()); + sup.join().unwrap(); + }); + assert_eq!(log.get(), vec!["init", "terminate", "init", "terminate"]); +} + +#[test] +fn named_run_name_clash_fails_before_init() { + let log = Log::default(); + let l = log.clone(); + run(move || { + let first = GenServerBuilder::new(counter(&l, false)) + .named(COUNTER) + .start() + .unwrap(); + let l2 = l.clone(); + let res = Arc::new(Mutex::new(None)); + let r2 = res.clone(); + let first_pid = first.pid(); + spawn(move || { + let r = GenServerBuilder::new(counter(&l2, false)) + .named(COUNTER) + .run(); + *r2.lock().unwrap() = Some(r); + }) + .join() + .unwrap(); + assert_eq!( + *res.lock().unwrap(), + Some(Err(RegisterError::NameTaken { holder: first_pid })) + ); + assert_eq!(l.get(), vec!["init"], "clashing server never ran init"); + first.shutdown(); + }); +} diff --git a/tests/gen_server_shutdown.rs b/tests/gen_server_shutdown.rs index 9491e49..c7421e6 100644 --- a/tests/gen_server_shutdown.rs +++ b/tests/gen_server_shutdown.rs @@ -196,12 +196,18 @@ fn linked_peer_death_reaches_handle_exit() { r.cast(Cast::Note("still-serving")).unwrap(); sleep(Duration::from_millis(20)); a.store(true, Ordering::SeqCst); - let mon = monitor(r.pid()); - drop(r); // inbox closes → clean exit - let _ = mon.rx.recv(); // don't let the root-exit sweep race terminate + r.shutdown(); // graceful; waits for terminate }); assert!(alive.load(Ordering::SeqCst)); - assert_eq!(log.get(), vec!["handle_exit", "still-serving", "terminate"]); + assert_eq!( + log.get(), + vec![ + "handle_exit", + "still-serving", + "handle_shutdown", + "terminate" + ] + ); } #[test] diff --git a/tests/gen_statem_lifetime.rs b/tests/gen_statem_lifetime.rs new file mode 100644 index 0000000..426801c --- /dev/null +++ b/tests/gen_statem_lifetime.rs @@ -0,0 +1,128 @@ +//! gen_statem lifetime parity with gen_server: a machine lives until it +//! stops, is shut down, or is killed — its refs are addresses. And +//! `gen_statem::run_named` runs a machine inline as the current actor, so it +//! is a direct `ChildSpec` child addressed by name. + +use smarm::gen_statem::{self, GenStatemName, Reply}; +use smarm::supervisor::{ChildSpec, OneForOne, Restart, Shutdown}; +use smarm::{monitor, request_shutdown, run, sleep, spawn, DownReason}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +#[derive(Default, Clone)] +struct Log(Arc>>); +impl Log { + fn push(&self, e: &'static str) { + self.0.lock().unwrap().push(e); + } + fn get(&self) -> Vec<&'static str> { + self.0.lock().unwrap().clone() + } +} + +#[derive(Clone, Copy, PartialEq, Eq, Debug)] +enum S { + On, +} + +struct D { + log: Log, + trap: bool, + n: u64, +} + +enum Cast { + Inc, + StopNow, +} +enum Call { + Get(Reply), +} + +smarm::gen_statem! { + machine: Sm { state: S, data: D }; + event: Ev { cast: Cast, call: Call, info: () }; + context(data, prev, cx); + + enter { + S::On => { data.log.push("enter"); if data.trap { cx.trap_exit() } }, + } + + on S::On => { + cast Cast::Inc => { data.n += 1; prev }, + cast Cast::StopNow => stop, + call Call::Get(r) => { r.reply(data.n); prev }, + shutdown => { data.log.push("shutdown"); cx.stop(); prev }, + state_timeout => unhandled, + timeout _ => unhandled, + } + + terminate { data.log.push("terminate"); } +} + +fn d(log: &Log, trap: bool) -> D { + D { + log: log.clone(), + trap, + n: 0, + } +} + +#[test] +fn dropping_last_ref_does_not_end_machine() { + let log = Log::default(); + let l = log.clone(); + run(move || { + let m = Sm::start(S::On, d(&l, true)); + let pid = m.pid(); + m.send(Ev::Cast(Cast::Inc)).unwrap(); + assert_eq!(m.call(|r| Ev::Call(Call::Get(r))).unwrap(), 1); + let mon = monitor(pid); + drop(m); + sleep(Duration::from_millis(30)); + assert!( + mon.rx.try_recv().unwrap().is_none(), + "machine must outlive its refs" + ); + assert_eq!(l.get(), vec!["enter"]); + request_shutdown(pid); + assert_eq!(mon.rx.recv().unwrap().reason, DownReason::Exit); + }); + assert_eq!(log.get(), vec!["enter", "shutdown", "terminate"]); +} + +const SM: GenStatemName = GenStatemName::new("lifetime-sm"); + +#[test] +fn run_named_is_a_direct_supervised_child() { + let log = Log::default(); + let l = log.clone(); + run(move || { + let l2 = l.clone(); + let sup = spawn(move || { + let l3 = l2.clone(); + OneForOne::new() + .child( + ChildSpec::new(Restart::Permanent, move || { + gen_statem::run_named(SM, Sm::new(S::On, d(&l3, true))).expect("name free"); + }) + .shutdown(Shutdown::Infinity), + ) + .run(); + }); + sleep(Duration::from_millis(10)); + gen_statem::send(SM, Ev::Cast(Cast::Inc)).unwrap(); + assert_eq!(gen_statem::call(SM, |r| Ev::Call(Call::Get(r))).unwrap(), 1); + // Normal self-exit → Permanent restart → fresh data, same name. + gen_statem::send(SM, Ev::Cast(Cast::StopNow)).unwrap(); + sleep(Duration::from_millis(30)); + assert_eq!(gen_statem::call(SM, |r| Ev::Call(Call::Get(r))).unwrap(), 0); + request_shutdown(sup.pid()); + sup.join().unwrap(); + assert!(gen_statem::whereis_machine(SM).is_none()); + }); + assert_eq!( + log.get(), + vec!["enter", "terminate", "enter", "shutdown", "terminate"] + ); +} diff --git a/tests/root_sweep_trace.rs b/tests/root_sweep_trace.rs new file mode 100644 index 0000000..7514ea4 --- /dev/null +++ b/tests/root_sweep_trace.rs @@ -0,0 +1,36 @@ +//! Under `smarm-trace`, every actor the root-exit sweep reaches is recorded +//! as a `root_sweep` event — the way to *see* unsupervised leftovers. One test +//! per binary: the trace file is process-global. +#![cfg(feature = "smarm-trace")] + +use smarm::gen_server::{self, GenServer, GenServerCtx}; +use smarm::run; + +struct Quiet; +impl GenServer for Quiet { + type Call = (); + type Reply = (); + type Cast = (); + type Info = (); + type Timer = (); + fn init(&mut self, _: &GenServerCtx) {} + fn handle_call(&mut self, _: ()) {} + fn handle_cast(&mut self, _: ()) {} +} + +#[test] +fn forgotten_server_shows_up_as_root_sweep() { + let path = std::env::temp_dir().join(format!("smarm_root_sweep_{}.json", std::process::id())); + std::env::set_var("SMARM_TRACE_FILE", &path); + run(|| { + let srv = gen_server::start(Quiet); + srv.call(()).unwrap(); + drop(srv); // forgotten: nobody supervises it, nobody holds it + }); + let trace = std::fs::read_to_string(&path).expect("trace file written"); + let _ = std::fs::remove_file(&path); + assert!( + trace.contains("root_sweep stopped"), + "expected a root_sweep line for the non-trapping leftover; got:\n{trace}" + ); +}