feat(gen_statem): graceful-shutdown parity — trap_exit, shutdown/exit rows, stop, terminate

Mirrors the gen_server surface in gen_statem's event model:
- Cx::trap_exit() (in the initial enter): a shutdown request then arrives
  as the Shutdown event, routed by state through `shutdown` rows (default
  for a state with no row: stop); linked-peer deaths as `exit <pat>` rows
  (default: drop). Non-trapping machines are stopped outright, as before.
- Cx::stop(): normal self-exit after the current event; `stop` tail keyword
  is sugar for { cx.stop(); prev }.
- Machine::terminate (optional `terminate { … }` macro block), run from a
  Drop guard on every exit path; the guard also drains armed timers.
- Machine::shutdown_ev / exit_ev (defaults None) so hand-written machines
  keep compiling; GenStatemRef::shutdown() is graceful and waits.
- Loop selects exits > timers > inbox.

Tests: tests/gen_statem_shutdown.rs.
This commit is contained in:
Claude (sandbox)
2026-08-19 16:25:15 +00:00
parent 250f31265b
commit 6ceb138f5f
2 changed files with 579 additions and 61 deletions
+360 -61
View File
@@ -68,10 +68,33 @@
//! inbox or timer event; a replayed event may postpone again (it re-queues for
//! the next transition). See the macro docs for the row surface and [`Step`]
//! for how a postpone surfaces to the loop.
//!
//! ## 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
//! 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
//! restarted). [`terminate`](Machine::terminate) — the optional `terminate
//! { … }` macro block — runs on every exit path.
//!
//! From outside, [`GenStatemRef::shutdown`] (or a plain
//! [`request_shutdown`](crate::request_shutdown), which is what a supervisor
//! sends) asks the machine to stop. By default a machine does not trap exits
//! and the request stops it outright. A machine that calls
//! [`cx.trap_exit()`](Cx::trap_exit) in its initial `enter` instead receives
//! it as the **`shutdown` event**, routed by state like any other — a
//! `Connected` state may transition into `Draining` and stop later from a
//! timeout row, while a state with no `shutdown` row takes the macro's
//! default, `stop`. Linked-peer deaths reach a trapping machine as
//! `exit <pat>` events; an unmatched one is dropped like an unmatched info.
use crate::channel::{channel, select, Receiver, Sender};
use crate::channel::{channel, select, Receiver, Selectable, Sender};
use crate::link::ExitSignal;
use crate::monitor::{monitor, DownReason};
use crate::pid::Pid;
use crate::scheduler::{cancel_timer, send_after_to};
use crate::scheduler::{cancel_timer, request_shutdown, send_after_to};
use crate::timer::TimerId;
use std::collections::{HashMap, VecDeque};
use std::marker::PhantomData;
@@ -119,6 +142,32 @@ pub trait Machine: Send + 'static {
/// state's `enter`, and returns [`Step::Transitioned`]; a stay or unmatched
/// event returns [`Step::Stayed`].
fn handle(&mut self, ev: Self::Ev, cx: &mut Cx<Self::Ev>) -> Step<Self::Ev>;
/// Wrap a graceful shutdown request into this machine's event, so a
/// trapping machine (see [`Cx::trap_exit`]) can route it **by state**. The
/// macro generates it as `Ev::Shutdown` and matches it in `shutdown` rows;
/// its default for a state that writes no such row is `stop`. A
/// hand-written machine that returns `None` (the default) is simply
/// stopped — the loop breaks and `terminate` runs on the normal path.
fn shutdown_ev() -> Option<Self::Ev> {
None
}
/// Wrap a linked peer's death (an [`ExitSignal`] that is not a shutdown
/// request, delivered only when trapping) into this machine's event. The
/// macro generates it as `Ev::Exit(sig)` and matches it in `exit <pat>`
/// rows; an unmatched exit is silently dropped, like an unmatched info.
/// A hand-written machine that returns `None` (the default) drops it.
fn exit_ev(_sig: ExitSignal) -> Option<Self::Ev> {
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.
fn terminate(&mut self) {}
}
// ---------------------------------------------------------------------------
@@ -247,6 +296,11 @@ impl Timers {
pub struct Cx<Ev> {
sys_tx: Sender<Sys>,
reg: Arc<Mutex<Timers>>,
/// Set by [`trap_exit`](Self::trap_exit) during `on_start`; read once by
/// the loop right after, fixed for the machine's life.
trap: bool,
/// Set by [`stop`](Self::stop); the loop breaks after the current event.
stop: bool,
_ev: PhantomData<fn() -> Ev>,
}
@@ -255,10 +309,30 @@ impl<Ev> Cx<Ev> {
Cx {
sys_tx,
reg,
trap: false,
stop: false,
_ev: PhantomData,
}
}
/// Trap exits for the machine's lifetime: a shutdown request then arrives
/// as the `shutdown` event (routed by state) and linked-peer deaths as
/// `exit` events, instead of stopping the machine outright. Call it in the
/// initial state's `enter` (i.e. during `on_start`); later calls have no
/// effect.
pub fn trap_exit(&mut self) {
self.trap = true;
}
/// End the machine after the current event: the loop breaks and the actor
/// exits *normally* (OTP's `{stop, normal}`); `terminate` runs on the
/// normal path. Anything still queued or postponed is dropped. The `stop`
/// tail keyword in a macro row is sugar for `{ cx.stop(); prev }`.
/// Mirrors gen_server's [`StopHandle`](crate::gen_server::StopHandle).
pub fn stop(&mut self) {
self.stop = true;
}
/// Arm the **state timeout**: fire a `state_timeout` event after `after` in
/// the current state. Auto-reset on any state change (the loop cancels and
/// clears it on every real transition), so it measures quiet time *within* a
@@ -433,6 +507,19 @@ impl<M: Machine> GenStatemRef<M> {
self.send(ev).map_err(|_| CallError::Down)?;
rx.recv().map_err(|_| CallError::Down)
}
/// Ask the machine to shut down and block until it has fully exited, so
/// [`Machine::terminate`] has run by the time this returns. A trapping
/// machine winds down through its `shutdown` rows; any other is stopped
/// outright. Returns immediately if the machine is already gone. Waits as
/// long as the machine takes — a supervisor bounds that with its child's
/// [`Shutdown`](crate::supervisor::Shutdown) policy. Mirrors
/// [`GenServerRef::shutdown`](crate::gen_server::GenServerRef::shutdown).
pub fn shutdown(&self) {
let mon = monitor(self.pid);
request_shutdown(self.pid);
let _ = mon.rx.recv();
}
}
// ---------------------------------------------------------------------------
@@ -464,32 +551,88 @@ pub fn spawn_with<M: Machine>(opts: crate::scheduler::SpawnOpts, machine: M) ->
}
/// The machine actor body: `on_start`, then one `handle` per event until the
/// inbox closes (all refs dropped → graceful shutdown).
/// inbox closes (all refs dropped), a row resolves to `stop`, or the actor is
/// stopped from outside.
///
/// Two intake sources are selected each iteration with the **timer arm above
/// the inbox**, so a timeout fire is never starved by inbox traffic: `sys_rx`
/// carries timer fires armed through `cx`, `rx` is the user inbox. A fire is
/// turned into the matching internal event (`state_timeout` / `timeout(name)`)
/// and run through the same `handle` dispatch as an inbox event — the
/// gen_statem model, where timeouts surface as ordinary events.
/// Intake arms are selected each iteration in priority order — **exits**
/// (only when trapping) above **timers** above the **inbox** — so a shutdown
/// request or a timeout fire is never starved by inbox traffic. `sys_rx`
/// carries timer fires armed through `cx`; a fire is turned into the matching
/// internal event (`state_timeout` / `timeout(name)`) and run through the same
/// `handle` dispatch as an inbox event — the gen_statem model, where timeouts
/// (and, when trapping, shutdown and exits) surface as ordinary events.
///
/// The loop owns the **postpone queue**: a `handle` that defers its event hands
/// 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<M: Machine>(rx: Receiver<M::Ev>, mut machine: M) {
fn statem_loop<M: Machine>(rx: Receiver<M::Ev>, machine: M) {
// 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
// gen_server's guard; see the rationale there.
struct Terminate<M: Machine>(M, Arc<Mutex<Timers>>);
impl<M: Machine> Drop for Terminate<M> {
fn drop(&mut self) {
{
let mut reg = match self.1.lock() {
Ok(g) => g,
Err(e) => panic!("smarm: gen_statem reg lock poisoned (core corrupt): {e}"),
};
if let Some((_, sub)) = reg.state.take() {
cancel_timer(sub);
}
for (_, (_, sub)) in reg.named.drain() {
cancel_timer(sub);
}
}
self.0.terminate();
}
}
let (sys_tx, sys_rx) = channel::<Sys>();
let reg = Arc::new(Mutex::new(Timers::new()));
let mut guard = Terminate(machine, reg.clone());
// The loop owns `cx` (and through it a `sys_tx` clone) for its whole life,
// so the sys arm never closes from under us — no auto-close dance needed.
let mut cx = Cx::new(sys_tx, reg.clone());
// Events deferred by `postpone` rows, replayed FIFO on the next transition.
let mut postpone: VecDeque<M::Ev> = VecDeque::new();
machine.on_start(&mut cx);
guard.0.on_start(&mut cx);
// Trapping is opted into during on_start and fixed for the loop's life.
// The inbox is armed only when set: an untrapped machine keeps the
// two-arm select, and a shutdown request simply stops it as
// `request_stop` would.
let exits: Option<Receiver<ExitSignal>> = cx.trap.then(crate::link::trap_exit);
loop {
// Timer arm first: a ready fire is taken in preference to the inbox.
let i = select(&[&sys_rx, &rx]);
if i == 0 {
// Arm order encodes priority: exits → timers → inbox.
let ne = exits.is_some() as usize;
let i = {
let mut arms: Vec<&dyn Selectable> = Vec::with_capacity(3);
if let Some(e) = &exits {
arms.push(e);
}
arms.push(&sys_rx);
arms.push(&rx);
select(&arms)
};
if i < ne {
// Exit arm: a shutdown request or a linked peer's death. The trap
// inbox lives for the loop's life, so it never closes.
let sig = exits.as_ref().and_then(|e| e.try_recv().ok().flatten());
match sig {
Some(sig) if sig.reason == DownReason::Shutdown => match M::shutdown_ev() {
Some(ev) => dispatch(&mut guard.0, &mut cx, &mut postpone, ev),
None => cx.stop(),
},
Some(sig) => {
if let Some(ev) = M::exit_ev(sig) {
dispatch(&mut guard.0, &mut cx, &mut postpone, ev);
}
}
None => {}
}
} else if i == ne {
match sys_rx.try_recv() {
Ok(Some(fire)) => {
// Confirm the fire is still the live one before dispatching:
@@ -528,7 +671,7 @@ fn statem_loop<M: Machine>(rx: Receiver<M::Ev>, mut machine: M) {
}
};
if let Some(ev) = ev {
dispatch(&mut machine, &mut cx, &mut postpone, ev);
dispatch(&mut guard.0, &mut cx, &mut postpone, ev);
}
}
// Single-receiver: nothing can drain the arm between select's
@@ -540,12 +683,16 @@ fn statem_loop<M: Machine>(rx: Receiver<M::Ev>, mut machine: M) {
}
} else {
match rx.try_recv() {
Ok(Some(ev)) => dispatch(&mut machine, &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"),
// All GenStatemRefs dropped → inbox closed → shutdown.
// All GenStatemRefs dropped → inbox closed → normal exit.
Err(_) => break,
}
}
// A handler (or a replay) asked to stop: a normal exit.
if cx.stop {
break;
}
// Observation point so a machine fed a hot inbox stays preemptible and
// cancellable.
crate::check!();
@@ -554,7 +701,8 @@ fn statem_loop<M: Machine>(rx: Receiver<M::Ev>, mut machine: M) {
/// Run one event through `handle` and act on its [`Step`]: stash a deferred
/// event on the postpone queue, or — on a real transition — [`replay`] the
/// queue in the new state. A stay/unmatched event needs nothing further.
/// queue in the new state. A stay/unmatched event needs nothing further. A
/// [`Cx::stop`] raised by the handler skips the replay; the loop breaks next.
fn dispatch<M: Machine>(
machine: &mut M,
cx: &mut Cx<M::Ev>,
@@ -564,14 +712,19 @@ fn dispatch<M: Machine>(
match machine.handle(ev, cx) {
Step::Postponed(ev) => postpone.push_back(ev),
Step::Stayed => {}
Step::Transitioned => replay(machine, cx, postpone),
Step::Transitioned => {
if !cx.stop {
replay(machine, cx, postpone)
}
}
}
}
/// Replay deferred events after a real transition: each goes back through
/// `handle` in FIFO order, in the now-current state. An event that postpones
/// again re-queues (to wait for the *next* transition); one that transitions
/// re-arms the replay, so a later state can in turn drain what is still pending.
/// re-arms the replay, so a later state can in turn drain what is still pending;
/// one that raises [`Cx::stop`] ends the replay (and the machine) at once.
/// Subsequent events in a batch already see the post-transition state, since
/// `handle` reads the live state cell — the outer loop only re-runs to give
/// re-queued events another pass once a transition has occurred within a batch.
@@ -590,6 +743,9 @@ fn replay<M: Machine>(machine: &mut M, cx: &mut Cx<M::Ev>, postpone: &mut VecDeq
Step::Stayed => {}
Step::Transitioned => transitioned = true,
}
if cx.stop {
return;
}
}
if !transitioned {
return;
@@ -691,8 +847,9 @@ fn replay<M: Machine>(machine: &mut M, cx: &mut Cx<M::Ev>, postpone: &mut VecDeq
/// // The transition table. Group rows by current state with `on <pat>`.
/// // A row is: <kind> <event-pattern> [if <guard>] => <tail> ,
/// // where <kind> is one of `cast`, `call`, `info`, `state_timeout`
/// // (no pattern — it is a unit event), or `timeout <name-pattern>`,
/// // and the tail is one of:
/// // (no pattern — it is a unit event), `timeout <name-pattern>`,
/// // `shutdown` (unit; trapping machines only) or `exit <sig-pattern>`
/// // (trapping only), and the tail is one of:
/// // * a state tag `Door::Closed` (transition, or "stay"
/// // if it equals current)
/// // * a block ending in one `{ data.enters += 1; Door::Closed }`
@@ -701,6 +858,9 @@ fn replay<M: Machine>(machine: &mut M, cx: &mut Cx<M::Ev>, postpone: &mut VecDeq
/// // * the keyword `postpone` (defer until next
/// // transition; cast/call/
/// // info only)
/// // * the keyword `stop` (end the machine
/// // normally; sugar for
/// // `{ cx.stop(); prev }`)
/// on Door::Open => {
/// cast DoorCast::Push => Door::Closed,
/// // An armed state-timeout surfaces as an ordinary event:
@@ -724,6 +884,9 @@ fn replay<M: Machine>(machine: &mut M, cx: &mut Cx<M::Ev>, postpone: &mut VecDeq
/// on _ => {
/// call DoorCall::GetState(r) => { r.reply(prev); prev },
/// }
///
/// // Optional: runs as the machine exits, on every exit path.
/// terminate { data.enters = 0; }
/// }
/// ```
///
@@ -748,6 +911,12 @@ fn replay<M: Machine>(machine: &mut M, cx: &mut Cx<M::Ev>, postpone: &mut VecDeq
/// ignores info writes no `info` rows at all. State-timeouts and named
/// timeouts have **no** such default: a state that can see one must handle it
/// (or `unhandled` it) or the match is non-exhaustive.
/// * **`shutdown` defaults to `stop`, `exit` to a silent drop.** Both reach a
/// machine only if its initial `enter` called `cx.trap_exit()`; a
/// non-trapping machine is stopped outright by a shutdown request. Write
/// `shutdown => …` rows only in the states that want to wind down first
/// (transition into a draining state, `stop` later); write `exit sig => …`
/// rows to react to linked-peer deaths. Neither is postponable.
/// * **Stay** = return the current tag. The `prev` you named in `context` is
/// bound to the pre-handler state for exactly this — handy in any-state
/// (`on _`) rows where there is no single literal tag to write.
@@ -773,12 +942,12 @@ fn replay<M: Machine>(machine: &mut M, cx: &mut Cx<M::Ev>, postpone: &mut VecDeq
/// # What it emits
///
/// The unified `enum $Ev` (the `Cast`/`Call`/`Info` wrappers plus the internal
/// `StateTimeout` / `Timeout` events), `struct $Sm { state, data }`,
/// `$Sm::start(init, data) -> GenStatemRef<$Sm>`, and the `Machine` impl:
/// `on_start` runs the initial `enter`; `handle` is the dispatch match plus the
/// stay/transition/unhandled apply-tail (the cell's sole writer, which also
/// auto-resets the state-timeout on every real transition); and the `enter`
/// dispatch.
/// `StateTimeout` / `Timeout` / `Shutdown` / `Exit` events), `struct $Sm {
/// state, data }`, `$Sm::start(init, data) -> GenStatemRef<$Sm>`, and the
/// `Machine` impl: `on_start` runs the initial `enter`; `handle` is the
/// dispatch match plus the stay/transition/unhandled apply-tail (the cell's
/// sole writer, which also auto-resets the state-timeout on every real
/// transition); the `enter` dispatch; and `terminate` when the block is given.
///
/// # Limitation
///
@@ -794,9 +963,10 @@ macro_rules! gen_statem {
context ( $data:ident , $cur:ident , $cx:ident ) ;
enter { $( $est:pat => $ebody:expr ),+ $(,)? }
$( on $st:pat => { $($rows:tt)* } )+
$( terminate { $($tbody:tt)* } )?
) => {
/// Unified inbox payload: the user's `cast`/`call`/`info` enums folded
/// together with the runtime's internal timeout events.
/// together with the runtime's internal events.
enum $Ev {
Cast($Cast),
Call($Call),
@@ -807,6 +977,12 @@ macro_rules! gen_statem {
StateTimeout,
/// A named timeout fired (matched in `timeout <pat>` rows).
Timeout(&'static str),
/// A graceful shutdown request reached this (trapping) machine
/// (matched in `shutdown` rows). A state with no such row `stop`s.
Shutdown,
/// A linked peer died (trapping only; matched in `exit <pat>`
/// rows). An unmatched exit is silently dropped.
Exit($crate::ExitSignal),
}
struct $sm {
@@ -840,11 +1016,27 @@ macro_rules! gen_statem {
$Ev::Timeout(name)
}
fn shutdown_ev() -> Option<$Ev> {
Some($Ev::Shutdown)
}
fn exit_ev(sig: $crate::ExitSignal) -> Option<$Ev> {
Some($Ev::Exit(sig))
}
fn on_start(&mut self, $cx: &mut $crate::gen_statem::Cx<$Ev>) {
let s = self.state;
self.enter(s, $cx);
}
$(
#[allow(unused_variables)]
fn terminate(&mut self) {
let $data = &mut self.data;
$($tbody)*
}
)?
#[allow(unused_variables)]
#[deny(unreachable_patterns)] // conflicting rows must fail even though
// this match is external-macro-expanded
@@ -860,7 +1052,7 @@ macro_rules! gen_statem {
// untouched), then the consuming `match (state, event)` whose
// value is this `Resolution`.
let next: $crate::gen_statem::Resolution<$State> =
$crate::gen_statem!(@arms ($Ev) ($cur, ev) [ ] [ ]
$crate::gen_statem!(@arms ($Ev) ($cur, ev, $cx) [ ] [ ]
$( on $st => { $($rows)* } )+);
match next {
$crate::gen_statem::Resolution::To(s) if s == $cur => {
@@ -896,7 +1088,7 @@ macro_rules! gen_statem {
// is the `Info` silent-drop — last and broadest, so per-state `info` rows
// stay reachable; cast/call/timeouts get no fallback, so a forgotten pair is
// still E0004).
(@arms ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ]) => {
(@arms ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ]) => {
{
// Phase 1 — postpone routing (borrow-only). A guard on a postpone
// row runs here, against by-ref bindings, so it must not depend on
@@ -913,140 +1105,247 @@ macro_rules! gen_statem {
// returned for it.
match ($ss, $se) {
$($arms)*
// Macro-injected defaults, last and broadest so per-state rows
// stay reachable; a user's own catch-all row may shadow them
// entirely, hence the allow.
#[allow(unreachable_patterns)]
(_, $Ev::Info(_)) => $crate::gen_statem::Resolution::Unhandled,
#[allow(unreachable_patterns)]
(_, $Ev::Exit(_)) => $crate::gen_statem::Resolution::Unhandled,
#[allow(unreachable_patterns)]
(_, $Ev::Shutdown) => { $cx.stop(); $crate::gen_statem::Resolution::To($ss) },
}
}
};
// Open an on-block: remember its state pat, drain its rows, then continue.
(@arms ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ]
(@arms ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ]
on $st:pat => { $($rows:tt)* } $($more:tt)*
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se) [ $($arms)* ] [ $($post)* ] ($st)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx) [ $($arms)* ] [ $($post)* ] ($st)
{ $($rows)* } { $($more)* })
};
// ===== @rows: drain one on-block's rows, threading both accs =============
// cast, explicit refusal
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ cast $ev:pat $(if $g:expr)? => unhandled , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Cast($ev)) $(if $g)? => $crate::gen_statem::Resolution::Unhandled, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// cast, stop (end the machine normally after this event)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ cast $ev:pat $(if $g:expr)? => stop , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Cast($ev)) $(if $g)? => { $cx.stop(); $crate::gen_statem::Resolution::To($ss) }, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// cast, postpone (defer the event; the replay in a later state handles it)
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ cast $ev:pat $(if $g:expr)? => postpone , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Cast($ev)) $(if $g)? => unreachable!("postponed event is replayed, not dispatched here"), ]
[ $($post)* ($st, $Ev::Cast($ev)) $(if $g)? => true, ]
($st) { $($rows)* } { $($more)* })
};
// cast, transition / stay / branch
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ cast $ev:pat $(if $g:expr)? => $tail:expr , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Cast($ev)) $(if $g)? => $crate::gen_statem::Resolution::To($tail.into()), ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// call, explicit refusal
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ call $ev:pat $(if $g:expr)? => unhandled , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Call($ev)) $(if $g)? => $crate::gen_statem::Resolution::Unhandled, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// call, stop (end the machine normally after this event)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ call $ev:pat $(if $g:expr)? => stop , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Call($ev)) $(if $g)? => { $cx.stop(); $crate::gen_statem::Resolution::To($ss) }, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// call, postpone (the Reply rides inside the event onto the queue)
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ call $ev:pat $(if $g:expr)? => postpone , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Call($ev)) $(if $g)? => unreachable!("postponed event is replayed, not dispatched here"), ]
[ $($post)* ($st, $Ev::Call($ev)) $(if $g)? => true, ]
($st) { $($rows)* } { $($more)* })
};
// call, transition / stay / branch
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ call $ev:pat $(if $g:expr)? => $tail:expr , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Call($ev)) $(if $g)? => $crate::gen_statem::Resolution::To($tail.into()), ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// info, explicit refusal
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ info $ev:pat $(if $g:expr)? => unhandled , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Info($ev)) $(if $g)? => $crate::gen_statem::Resolution::Unhandled, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// info, stop (end the machine normally after this event)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ info $ev:pat $(if $g:expr)? => stop , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Info($ev)) $(if $g)? => { $cx.stop(); $crate::gen_statem::Resolution::To($ss) }, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// info, postpone
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ info $ev:pat $(if $g:expr)? => postpone , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Info($ev)) $(if $g)? => unreachable!("postponed event is replayed, not dispatched here"), ]
[ $($post)* ($st, $Ev::Info($ev)) $(if $g)? => true, ]
($st) { $($rows)* } { $($more)* })
};
// info, transition / stay / branch
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ info $ev:pat $(if $g:expr)? => $tail:expr , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Info($ev)) $(if $g)? => $crate::gen_statem::Resolution::To($tail.into()), ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// state_timeout, explicit refusal (unit event — no pattern; not postponable)
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ state_timeout $(if $g:expr)? => unhandled , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::StateTimeout) $(if $g)? => $crate::gen_statem::Resolution::Unhandled, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// state_timeout, stop (end the machine normally after this event)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ state_timeout $(if $g:expr)? => stop , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::StateTimeout) $(if $g)? => { $cx.stop(); $crate::gen_statem::Resolution::To($ss) }, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// state_timeout, transition / stay / branch
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ state_timeout $(if $g:expr)? => $tail:expr , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::StateTimeout) $(if $g)? => $crate::gen_statem::Resolution::To($tail.into()), ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// timeout, explicit refusal (pattern matches the name; not postponable)
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ timeout $ev:pat $(if $g:expr)? => unhandled , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Timeout($ev)) $(if $g)? => $crate::gen_statem::Resolution::Unhandled, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// timeout, stop (end the machine normally after this event)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ timeout $ev:pat $(if $g:expr)? => stop , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Timeout($ev)) $(if $g)? => { $cx.stop(); $crate::gen_statem::Resolution::To($ss) }, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// timeout, transition / stay / branch
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ timeout $ev:pat $(if $g:expr)? => $tail:expr , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se)
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Timeout($ev)) $(if $g)? => $crate::gen_statem::Resolution::To($tail.into()), ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// shutdown, explicit refusal (unit event — no pattern; not postponable; a state with no shutdown row defaults to `stop`)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ shutdown $(if $g:expr)? => unhandled , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Shutdown) $(if $g)? => $crate::gen_statem::Resolution::Unhandled, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// shutdown, stop (end the machine normally after this event)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ shutdown $(if $g:expr)? => stop , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Shutdown) $(if $g)? => { $cx.stop(); $crate::gen_statem::Resolution::To($ss) }, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// shutdown, transition / stay / branch
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ shutdown $(if $g:expr)? => $tail:expr , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Shutdown) $(if $g)? => $crate::gen_statem::Resolution::To($tail.into()), ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// exit, explicit refusal (pattern matches the ExitSignal; not postponable)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ exit $ev:pat $(if $g:expr)? => unhandled , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Exit($ev)) $(if $g)? => $crate::gen_statem::Resolution::Unhandled, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// exit, stop (end the machine normally after this event)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ exit $ev:pat $(if $g:expr)? => stop , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Exit($ev)) $(if $g)? => { $cx.stop(); $crate::gen_statem::Resolution::To($ss) }, ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// exit, transition / stay / branch
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ exit $ev:pat $(if $g:expr)? => $tail:expr , $($rows:tt)* } { $($more:tt)* }
) => {
$crate::gen_statem!(@rows ($Ev) ($ss, $se, $cx)
[ $($arms)* ($st, $Ev::Exit($ev)) $(if $g)? => $crate::gen_statem::Resolution::To($tail.into()), ]
[ $($post)* ]
($st) { $($rows)* } { $($more)* })
};
// this block is drained: hand the remaining on-blocks back to @arms
(@rows ($Ev:ident) ($ss:expr, $se:expr) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
(@rows ($Ev:ident) ($ss:expr, $se:expr, $cx:ident) [ $($arms:tt)* ] [ $($post:tt)* ] ($st:pat)
{ } { $($more:tt)* }
) => {
$crate::gen_statem!(@arms ($Ev) ($ss, $se) [ $($arms)* ] [ $($post)* ] $($more)*)
$crate::gen_statem!(@arms ($Ev) ($ss, $se, $cx) [ $($arms)* ] [ $($post)* ] $($more)*)
};
}
+219
View File
@@ -0,0 +1,219 @@
//! gen_statem graceful shutdown — the gen_server surface, in state-machine
//! clothes. Where gen_server routes a shutdown request to a `handle_shutdown`
//! method, a gen_statem gets it as an **event** so it can be routed by state:
//!
//! - A machine that does not opt in (`cx.trap_exit()` in the initial `enter`)
//! is stopped outright by `request_shutdown`, exactly as by `request_stop`.
//! - A trapping machine sees the request as a `shutdown` row (a unit event
//! like `state_timeout`). The macro's default, when a state writes no
//! `shutdown` row, is `stop` — the loop breaks and `terminate` runs on the
//! normal path. A row may instead transition (e.g. into a Draining state)
//! and `stop` later from any row via the `stop` tail keyword.
//! - Linked-peer deaths reach a trapping machine as `exit <pat>` rows; an
//! unmatched exit is silently dropped, like an unmatched info.
//! - `terminate { … }` is an optional macro block, run on every exit path.
use smarm::gen_statem;
use smarm::gen_statem::{GenStatemRef, Reply};
use smarm::{link, monitor, request_shutdown, run, sleep, spawn, DownReason, ExitSignal};
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 {
Idle,
Draining,
}
struct D {
log: Log,
trap: bool,
exits: Vec<ExitSignal>,
}
enum Cast {
Note(&'static str),
StopNow,
}
enum Call {
Exits(Reply<usize>),
}
gen_statem! {
machine: Sm { state: S, data: D };
event: Ev { cast: Cast, call: Call, info: () };
context(data, prev, cx);
enter {
S::Idle => if data.trap { cx.trap_exit() },
S::Draining => { data.log.push("draining"); cx.state_timeout(Duration::from_millis(30)); },
}
on S::Idle => {
// Shutdown in Idle: go drain first, stop later.
shutdown => S::Draining,
cast Cast::StopNow => stop,
state_timeout => unhandled,
}
on S::Draining => {
// Drained: end the machine normally.
state_timeout => { data.log.push("drained"); cx.stop(); prev },
// A second request while draining is ignored.
shutdown => unhandled,
cast Cast::StopNow => stop,
}
on _ => {
cast Cast::Note(s) => { data.log.push(s); prev },
call Call::Exits(r) => { r.reply(data.exits.len()); prev },
exit sig => { data.log.push("exit"); data.exits.push(sig); prev },
timeout _ => unhandled,
}
terminate {
data.log.push("terminate");
}
}
/// A machine with no `shutdown` rows at all: the macro default applies.
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum P {
On,
}
struct PD {
log: Log,
}
enum PCast {}
enum PCall {}
gen_statem! {
machine: Plain { state: P, data: PD };
event: PEv { cast: PCast, call: PCall, info: () };
context(data, prev, cx);
enter { P::On => cx.trap_exit(), }
on P::On => {
cast _ => unhandled,
call _ => unhandled,
state_timeout => unhandled,
timeout _ => unhandled,
}
terminate { data.log.push("terminate"); }
}
fn settled(log: &Log, trap: bool) -> GenStatemRef<Sm> {
let r = Sm::start(
S::Idle,
D {
log: log.clone(),
trap,
exits: Vec::new(),
},
);
sleep(Duration::from_millis(20)); // let on_start (trap_exit) run
r
}
#[test]
fn non_trapping_machine_is_stopped_outright() {
let log = Log::default();
let l = log.clone();
run(move || {
let r = settled(&l, false);
let mon = monitor(r.pid());
request_shutdown(r.pid());
assert_eq!(mon.rx.recv().unwrap().reason, DownReason::Stopped);
});
assert_eq!(log.get(), vec!["terminate"]);
}
#[test]
fn shutdown_row_routes_by_state_and_stop_tail_exits_normally() {
let log = Log::default();
let l = log.clone();
run(move || {
let r = settled(&l, true);
let mon = monitor(r.pid());
request_shutdown(r.pid());
// The second request lands in Draining and is `unhandled` (ignored).
sleep(Duration::from_millis(5));
request_shutdown(r.pid());
assert_eq!(mon.rx.recv().unwrap().reason, DownReason::Exit);
});
assert_eq!(log.get(), vec!["draining", "drained", "terminate"]);
}
#[test]
fn default_shutdown_is_stop() {
let log = Log::default();
let l = log.clone();
run(move || {
let r = Plain::start(P::On, PD { log: l });
sleep(Duration::from_millis(20));
let mon = monitor(r.pid());
request_shutdown(r.pid());
assert_eq!(mon.rx.recv().unwrap().reason, DownReason::Exit);
});
assert_eq!(log.get(), vec!["terminate"]);
}
#[test]
fn stop_tail_from_a_cast_is_a_normal_exit() {
let log = Log::default();
let l = log.clone();
run(move || {
let r = settled(&l, false);
let mon = monitor(r.pid());
r.send(Ev::Cast(Cast::Note("a"))).unwrap();
r.send(Ev::Cast(Cast::StopNow)).unwrap();
r.send(Ev::Cast(Cast::Note("after-stop"))).unwrap(); // never dispatched
assert_eq!(mon.rx.recv().unwrap().reason, DownReason::Exit);
});
assert_eq!(log.get(), vec!["a", "terminate"]);
}
#[test]
fn linked_peer_death_reaches_exit_row() {
let log = Log::default();
let l = log.clone();
run(move || {
let r = settled(&l, true);
let pid = r.pid();
let peer = spawn(move || {
link(pid);
panic!("peer dies");
});
let _ = peer.join();
sleep(Duration::from_millis(20));
r.send(Ev::Cast(Cast::Note("still-running"))).unwrap();
assert_eq!(r.call(|r| Ev::Call(Call::Exits(r))).unwrap(), 1);
r.shutdown();
});
assert_eq!(
log.get(),
vec!["exit", "still-running", "draining", "drained", "terminate"]
);
}
#[test]
fn ref_shutdown_is_graceful_and_waits() {
let log = Log::default();
let l = log.clone();
run(move || {
let r = settled(&l, true);
r.shutdown();
// terminate has run by the time shutdown() returns.
assert_eq!(l.get(), vec!["draining", "drained", "terminate"]);
});
}