feat(gen_server,gen_statem): lifetime is the actor's — refs are addresses; inline named run
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.
This commit is contained in:
@@ -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
|
runtime (a signal thread), `Runtime::handle().request_shutdown(pid)` does the
|
||||||
same. `examples/graceful_shutdown.rs` shows all of it.
|
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
|
## Layout
|
||||||
|
|
||||||
```
|
```
|
||||||
|
|||||||
+5
-8
@@ -273,14 +273,11 @@ outright — `join` what you need finished. No forcing sweep follows.
|
|||||||
---
|
---
|
||||||
|
|
||||||
### Open items from the graceful-shutdown work (not scheduled)
|
### Open items from the graceful-shutdown work (not scheduled)
|
||||||
- **gen_server / gen_statem as a direct supervised child.** A `ChildSpec` start
|
- ~~gen_server / gen_statem as a direct supervised child.~~ Done: lifetime is
|
||||||
fn is `Fn()`, and `GenServerBuilder::start` spawns a *new* actor and hands
|
the actor's (refs are addresses, the loop holds an inbox sender);
|
||||||
back a ref, so a supervised server today needs a trapping wrapper actor that
|
`NamedGenServerBuilder::run` / `gen_statem::run_named` run the loop inline as
|
||||||
starts it `under(self_pid())`, forwards the shutdown and waits
|
the `ChildSpec` child; the root-exit sweep traces each leftover as
|
||||||
(`examples/graceful_shutdown.rs::drainer_child`). An inline
|
`root_sweep` under `smarm-trace`.
|
||||||
`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.
|
|
||||||
- Supervisor `Live` drop-guard sweep is `request_stop` (kill propagates as
|
- Supervisor `Live` drop-guard sweep is `request_stop` (kill propagates as
|
||||||
kill); OTP would deliver a trappable `killed`. Chosen for boundedness.
|
kill); OTP would deliver a trappable `killed`. Chosen for boundedness.
|
||||||
- A root-exit shutdown reaches only actors live *at that instant*; a
|
- A root-exit shutdown reaches only actors live *at that instant*; a
|
||||||
|
|||||||
@@ -1620,8 +1620,9 @@
|
|||||||
wait: <code>select</code> priority is <strong>Down arms › Watcher arm › info channels (declaration order) ›
|
wait: <code>select</code> priority is <strong>Down arms › Watcher arm › info channels (declaration order) ›
|
||||||
inbox</strong>, rebuilt each turn. A hot inbox can't starve a death notice or a system message;
|
inbox</strong>, rebuilt each turn. A hot inbox can't starve a death notice or a system message;
|
||||||
conversely a hot info channel <em>can</em> starve the inbox — deliberately. A closed info arm is
|
conversely a hot info channel <em>can</em> starve the inbox — deliberately. A closed info arm is
|
||||||
silently dropped from the set; a closed <em>inbox</em> (every <code>ServerRef</code> gone) is graceful
|
silently dropped from the set. The inbox never closes — the loop holds one sender for its whole
|
||||||
shutdown.</p>
|
life, so a <code>GenServerRef</code> is an address, not an owner: the server ends only by
|
||||||
|
<code>StopHandle::stop</code>, a shutdown, a hard stop, or a panic.</p>
|
||||||
|
|
||||||
<h3>Death needs no monitor</h3>
|
<h3>Death needs no monitor</h3>
|
||||||
<p>Server death detection falls out of channel closure. Already dead → the inbox is closed and
|
<p>Server death detection falls out of channel closure. Already dead → the inbox is closed and
|
||||||
@@ -1700,7 +1701,7 @@
|
|||||||
</div>
|
</div>
|
||||||
<div class="module-card">
|
<div class="module-card">
|
||||||
<div class="module-name" style="color:var(--red)">Panics in <code>terminate()</code></div>
|
<div class="module-name" style="color:var(--red)">Panics in <code>terminate()</code></div>
|
||||||
<p>gen_server's <code>terminate()</code> 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 (<code>handle_shutdown → Exit</code>, <code>StopHandle::stop</code>, inbox close) runs it outside an unwind, where it may do real work.</p>
|
<p>gen_server's <code>terminate()</code> 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 (<code>handle_shutdown → Exit</code>, <code>StopHandle::stop</code>) runs it outside an unwind, where it may do real work.</p>
|
||||||
</div>
|
</div>
|
||||||
<div class="module-card">
|
<div class="module-card">
|
||||||
<div class="module-name" style="color:var(--yellow)">Cold locks are leaf locks</div>
|
<div class="module-name" style="color:var(--yellow)">Cold locks are leaf locks</div>
|
||||||
|
|||||||
@@ -24,10 +24,11 @@
|
|||||||
//! just waits for the tree to come down.
|
//! just waits for the tree to come down.
|
||||||
|
|
||||||
use smarm::gen_server::{
|
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::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::thread;
|
||||||
use std::time::Duration;
|
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
|
/// The server's name: how the rest of the app reaches it (and the only handle
|
||||||
/// usable as a `ChildSpec` start fn; the wrapper traps, forwards the shutdown,
|
/// that survives a restart).
|
||||||
/// and waits for the server to finish. See ROADMAP "open items".)
|
const DRAINER: GenServerName<Drainer> = GenServerName::new("drainer");
|
||||||
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;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn app_tree() -> OneForOne {
|
fn app_tree() -> OneForOne {
|
||||||
OneForOne::new()
|
OneForOne::new()
|
||||||
@@ -111,7 +93,22 @@ fn app_tree() -> OneForOne {
|
|||||||
})
|
})
|
||||||
.shutdown(Shutdown::Timeout(Duration::from_millis(100))),
|
.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() {
|
fn main() {
|
||||||
|
|||||||
+79
-29
@@ -127,8 +127,8 @@
|
|||||||
//! - [`GenServer::init`] runs once before the first message. Use it to start
|
//! - [`GenServer::init`] runs once before the first message. Use it to start
|
||||||
//! timers or set up monitors; see the [`GenServerCtx`] it receives.
|
//! timers or set up monitors; see the [`GenServerCtx`] it receives.
|
||||||
//! - [`GenServer::terminate`] runs when the server is about to exit. It fires
|
//! - [`GenServer::terminate`] runs when the server is about to exit. It fires
|
||||||
//! on every exit path (all `GenServerRef`s dropped, a handler panic, a
|
//! on every exit path (a self-stop, a graceful shutdown, a cooperative hard
|
||||||
//! cooperative stop, or a graceful shutdown), not only on clean shutdown.
|
//! stop, or a handler panic), not only on clean shutdown.
|
||||||
//! Keep it non-panicking: on the panic and hard-stop paths it runs
|
//! 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
|
//! 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
|
//! not block (any park re-observes the stop). Only on the graceful path
|
||||||
@@ -136,13 +136,22 @@
|
|||||||
//!
|
//!
|
||||||
//! ## When the server stops
|
//! ## When the server stops
|
||||||
//!
|
//!
|
||||||
//! The server runs as long as at least one [`GenServerRef`] exists. When the last
|
//! A server lives until it stops, is shut down, or is killed — as an OTP
|
||||||
//! one is dropped, the inbox closes and the loop exits normally. It can also
|
//! process does. A [`GenServerRef`] is an *address*: cloning and dropping it
|
||||||
//! end itself: clone a [`StopHandle`] from [`GenServerCtx::stop_handle`] in
|
//! never changes the server's lifetime, and a ref that nobody holds is not a
|
||||||
//! `init` and call [`StopHandle::stop`] from any handler — the loop breaks
|
//! leak — a forgotten server idles until the run ends, when the root-exit
|
||||||
//! after the current message and exits *normally* (OTP's `{stop, normal}`).
|
//! shutdown (see [`Runtime::run`](crate::Runtime::run)) takes it down with
|
||||||
//! This is distinct from `request_stop(self_pid())`, which is an abnormal
|
//! every other unsupervised actor. Anything meant to live long should be
|
||||||
//! `Stopped` and gets a `Transient` child restarted.
|
//! 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
|
//! ## Graceful shutdown
|
||||||
//!
|
//!
|
||||||
@@ -193,9 +202,23 @@
|
|||||||
//! Names and registration: a server can be given a static name so other
|
//! Names and registration: a server can be given a static name so other
|
||||||
//! actors can reach it without holding a `GenServerRef`. Use
|
//! actors can reach it without holding a `GenServerRef`. Use
|
||||||
//! [`GenServerBuilder::named`] to register on start, and the free functions
|
//! [`GenServerBuilder::named`] to register on start, and the free functions
|
||||||
//! [`call`], [`cast`], and [`whereis_server`] to address it by name. Registered
|
//! [`call`], [`cast`], and [`whereis_server`] to address it by name.
|
||||||
//! servers are a natural fit for supervision; see `supervisor` for how to
|
//!
|
||||||
//! build a tree that restarts servers on failure.
|
//! ## 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<Counter> = GenServerName::new("counter");
|
||||||
|
//! OneForOne::new().child(ChildSpec::new(Restart::Permanent, || {
|
||||||
|
//! GenServerBuilder::new(Counter::default()).named(COUNTER).run().unwrap();
|
||||||
|
//! }));
|
||||||
|
//! ```
|
||||||
//!
|
//!
|
||||||
//! ## Limitations
|
//! ## Limitations
|
||||||
//!
|
//!
|
||||||
@@ -317,9 +340,9 @@ enum Envelope<G: GenServer> {
|
|||||||
Cast(G::Cast),
|
Cast(G::Cast),
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A clonable handle to a running server. Cloning yields another sender to the
|
/// A clonable handle to a running server: an *address*, not an owner. Cloning
|
||||||
/// same inbox; the server lives until the last `GenServerRef` is dropped, at which
|
/// yields another sender to the same inbox; dropping refs never ends the
|
||||||
/// point its inbox closes and the loop exits normally.
|
/// server (see the module docs, *When the server stops*).
|
||||||
pub struct GenServerRef<G: GenServer> {
|
pub struct GenServerRef<G: GenServer> {
|
||||||
tx: Sender<Envelope<G>>,
|
tx: Sender<Envelope<G>>,
|
||||||
pid: Pid,
|
pid: Pid,
|
||||||
@@ -423,9 +446,8 @@ impl<G: GenServer> GenServerRef<G> {
|
|||||||
/// a supervisor uses its child's [`Shutdown`](crate::supervisor::Shutdown)
|
/// a supervisor uses its child's [`Shutdown`](crate::supervisor::Shutdown)
|
||||||
/// policy to bound it.
|
/// policy to bound it.
|
||||||
///
|
///
|
||||||
/// This is the right teardown for a server kept alive by a registered
|
/// This is the explicit close: dropping refs never ends a server. Like
|
||||||
/// [`GenServerName`], where dropping every external `GenServerRef` is not enough
|
/// all cooperative cancellation, it is best-effort:
|
||||||
/// 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
|
/// a server wedged in a tight loop with no observation point cannot be
|
||||||
/// stopped this way. Panics if called outside `Runtime::run()`.
|
/// stopped this way. Panics if called outside `Runtime::run()`.
|
||||||
pub fn shutdown(&self) {
|
pub fn shutdown(&self) {
|
||||||
@@ -807,9 +829,10 @@ impl<G: GenServer> GenServerBuilder<G> {
|
|||||||
self
|
self
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Spawn the server actor and hand back its [`GenServerRef`]. The server's
|
/// Spawn the server actor and hand back its [`GenServerRef`] (an address;
|
||||||
/// lifetime is governed by its refs, not by joining, so the backing join
|
/// the server's lifetime is its own, see the module docs). The backing
|
||||||
/// handle is dropped.
|
/// join handle is dropped. For a supervised server use
|
||||||
|
/// [`named`](Self::named) + [`NamedGenServerBuilder::run`] instead.
|
||||||
pub fn start(self) -> GenServerRef<G> {
|
pub fn start(self) -> GenServerRef<G> {
|
||||||
self.spawn_server()
|
self.spawn_server()
|
||||||
}
|
}
|
||||||
@@ -837,13 +860,14 @@ impl<G: GenServer> GenServerBuilder<G> {
|
|||||||
supervisor,
|
supervisor,
|
||||||
stack_opts,
|
stack_opts,
|
||||||
} = self;
|
} = self;
|
||||||
|
let keep = tx.clone();
|
||||||
let handle = match supervisor {
|
let handle = match supervisor {
|
||||||
Some(sup) => crate::scheduler::spawn_under_with(sup, stack_opts, move || {
|
Some(sup) => crate::scheduler::spawn_under_with(sup, stack_opts, move || {
|
||||||
server_loop::<G>(rx, state, infos)
|
server_loop::<G>(keep, rx, state, infos)
|
||||||
|
}),
|
||||||
|
None => crate::scheduler::spawn_with(stack_opts, move || {
|
||||||
|
server_loop::<G>(keep, rx, state, infos)
|
||||||
}),
|
}),
|
||||||
None => {
|
|
||||||
crate::scheduler::spawn_with(stack_opts, move || server_loop::<G>(rx, state, infos))
|
|
||||||
}
|
|
||||||
};
|
};
|
||||||
GenServerRef {
|
GenServerRef {
|
||||||
tx,
|
tx,
|
||||||
@@ -931,19 +955,38 @@ impl<G: GenServer> NamedGenServerBuilder<G> {
|
|||||||
/// The inbox sender is published under the name **from the parent side,
|
/// The inbox sender is published under the name **from the parent side,
|
||||||
/// before this returns**, so a by-name `call` / `cast` resolves the instant
|
/// 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
|
/// `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
|
/// just-spawned server is stopped, so a failed bind leaks no actor.
|
||||||
/// inbox), so a failed bind leaks no actor.
|
|
||||||
pub fn start(self) -> Result<GenServerRef<G>, RegisterError> {
|
pub fn start(self) -> Result<GenServerRef<G>, RegisterError> {
|
||||||
let NamedGenServerBuilder { builder, name } = self;
|
let NamedGenServerBuilder { builder, name } = self;
|
||||||
let server = builder.spawn_server();
|
let server = builder.spawn_server();
|
||||||
match register_with::<Envelope<G>>(server.pid, name, server.tx.clone()) {
|
match register_with::<Envelope<G>>(server.pid, name, server.tx.clone()) {
|
||||||
Ok(()) => Ok(server),
|
Ok(()) => Ok(server),
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
drop(server); // inbox closes → loop exits gracefully
|
crate::scheduler::request_stop(server.pid); // never ran init
|
||||||
Err(e)
|
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::<Envelope<G>>();
|
||||||
|
register_with::<Envelope<G>>(crate::scheduler::self_pid(), name, tx.clone())?;
|
||||||
|
server_loop::<G>(tx, rx, state, infos);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Resolve a [`GenServerName`] to a [`GenServerRef`] when you want a handle to hold or
|
/// Resolve a [`GenServerName`] to a [`GenServerRef`] when you want a handle to hold or
|
||||||
@@ -1002,10 +1045,17 @@ pub fn start_under<G: GenServer>(supervisor: Pid, state: G) -> GenServerRef<G> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn server_loop<G: GenServer>(
|
fn server_loop<G: GenServer>(
|
||||||
|
keep: Sender<Envelope<G>>,
|
||||||
rx: Receiver<Envelope<G>>,
|
rx: Receiver<Envelope<G>>,
|
||||||
state: G,
|
state: G,
|
||||||
mut infos: Vec<Receiver<G::Info>>,
|
mut infos: Vec<Receiver<G::Info>>,
|
||||||
) {
|
) {
|
||||||
|
// 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.
|
// Drop guard — owns the server state and the timer registry.
|
||||||
//
|
//
|
||||||
// Why a guard rather than code after the loop:
|
// Why a guard rather than code after the loop:
|
||||||
@@ -1117,7 +1167,7 @@ fn server_loop<G: GenServer>(
|
|||||||
guard.0.handle_idle();
|
guard.0.handle_idle();
|
||||||
reset_idle(&mut idle_deadline);
|
reset_idle(&mut idle_deadline);
|
||||||
}
|
}
|
||||||
// All ServerRefs dropped → inbox closed → shutdown.
|
// Defensive: the loop holds a sender, so unreachable.
|
||||||
Err(RecvTimeoutError::Disconnected) => break,
|
Err(RecvTimeoutError::Disconnected) => break,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+120
-16
@@ -71,8 +71,12 @@
|
|||||||
//!
|
//!
|
||||||
//! ## Stopping, shutdown, and exits
|
//! ## Stopping, shutdown, and exits
|
||||||
//!
|
//!
|
||||||
//! A machine runs while a [`GenStatemRef`] exists; dropping the last one
|
//! A machine lives until it stops, is shut down, or is killed; a
|
||||||
//! closes the inbox and ends it normally. It can also end itself: any row
|
//! [`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,
|
//! 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
|
//! 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
|
//! 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::link::ExitSignal;
|
||||||
use crate::monitor::{monitor, DownReason};
|
use crate::monitor::{monitor, DownReason};
|
||||||
use crate::pid::Pid;
|
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::scheduler::{cancel_timer, request_shutdown, send_after_to};
|
||||||
use crate::timer::TimerId;
|
use crate::timer::TimerId;
|
||||||
use std::collections::{HashMap, VecDeque};
|
use std::collections::{HashMap, VecDeque};
|
||||||
@@ -162,11 +167,12 @@ pub trait Machine: Send + 'static {
|
|||||||
None
|
None
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Runs as the machine actor exits, on any exit path (inbox closed, a
|
/// Runs as the machine actor exits, on any exit path (a `stop`, a
|
||||||
/// `stop`, a handler panic, a hard stop). Like `gen_server::terminate`:
|
/// graceful shutdown, a handler panic, a hard stop). Like
|
||||||
/// on the panic and hard-stop paths it runs mid-unwind — do not panic or
|
/// `gen_server::terminate`: on the panic and hard-stop paths it runs
|
||||||
/// park there; only on the normal path (`stop`, inbox close) may it do real
|
/// mid-unwind — do not panic or park there; only on the normal path
|
||||||
/// work. The macro's optional `terminate { … }` block generates it.
|
/// (`stop`, shutdown rows) may it do real work. The macro's optional
|
||||||
|
/// `terminate { … }` block generates it.
|
||||||
fn terminate(&mut self) {}
|
fn terminate(&mut self) {}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -460,8 +466,7 @@ pub enum SendError {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// A clonable handle to a running machine. Cloning yields another sender to the
|
/// 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
|
/// same inbox. An address, not an owner: dropping refs never ends the machine.
|
||||||
/// point its inbox closes and the loop exits.
|
|
||||||
pub struct GenStatemRef<M: Machine> {
|
pub struct GenStatemRef<M: Machine> {
|
||||||
tx: Sender<M::Ev>,
|
tx: Sender<M::Ev>,
|
||||||
pid: Pid,
|
pid: Pid,
|
||||||
@@ -528,7 +533,8 @@ impl<M: Machine> GenStatemRef<M> {
|
|||||||
|
|
||||||
/// Spawn `machine` as an actor and hand back its [`GenStatemRef`]. Shape mirrors
|
/// 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
|
/// `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()`.
|
/// Panics if called outside `Runtime::run()`.
|
||||||
pub fn spawn<M: Machine>(machine: M) -> GenStatemRef<M> {
|
pub fn spawn<M: Machine>(machine: M) -> GenStatemRef<M> {
|
||||||
@@ -543,16 +549,104 @@ pub fn spawn<M: Machine>(machine: M) -> GenStatemRef<M> {
|
|||||||
/// Panics if called outside `Runtime::run()`.
|
/// Panics if called outside `Runtime::run()`.
|
||||||
pub fn spawn_with<M: Machine>(opts: crate::scheduler::SpawnOpts, machine: M) -> GenStatemRef<M> {
|
pub fn spawn_with<M: Machine>(opts: crate::scheduler::SpawnOpts, machine: M) -> GenStatemRef<M> {
|
||||||
let (tx, rx) = channel::<M::Ev>();
|
let (tx, rx) = channel::<M::Ev>();
|
||||||
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 {
|
GenStatemRef {
|
||||||
tx,
|
tx,
|
||||||
pid: handle.pid(),
|
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<M> {
|
||||||
|
name: &'static str,
|
||||||
|
_marker: PhantomData<fn() -> M>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<M> GenStatemName<M> {
|
||||||
|
/// 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<M> Copy for GenStatemName<M> {}
|
||||||
|
impl<M> Clone for GenStatemName<M> {
|
||||||
|
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<M: Machine>(name: GenStatemName<M>, machine: M) -> Result<(), RegisterError> {
|
||||||
|
let (tx, rx) = channel::<M::Ev>();
|
||||||
|
register_with::<M::Ev>(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<M: Machine>(name: GenStatemName<M>) -> Option<GenStatemRef<M>> {
|
||||||
|
resolve_named_sender::<M::Ev>(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<M: Machine>(name: GenStatemName<M>, 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<M, T, F>(name: GenStatemName<M>, make: F) -> Result<T, CallError>
|
||||||
|
where
|
||||||
|
M: Machine,
|
||||||
|
T: Send + 'static,
|
||||||
|
F: FnOnce(Reply<T>) -> 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<M: Machine>(name: GenStatemName<M>) {
|
||||||
|
if let Some(m) = whereis_machine(name) {
|
||||||
|
m.shutdown();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// The machine actor body: `on_start`, then one `handle` per event until the
|
/// 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
|
/// row resolves to `stop`, a shutdown row stops it, or the actor is stopped
|
||||||
/// stopped from outside.
|
/// from outside.
|
||||||
///
|
///
|
||||||
/// Intake arms are selected each iteration in priority order — **exits**
|
/// Intake arms are selected each iteration in priority order — **exits**
|
||||||
/// (only when trapping) above **timers** above the **inbox** — so a shutdown
|
/// (only when trapping) above **timers** above the **inbox** — so a shutdown
|
||||||
@@ -566,7 +660,11 @@ pub fn spawn_with<M: Machine>(opts: crate::scheduler::SpawnOpts, machine: M) ->
|
|||||||
/// it back ([`Step::Postponed`]) for the queue; a `handle` that transitions
|
/// it back ([`Step::Postponed`]) for the queue; a `handle` that transitions
|
||||||
/// ([`Step::Transitioned`]) triggers a [`replay`] of the queue in the new
|
/// ([`Step::Transitioned`]) triggers a [`replay`] of the queue in the new
|
||||||
/// state, ahead of the next intake.
|
/// state, ahead of the next intake.
|
||||||
fn statem_loop<M: Machine>(rx: Receiver<M::Ev>, machine: M) {
|
fn statem_loop<M: Machine>(keep: Sender<M::Ev>, rx: Receiver<M::Ev>, 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`
|
// Drop guard — owns the machine and the timer registry, so `terminate`
|
||||||
// fires on every exit path (clean close, `stop`, a handler panic, a hard
|
// 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
|
// stop) and the timer drain is sequenced before it. Same shape as
|
||||||
@@ -685,7 +783,7 @@ fn statem_loop<M: Machine>(rx: Receiver<M::Ev>, machine: M) {
|
|||||||
match rx.try_recv() {
|
match rx.try_recv() {
|
||||||
Ok(Some(ev)) => dispatch(&mut guard.0, &mut cx, &mut postpone, ev),
|
Ok(Some(ev)) => dispatch(&mut guard.0, &mut cx, &mut postpone, ev),
|
||||||
Ok(None) => debug_assert!(false, "ready inbox was empty"),
|
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,
|
Err(_) => break,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -991,8 +1089,14 @@ macro_rules! gen_statem {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl $sm {
|
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> {
|
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)]
|
#[allow(unused_variables)]
|
||||||
|
|||||||
+1
-1
@@ -63,7 +63,7 @@ pub use gen_server::{
|
|||||||
ShutdownAction, StopHandle, TimerHandle, Watcher,
|
ShutdownAction, StopHandle, TimerHandle, Watcher,
|
||||||
};
|
};
|
||||||
pub use gen_statem::{
|
pub use gen_statem::{
|
||||||
CallError as GenStatemCallError, Cx, GenStatemRef, Machine, Reply, Resolution,
|
CallError as GenStatemCallError, Cx, GenStatemName, GenStatemRef, Machine, Reply, Resolution,
|
||||||
SendError as GenStatemSendError,
|
SendError as GenStatemSendError,
|
||||||
};
|
};
|
||||||
pub use introspect::{
|
pub use introspect::{
|
||||||
|
|||||||
+12
-1
@@ -1965,7 +1965,18 @@ fn shutdown_forest_roots(inner: &Arc<RuntimeInner>, root: Pid) {
|
|||||||
.slot_at(parent)
|
.slot_at(parent)
|
||||||
.is_some_and(|ps| ps.is_live_for(parent));
|
.is_some_and(|ps| ps.is_live_for(parent));
|
||||||
if !parent_live {
|
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
|
||||||
|
});
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+17
-2
@@ -662,6 +662,17 @@ pub fn request_shutdown<A>(pid: Pid<A>) {
|
|||||||
// cold lock (generation-verified), then acts outside the lock: a trap send
|
// cold lock (generation-verified), then acts outside the lock: a trap send
|
||||||
// may unpark the receiver, and `request_stop_inner` re-takes the lock.
|
// may unpark the receiver, and `request_stop_inner` re-takes the lock.
|
||||||
pub(crate) fn request_shutdown_inner(inner: &RuntimeInner, pid: Pid, from: Pid) {
|
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<bool> {
|
||||||
let trap = match inner.slot_at(pid) {
|
let trap = match inner.slot_at(pid) {
|
||||||
Some(slot) => {
|
Some(slot) => {
|
||||||
let cold = slot.cold.lock();
|
let cold = slot.cold.lock();
|
||||||
@@ -679,9 +690,13 @@ pub(crate) fn request_shutdown_inner(inner: &RuntimeInner, pid: Pid, from: Pid)
|
|||||||
from,
|
from,
|
||||||
reason: crate::monitor::DownReason::Shutdown,
|
reason: crate::monitor::DownReason::Shutdown,
|
||||||
});
|
});
|
||||||
|
Some(true)
|
||||||
}
|
}
|
||||||
Some(None) => request_stop_inner(inner, pid),
|
Some(None) => {
|
||||||
None => {}
|
request_stop_inner(inner, pid);
|
||||||
|
Some(false)
|
||||||
|
}
|
||||||
|
None => None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+24
-2
@@ -46,17 +46,32 @@ mod inner {
|
|||||||
#[derive(Clone, Debug)]
|
#[derive(Clone, Debug)]
|
||||||
pub enum Event {
|
pub enum Event {
|
||||||
// Actor lifecycle
|
// Actor lifecycle
|
||||||
Spawn { parent: Pid, child: Pid },
|
Spawn {
|
||||||
|
parent: Pid,
|
||||||
|
child: Pid,
|
||||||
|
},
|
||||||
Resume(Pid),
|
Resume(Pid),
|
||||||
Yield(Pid),
|
Yield(Pid),
|
||||||
Park(Pid),
|
Park(Pid),
|
||||||
Done(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
|
// Wakeup paths
|
||||||
UnparkDirect(Pid), // unpark() saw Parked -> re-queued immediately
|
UnparkDirect(Pid), // unpark() saw Parked -> re-queued immediately
|
||||||
UnparkDeferred(Pid), // unpark() saw Runnable -> set pending_unpark flag
|
UnparkDeferred(Pid), // unpark() saw Runnable -> set pending_unpark flag
|
||||||
UnparkFlagConsumed(Pid), // scheduler saw flag on Park -> re-queued instead
|
UnparkFlagConsumed(Pid), // scheduler saw flag on Park -> re-queued instead
|
||||||
// Channel
|
// Channel
|
||||||
Send { sender: Pid, receiver: Option<Pid> },
|
Send {
|
||||||
|
sender: Pid,
|
||||||
|
receiver: Option<Pid>,
|
||||||
|
},
|
||||||
RecvPark(Pid),
|
RecvPark(Pid),
|
||||||
RecvWake(Pid),
|
RecvWake(Pid),
|
||||||
// Queue
|
// Queue
|
||||||
@@ -253,6 +268,13 @@ mod inner {
|
|||||||
Event::Yield(p) => ("yield".into(), p.index()),
|
Event::Yield(p) => ("yield".into(), p.index()),
|
||||||
Event::Park(p) => ("park".into(), p.index()),
|
Event::Park(p) => ("park".into(), p.index()),
|
||||||
Event::Done(p) => ("done".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::UnparkDirect(p) => ("unpark_direct".into(), p.index()),
|
||||||
Event::UnparkDeferred(p) => ("unpark_deferred".into(), p.index()),
|
Event::UnparkDeferred(p) => ("unpark_deferred".into(), p.index()),
|
||||||
Event::UnparkFlagConsumed(p) => ("unpark_flag_consumed".into(), p.index()),
|
Event::UnparkFlagConsumed(p) => ("unpark_flag_consumed".into(), p.index()),
|
||||||
|
|||||||
+5
-5
@@ -88,7 +88,7 @@ impl GenServer for Lifecycle {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// init -> handle_call -> (drop last ref closes inbox) -> terminate.
|
// init -> handle_call -> shutdown -> terminate.
|
||||||
#[test]
|
#[test]
|
||||||
fn init_and_terminate_run() {
|
fn init_and_terminate_run() {
|
||||||
let log = Arc::new(Mutex::new(Vec::new()));
|
let log = Arc::new(Mutex::new(Vec::new()));
|
||||||
@@ -96,9 +96,9 @@ fn init_and_terminate_run() {
|
|||||||
run(move || {
|
run(move || {
|
||||||
let server = start(Lifecycle { log: log2 });
|
let server = start(Lifecycle { log: log2 });
|
||||||
server.call(()).unwrap();
|
server.call(()).unwrap();
|
||||||
// Dropping the only ref closes the inbox; the server breaks out of its
|
// Refs are addresses: dropping one does not end the server. The
|
||||||
// recv loop and runs terminate. run() will not return until it has.
|
// explicit close does, and waits for terminate.
|
||||||
drop(server);
|
server.shutdown();
|
||||||
});
|
});
|
||||||
assert_eq!(*log.lock().unwrap(), vec!["init", "call", "terminate"]);
|
assert_eq!(*log.lock().unwrap(), vec!["init", "call", "terminate"]);
|
||||||
}
|
}
|
||||||
@@ -702,7 +702,7 @@ fn no_timer_survives_exit() {
|
|||||||
let _ = server.call(()).unwrap(); // sync: periodic armed
|
let _ = server.call(()).unwrap(); // sync: periodic armed
|
||||||
smarm::sleep(Duration::from_millis(45)); // a couple of ticks
|
smarm::sleep(Duration::from_millis(45)); // a couple of ticks
|
||||||
let mon = smarm::monitor(server.pid());
|
let mon = smarm::monitor(server.pid());
|
||||||
drop(server); // inbox closes → loop exits → guard drains timers
|
server.shutdown(); // loop exits → guard drains timers
|
||||||
// Clean Down ⇒ the loop returned without the no-leak assert aborting.
|
// Clean Down ⇒ the loop returned without the no-leak assert aborting.
|
||||||
assert!(mon.rx.recv().is_ok());
|
assert!(mon.rx.recv().is_ok());
|
||||||
let at_exit = f_read.lock().unwrap().len();
|
let at_exit = f_read.lock().unwrap().len();
|
||||||
|
|||||||
@@ -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<Mutex<Vec<String>>>);
|
||||||
|
impl Log {
|
||||||
|
fn push(&self, s: impl Into<String>) {
|
||||||
|
self.0.lock().unwrap().push(s.into());
|
||||||
|
}
|
||||||
|
fn get(&self) -> Vec<String> {
|
||||||
|
self.0.lock().unwrap().clone()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
struct Counter {
|
||||||
|
log: Log,
|
||||||
|
n: u64,
|
||||||
|
trap: bool,
|
||||||
|
stop: Option<StopHandle<Counter>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<Self>) {
|
||||||
|
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<Counter> = 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();
|
||||||
|
});
|
||||||
|
}
|
||||||
@@ -196,12 +196,18 @@ fn linked_peer_death_reaches_handle_exit() {
|
|||||||
r.cast(Cast::Note("still-serving")).unwrap();
|
r.cast(Cast::Note("still-serving")).unwrap();
|
||||||
sleep(Duration::from_millis(20));
|
sleep(Duration::from_millis(20));
|
||||||
a.store(true, Ordering::SeqCst);
|
a.store(true, Ordering::SeqCst);
|
||||||
let mon = monitor(r.pid());
|
r.shutdown(); // graceful; waits for terminate
|
||||||
drop(r); // inbox closes → clean exit
|
|
||||||
let _ = mon.rx.recv(); // don't let the root-exit sweep race terminate
|
|
||||||
});
|
});
|
||||||
assert!(alive.load(Ordering::SeqCst));
|
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]
|
#[test]
|
||||||
|
|||||||
@@ -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<Mutex<Vec<&'static str>>>);
|
||||||
|
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<u64>),
|
||||||
|
}
|
||||||
|
|
||||||
|
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<Sm> = 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"]
|
||||||
|
);
|
||||||
|
}
|
||||||
@@ -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<Self>) {}
|
||||||
|
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}"
|
||||||
|
);
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user