diff --git a/src/gen_statem.rs b/src/gen_statem.rs index 6533020..7fdf96a 100644 --- a/src/gen_statem.rs +++ b/src/gen_statem.rs @@ -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 ` 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) -> Step; + + /// 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 { + 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 ` + /// 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 { + 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 { sys_tx: Sender, reg: Arc>, + /// 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 Ev>, } @@ -255,10 +309,30 @@ impl Cx { 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 GenStatemRef { 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(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(rx: Receiver, mut machine: M) { +fn statem_loop(rx: Receiver, 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, Arc>); + impl Drop for Terminate { + 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::(); 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 = 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> = 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(rx: Receiver, 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(rx: Receiver, 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(rx: Receiver, 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( machine: &mut M, cx: &mut Cx, @@ -564,14 +712,19 @@ fn dispatch( 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(machine: &mut M, cx: &mut Cx, postpone: &mut VecDeq Step::Stayed => {} Step::Transitioned => transitioned = true, } + if cx.stop { + return; + } } if !transitioned { return; @@ -691,8 +847,9 @@ fn replay(machine: &mut M, cx: &mut Cx, postpone: &mut VecDeq /// // The transition table. Group rows by current state with `on `. /// // A row is: [if ] => , /// // where is one of `cast`, `call`, `info`, `state_timeout` -/// // (no pattern — it is a unit event), or `timeout `, -/// // and the tail is one of: +/// // (no pattern — it is a unit event), `timeout `, +/// // `shutdown` (unit; trapping machines only) or `exit ` +/// // (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(machine: &mut M, cx: &mut Cx, 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(machine: &mut M, cx: &mut Cx, 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(machine: &mut M, cx: &mut Cx, 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(machine: &mut M, cx: &mut Cx, 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 ` 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 ` + /// 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)*) }; } diff --git a/tests/gen_statem_shutdown.rs b/tests/gen_statem_shutdown.rs new file mode 100644 index 0000000..fa72988 --- /dev/null +++ b/tests/gen_statem_shutdown.rs @@ -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 ` 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>>); +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, +} + +enum Cast { + Note(&'static str), + StopNow, +} +enum Call { + Exits(Reply), +} + +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 { + 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"]); + }); +}