feat(scheduler,supervisor,gen_server): graceful shutdown — request_shutdown, child Shutdown policy, handle_shutdown
Lift OTP's `exit(Pid, shutdown)` + child-spec `shutdown` wholesale.
scheduler / runtime
- `request_shutdown(pid)`: the polite stop. A target trapping exits gets an
`ExitSignal { reason: DownReason::Shutdown }` on its trap inbox and keeps
running; a non-trapping target is stopped as by `request_stop`, which is
now documented as the hard stop (`exit(Pid, kill)`). Dead pid: no-op.
- `RuntimeHandle::request_shutdown` for the off-runtime (signal thread) path;
`from == ROOT_PID` there.
- `DownReason::Shutdown` — appears only in ExitSignal, never in Down (a
complying target exits *normally*).
supervisor
- `ChildSpec::shutdown(Shutdown::{BrutalKill, Timeout(d), Infinity})`,
default Timeout(5s). Every supervisor-initiated stop (ordered shutdown and
OneForAll/RestForOne sibling cycling) is: request_shutdown → await the
child's Signal up to the grace → request_stop → await. Sequential, reverse
start order.
- The supervisor traps exits; a Shutdown ExitSignal runs the ordered
shutdown and `run()` returns normally, so `request_shutdown(root_sup)`
tears a whole tree down top-down with each child's grace period.
- FIX: a hard `request_stop` on a supervisor previously orphaned its
children (the ordered shutdown lived after the loop, and the unwind
skipped it). `Live` (the by_pid map) now carries a drop guard that
fire-and-forget hard-stops live children when unwinding.
gen_server
- `GenServerCtx::trap_exit()` opt-in in `init`; the trap inbox becomes arm 0
of the loop's select. Shutdown ExitSignal → `handle_shutdown() ->
ShutdownAction::{Exit, Continue}` (default Exit: loop breaks, `terminate`
runs on the normal path and may block). Other ExitSignals →
`handle_exit(sig)`.
- `GenServerCtx::stop_handle() -> StopHandle`, `stop()` ends the server
after the current message with a *normal* exit — the missing
`{stop, normal, State}`; `request_stop(self_pid())` was the only self-exit
and it is abnormal (Transient restarts it).
- `GenServerRef::shutdown()` / `gen_server::shutdown(name)` now go through
`request_shutdown`.
Tests: tests/shutdown.rs, tests/supervisor_shutdown.rs,
tests/gen_server_shutdown.rs. Full suite green; fmt + clippy --lib clean.
This commit is contained in:
+149
-19
@@ -127,16 +127,42 @@
|
|||||||
//! - [`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, or an
|
//! on every exit path (all `GenServerRef`s dropped, a handler panic, a
|
||||||
//! explicit [`GenServerRef::shutdown`]), not only on clean shutdown. Keep it
|
//! cooperative stop, or a graceful shutdown), not only on clean shutdown.
|
||||||
//! short and non-blocking: if `terminate` panics while the server is already
|
//! Keep it non-panicking: on the panic and hard-stop paths it runs
|
||||||
//! unwinding from a handler panic, the process aborts.
|
//! 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
|
||||||
|
//! (see below) may it do real work.
|
||||||
//!
|
//!
|
||||||
//! ## When the server stops
|
//! ## When the server stops
|
||||||
//!
|
//!
|
||||||
//! The server runs as long as at least one [`GenServerRef`] exists. When the last
|
//! The server runs as long as at least one [`GenServerRef`] exists. When the last
|
||||||
//! one is dropped, the inbox closes and the loop exits gracefully. To stop a
|
//! one is dropped, the inbox closes and the loop exits normally. It can also
|
||||||
//! server explicitly and wait for it to finish, call [`GenServerRef::shutdown`].
|
//! 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
|
||||||
|
//!
|
||||||
|
//! From outside, [`GenServerRef::shutdown`] (or a plain
|
||||||
|
//! [`request_shutdown`](crate::request_shutdown), which is what a supervisor
|
||||||
|
//! sends) asks the server to stop. What happens next is the server's choice:
|
||||||
|
//!
|
||||||
|
//! - By default a server does not trap exits, and the request stops it
|
||||||
|
//! outright at its next observation point — `terminate` runs mid-unwind.
|
||||||
|
//! - A server that calls [`GenServerCtx::trap_exit`] in `init` receives the
|
||||||
|
//! request as [`GenServer::handle_shutdown`]. Return
|
||||||
|
//! [`ShutdownAction::Exit`] (the default) to have the loop break and
|
||||||
|
//! `terminate` run on the normal path, where it may block; return
|
||||||
|
//! [`ShutdownAction::Continue`] to keep serving — e.g. to drain in-flight
|
||||||
|
//! work — and end the server later with a [`StopHandle`]. The supervisor's
|
||||||
|
//! [`Shutdown`](crate::supervisor::Shutdown) policy bounds how long that
|
||||||
|
//! may take before it falls back to a hard stop.
|
||||||
|
//!
|
||||||
|
//! A trapping server also receives the deaths of its linked peers as
|
||||||
|
//! [`GenServer::handle_exit`] messages instead of dying with them.
|
||||||
//!
|
//!
|
||||||
//! If the server panics inside a handler, the panic unwinds the server thread.
|
//! If the server panics inside a handler, the panic unwinds the server thread.
|
||||||
//! Any caller currently waiting in `call` sees `Err(ServerDown)`: the reply
|
//! Any caller currently waiting in `call` sees `Err(ServerDown)`: the reply
|
||||||
@@ -181,10 +207,12 @@
|
|||||||
use crate::channel::{
|
use crate::channel::{
|
||||||
channel, select, select_timeout, Receiver, RecvTimeoutError, Selectable, Sender,
|
channel, select, select_timeout, Receiver, RecvTimeoutError, Selectable, Sender,
|
||||||
};
|
};
|
||||||
|
use crate::link::ExitSignal;
|
||||||
|
use crate::monitor::DownReason;
|
||||||
use crate::monitor::{demonitor, monitor, Down, Monitor};
|
use crate::monitor::{demonitor, monitor, Down, Monitor};
|
||||||
use crate::pid::Pid;
|
use crate::pid::Pid;
|
||||||
use crate::registry::{register_with, resolve_named_sender, RegisterError};
|
use crate::registry::{register_with, resolve_named_sender, RegisterError};
|
||||||
use crate::scheduler::{cancel_timer, request_stop, send_after_to};
|
use crate::scheduler::{cancel_timer, request_shutdown, send_after_to};
|
||||||
use crate::timer::TimerId;
|
use crate::timer::TimerId;
|
||||||
use std::cell::Cell;
|
use std::cell::Cell;
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
@@ -254,10 +282,34 @@ pub trait GenServer: Send + 'static {
|
|||||||
/// Default: no-op.
|
/// Default: no-op.
|
||||||
fn handle_idle(&mut self) {}
|
fn handle_idle(&mut self) {}
|
||||||
|
|
||||||
|
/// A graceful shutdown request (a [`request_shutdown`](crate::request_shutdown)
|
||||||
|
/// reaching this server), delivered only if `init` called
|
||||||
|
/// [`GenServerCtx::trap_exit`]. Return [`ShutdownAction::Exit`] to stop
|
||||||
|
/// now (the default), or [`ShutdownAction::Continue`] to keep serving and
|
||||||
|
/// end the server later with a [`StopHandle`].
|
||||||
|
fn handle_shutdown(&mut self) -> ShutdownAction {
|
||||||
|
ShutdownAction::Exit
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A linked peer's abnormal death (an [`ExitSignal`] that is not a
|
||||||
|
/// shutdown request), delivered only if `init` called
|
||||||
|
/// [`GenServerCtx::trap_exit`]. Default: drop it.
|
||||||
|
fn handle_exit(&mut self, _sig: ExitSignal) {}
|
||||||
|
|
||||||
/// Runs as the server actor exits, on any exit path (see module docs).
|
/// Runs as the server actor exits, on any exit path (see module docs).
|
||||||
fn terminate(&mut self) {}
|
fn terminate(&mut self) {}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// What a server does with a shutdown request; see [`GenServer::handle_shutdown`].
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum ShutdownAction {
|
||||||
|
/// Break the loop now. `terminate` runs on the normal path and may block.
|
||||||
|
Exit,
|
||||||
|
/// Keep dispatching. The server is expected to end itself with a
|
||||||
|
/// [`StopHandle`] once it is done winding down.
|
||||||
|
Continue,
|
||||||
|
}
|
||||||
|
|
||||||
/// What travels the server's single inbox channel: a synchronous call (with a
|
/// What travels the server's single inbox channel: a synchronous call (with a
|
||||||
/// reply sender) or an asynchronous cast. Private — callers use [`GenServerRef`].
|
/// reply sender) or an asynchronous cast. Private — callers use [`GenServerRef`].
|
||||||
enum Envelope<G: GenServer> {
|
enum Envelope<G: GenServer> {
|
||||||
@@ -362,9 +414,14 @@ impl<G: GenServer> GenServerRef<G> {
|
|||||||
|
|
||||||
/// Stop the server and block until it has fully exited.
|
/// Stop the server and block until it has fully exited.
|
||||||
///
|
///
|
||||||
/// Sends a cooperative stop signal to the server actor and waits for it to
|
/// Asks the server to shut down (a [`request_shutdown`](crate::request_shutdown))
|
||||||
/// exit, so [`GenServer::terminate`] has run by the time this returns.
|
/// and waits for it to exit, so [`GenServer::terminate`] has run by the
|
||||||
/// Returns immediately if the server is already gone.
|
/// time this returns. A trapping server gets to wind down via
|
||||||
|
/// [`GenServer::handle_shutdown`]; any other is stopped outright. Returns
|
||||||
|
/// immediately if the server is already gone. Waits as long as the server
|
||||||
|
/// takes — the caller, not the server, decides whether that is acceptable;
|
||||||
|
/// a supervisor uses its child's [`Shutdown`](crate::supervisor::Shutdown)
|
||||||
|
/// policy to bound it.
|
||||||
///
|
///
|
||||||
/// This is the right teardown for a server kept alive by a registered
|
/// This is the right teardown for a server kept alive by a registered
|
||||||
/// [`GenServerName`], where dropping every external `GenServerRef` is not enough
|
/// [`GenServerName`], where dropping every external `GenServerRef` is not enough
|
||||||
@@ -373,7 +430,7 @@ impl<G: GenServer> GenServerRef<G> {
|
|||||||
/// 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) {
|
||||||
let mon = monitor(self.pid);
|
let mon = monitor(self.pid);
|
||||||
request_stop(self.pid);
|
request_shutdown(self.pid);
|
||||||
// The Down lands when the server finalizes; an already-dead target makes
|
// The Down lands when the server finalizes; an already-dead target makes
|
||||||
// `monitor` deliver NoProc immediately, so this never blocks forever.
|
// `monitor` deliver NoProc immediately, so this never blocks forever.
|
||||||
let _ = mon.rx.recv();
|
let _ = mon.rx.recv();
|
||||||
@@ -396,6 +453,8 @@ enum Sys<G: GenServer> {
|
|||||||
/// the payload factory, dispatches it to [`GenServer::handle_timer`], and
|
/// the payload factory, dispatches it to [`GenServer::handle_timer`], and
|
||||||
/// re-arms the next tick before returning.
|
/// re-arms the next tick before returning.
|
||||||
Tick(crate::timer::TimerId),
|
Tick(crate::timer::TimerId),
|
||||||
|
/// The state asked to end the server (via [`StopHandle::stop`]).
|
||||||
|
Stop,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The server loop's runtime hook, passed to [`GenServer::init`]. Hands out the
|
/// The server loop's runtime hook, passed to [`GenServer::init`]. Hands out the
|
||||||
@@ -411,9 +470,29 @@ pub struct GenServerCtx<G: GenServer> {
|
|||||||
/// because `init` holds only `&ctx`; not `Send`, but `GenServerCtx` is only ever
|
/// because `init` holds only `&ctx`; not `Send`, but `GenServerCtx` is only ever
|
||||||
/// borrowed on the actor's own stack during `init`, never sent.
|
/// borrowed on the actor's own stack during `init`, never sent.
|
||||||
idle: Cell<Option<Duration>>,
|
idle: Cell<Option<Duration>>,
|
||||||
|
/// Whether the loop should trap exits (set via [`trap_exit`](Self::trap_exit)
|
||||||
|
/// during `init`, read by the loop after).
|
||||||
|
trap: Cell<bool>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<G: GenServer> GenServerCtx<G> {
|
impl<G: GenServer> GenServerCtx<G> {
|
||||||
|
/// Trap exits for the server's lifetime: shutdown requests then arrive as
|
||||||
|
/// [`GenServer::handle_shutdown`] and linked-peer deaths as
|
||||||
|
/// [`GenServer::handle_exit`], instead of stopping the server outright.
|
||||||
|
/// Call this once during `init`.
|
||||||
|
pub fn trap_exit(&self) {
|
||||||
|
self.trap.set(true);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A clonable handle that lets the state end the server from any handler
|
||||||
|
/// (a normal exit; see the module docs). Store it on the state during
|
||||||
|
/// `init`.
|
||||||
|
pub fn stop_handle(&self) -> StopHandle<G> {
|
||||||
|
StopHandle {
|
||||||
|
sys_tx: self.sys_tx.clone(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// A clonable handle to the loop's monitor intake. Store it in the state
|
/// A clonable handle to the loop's monitor intake. Store it in the state
|
||||||
/// during `init` to watch monitors from later handlers.
|
/// during `init` to watch monitors from later handlers.
|
||||||
pub fn watcher(&self) -> Watcher<G> {
|
pub fn watcher(&self) -> Watcher<G> {
|
||||||
@@ -452,6 +531,30 @@ impl<G: GenServer> GenServerCtx<G> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Lets a server's state end the server, cloned from
|
||||||
|
/// [`GenServerCtx::stop_handle`] during `init`. [`stop`](Self::stop) makes the
|
||||||
|
/// loop break after the current message and exit normally; `terminate` runs on
|
||||||
|
/// the normal path.
|
||||||
|
pub struct StopHandle<G: GenServer> {
|
||||||
|
sys_tx: Sender<Sys<G>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<G: GenServer> Clone for StopHandle<G> {
|
||||||
|
fn clone(&self) -> Self {
|
||||||
|
StopHandle {
|
||||||
|
sys_tx: self.sys_tx.clone(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<G: GenServer> StopHandle<G> {
|
||||||
|
/// End the server after the current message. Idempotent; a no-op once the
|
||||||
|
/// server is gone.
|
||||||
|
pub fn stop(&self) {
|
||||||
|
let _ = self.sys_tx.send(Sys::Stop);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Per-server timer bookkeeping, shared between the loop and every
|
/// Per-server timer bookkeeping, shared between the loop and every
|
||||||
/// [`TimerHandle`] clone. A gen_server actor is single-threaded — handlers
|
/// [`TimerHandle`] clone. A gen_server actor is single-threaded — handlers
|
||||||
/// and the loop never run concurrently — so this `Mutex` is always
|
/// and the loop never run concurrently — so this `Mutex` is always
|
||||||
@@ -974,9 +1077,14 @@ fn server_loop<G: GenServer>(
|
|||||||
sys_tx,
|
sys_tx,
|
||||||
reg: reg.clone(),
|
reg: reg.clone(),
|
||||||
idle: Cell::new(None),
|
idle: Cell::new(None),
|
||||||
|
trap: Cell::new(false),
|
||||||
};
|
};
|
||||||
guard.0.init(&ctx);
|
guard.0.init(&ctx);
|
||||||
let idle = ctx.idle.get();
|
let idle = ctx.idle.get();
|
||||||
|
// Trapping is opted into during init and fixed for the loop's life. The
|
||||||
|
// inbox is armed only when set: an untrapped server keeps the fast path,
|
||||||
|
// and a shutdown request simply stops it as `request_stop` would.
|
||||||
|
let exits: Option<Receiver<ExitSignal>> = ctx.trap.get().then(crate::link::trap_exit);
|
||||||
drop(ctx);
|
drop(ctx);
|
||||||
|
|
||||||
let mut monitors: Vec<Monitor> = Vec::new();
|
let mut monitors: Vec<Monitor> = Vec::new();
|
||||||
@@ -992,7 +1100,7 @@ fn server_loop<G: GenServer>(
|
|||||||
};
|
};
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
if monitors.is_empty() && !sys_open && infos.is_empty() {
|
if exits.is_none() && monitors.is_empty() && !sys_open && infos.is_empty() {
|
||||||
// Fast path: no extra arms, no select overhead — park directly on
|
// Fast path: no extra arms, no select overhead — park directly on
|
||||||
// the inbox. Mirrors the inbox arm of the select path below; any
|
// the inbox. Mirrors the inbox arm of the select path below; any
|
||||||
// change there must be applied here too.
|
// change there must be applied here too.
|
||||||
@@ -1020,16 +1128,21 @@ fn server_loop<G: GenServer>(
|
|||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
// Slow path: one or more extra arms live — build the arm slice and
|
// Slow path: one or more extra arms live — build the arm slice and
|
||||||
// select. Arm order encodes priority: downs → system → infos →
|
// select. Arm order encodes priority: exits → downs → system →
|
||||||
// inbox. The slice is rebuilt each iteration because the monitor
|
// infos → inbox (a shutdown request is noticed under any load).
|
||||||
// and info sets shrink/grow. Mirrors the fast-path inbox park
|
// The slice is rebuilt each iteration because the monitor and info
|
||||||
// above; keep them in sync.
|
// sets shrink/grow. Mirrors the fast-path inbox park above; keep
|
||||||
let nd = monitors.len(); // monitor band: [0, nd)
|
// them in sync.
|
||||||
|
let ne = exits.is_some() as usize; // exit arm: [0, ne)
|
||||||
|
let nd = ne + monitors.len(); // monitor band: [ne, nd)
|
||||||
let nw = sys_open as usize; // system arm: [nd, nd+nw)
|
let nw = sys_open as usize; // system arm: [nd, nd+nw)
|
||||||
// info band: [nd+nw, nd+nw+ni)
|
// info band: [nd+nw, nd+nw+ni)
|
||||||
// inbox arm: [nd+nw+ni]
|
// inbox arm: [nd+nw+ni]
|
||||||
let sel = {
|
let sel = {
|
||||||
let mut arms: Vec<&dyn Selectable> = Vec::with_capacity(nd + nw + infos.len() + 1);
|
let mut arms: Vec<&dyn Selectable> = Vec::with_capacity(nd + nw + infos.len() + 1);
|
||||||
|
if let Some(e) = &exits {
|
||||||
|
arms.push(e);
|
||||||
|
}
|
||||||
for m in &monitors {
|
for m in &monitors {
|
||||||
arms.push(&m.rx);
|
arms.push(&m.rx);
|
||||||
}
|
}
|
||||||
@@ -1055,10 +1168,25 @@ fn server_loop<G: GenServer>(
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
if i < nd {
|
if i < ne {
|
||||||
|
// Exit arm: a shutdown request or a linked peer's death.
|
||||||
|
// The inbox lives for the loop's life, so it never closes.
|
||||||
|
let sig = exits.as_ref().and_then(|e| e.try_recv().ok().flatten());
|
||||||
|
if let Some(sig) = sig {
|
||||||
|
if sig.reason == DownReason::Shutdown {
|
||||||
|
match guard.0.handle_shutdown() {
|
||||||
|
ShutdownAction::Exit => break,
|
||||||
|
ShutdownAction::Continue => {}
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
guard.0.handle_exit(sig);
|
||||||
|
}
|
||||||
|
reset_idle(&mut idle_deadline);
|
||||||
|
}
|
||||||
|
} else if i < nd {
|
||||||
// Monitor band: a Down retires its arm either way (one-shot)
|
// Monitor band: a Down retires its arm either way (one-shot)
|
||||||
// or closes without delivering (defensive; shouldn't happen).
|
// or closes without delivering (defensive; shouldn't happen).
|
||||||
let m = monitors.remove(i);
|
let m = monitors.remove(i - ne);
|
||||||
if let Ok(Some(down)) = m.rx.try_recv() {
|
if let Ok(Some(down)) = m.rx.try_recv() {
|
||||||
guard.0.handle_down(down);
|
guard.0.handle_down(down);
|
||||||
reset_idle(&mut idle_deadline);
|
reset_idle(&mut idle_deadline);
|
||||||
@@ -1067,6 +1195,8 @@ fn server_loop<G: GenServer>(
|
|||||||
match sys_rx.try_recv() {
|
match sys_rx.try_recv() {
|
||||||
// Control intake, not a dispatched message: no idle reset.
|
// Control intake, not a dispatched message: no idle reset.
|
||||||
Ok(Some(Sys::Watch(m))) => monitors.push(m),
|
Ok(Some(Sys::Watch(m))) => monitors.push(m),
|
||||||
|
// The state ended the server: a normal exit.
|
||||||
|
Ok(Some(Sys::Stop)) => break,
|
||||||
Ok(Some(Sys::Timer(id, msg))) => {
|
Ok(Some(Sys::Timer(id, msg))) => {
|
||||||
// The one-shot fired: retire its registry entry so the
|
// The one-shot fired: retire its registry entry so the
|
||||||
// live set tracks only still-pending timers, then
|
// live set tracks only still-pending timers, then
|
||||||
|
|||||||
+7
-7
@@ -60,7 +60,7 @@ pub use channel::{
|
|||||||
pub use gen_server::{
|
pub use gen_server::{
|
||||||
call, cast, shutdown, whereis_server, CallError, CallTimeoutError, CastError, GenServer,
|
call, cast, shutdown, whereis_server, CallError, CallTimeoutError, CastError, GenServer,
|
||||||
GenServerBuilder, GenServerCtx, GenServerName, GenServerRef, NamedGenServerBuilder,
|
GenServerBuilder, GenServerCtx, GenServerName, GenServerRef, NamedGenServerBuilder,
|
||||||
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, GenStatemRef, Machine, Reply, Resolution,
|
||||||
@@ -87,13 +87,13 @@ pub use registry::{
|
|||||||
};
|
};
|
||||||
pub use runtime::{init, Config, Runtime, RuntimeHandle};
|
pub use runtime::{init, Config, Runtime, RuntimeHandle};
|
||||||
pub use scheduler::{
|
pub use scheduler::{
|
||||||
block_on_io, cancel_timer, request_stop, run, self_pid, send_after, send_after_named,
|
block_on_io, cancel_timer, request_shutdown, request_stop, run, self_pid, send_after,
|
||||||
send_after_named_wall, send_after_wall, sleep, sleep_wall, spawn, spawn_addr, spawn_addr_with,
|
send_after_named, send_after_named_wall, send_after_wall, sleep, sleep_wall, spawn, spawn_addr,
|
||||||
spawn_under, spawn_under_with, spawn_with, try_spawn, try_spawn_under_with, wait_readable,
|
spawn_addr_with, spawn_under, spawn_under_with, spawn_with, try_spawn, try_spawn_under_with,
|
||||||
wait_readable_timeout, wait_writable, wait_writable_timeout, yield_now, FdArm, JoinError,
|
wait_readable, wait_readable_timeout, wait_writable, wait_writable_timeout, yield_now, FdArm,
|
||||||
JoinHandle, SpawnError, SpawnOpts,
|
JoinError, JoinHandle, SpawnError, SpawnOpts,
|
||||||
};
|
};
|
||||||
pub use supervisor::{ChildSpec, OneForOne, Restart, Signal, Strategy};
|
pub use supervisor::{ChildSpec, OneForOne, Restart, Shutdown, Signal, Strategy};
|
||||||
pub use timer::TimerId;
|
pub use timer::TimerId;
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|||||||
@@ -98,6 +98,11 @@ pub enum DownReason {
|
|||||||
Panic,
|
Panic,
|
||||||
/// The target was cooperatively cancelled via `request_stop`.
|
/// The target was cooperatively cancelled via `request_stop`.
|
||||||
Stopped,
|
Stopped,
|
||||||
|
/// A graceful shutdown was requested via `request_shutdown`. Only ever
|
||||||
|
/// appears in an [`ExitSignal`](crate::link::ExitSignal) delivered to a
|
||||||
|
/// trapping actor — never in a [`Down`]: a target that honours the request
|
||||||
|
/// exits *normally*, one that does not trap is `Stopped`.
|
||||||
|
Shutdown,
|
||||||
/// The target was already gone (finished and reclaimed, or never alive)
|
/// The target was already gone (finished and reclaimed, or never alive)
|
||||||
/// at the moment `monitor()` was called.
|
/// at the moment `monitor()` was called.
|
||||||
NoProc,
|
NoProc,
|
||||||
|
|||||||
@@ -1510,6 +1510,18 @@ impl RuntimeHandle {
|
|||||||
crate::scheduler::request_stop_inner(&inner, pid);
|
crate::scheduler::request_stop_inner(&inner, pid);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Ask `pid` to shut down gracefully, from any thread. The off-runtime
|
||||||
|
/// equivalent of [`scheduler::request_shutdown`](crate::request_shutdown);
|
||||||
|
/// the delivered [`ExitSignal`](crate::ExitSignal) carries `from ==
|
||||||
|
/// ROOT_PID`, since no actor made the request. A no-op if the runtime has
|
||||||
|
/// been dropped, or if the actor has already exited.
|
||||||
|
pub fn request_shutdown<A>(&self, pid: Pid<A>) {
|
||||||
|
let pid = pid.erase();
|
||||||
|
if let Some(inner) = self.inner.upgrade() {
|
||||||
|
crate::scheduler::request_shutdown_inner(&inner, pid, ROOT_PID);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|||||||
+53
-4
@@ -600,10 +600,12 @@ pub(crate) fn retire_wait() {
|
|||||||
/// [`JoinHandle::join`] reports it as a normal, non-error exit: cooperative
|
/// [`JoinHandle::join`] reports it as a normal, non-error exit: cooperative
|
||||||
/// stop is a controlled shutdown, not a failure.
|
/// stop is a controlled shutdown, not a failure.
|
||||||
///
|
///
|
||||||
/// This is exactly the mechanism `gen_server` shutdown, supervisor restarts,
|
/// This is the *hard* stop — OTP's `exit(Pid, kill)`. It is what a supervisor
|
||||||
/// and structured teardown are built from: reach for [`GenServerRef::shutdown`](crate::GenServerRef::shutdown)
|
/// falls back to when a child overstays its [`Shutdown`](crate::supervisor::Shutdown)
|
||||||
/// or a [`supervisor`](crate::supervisor) instead of calling this directly
|
/// grace period. For a stop the target gets to prepare for, use
|
||||||
/// where those apply.
|
/// [`request_shutdown`]; for structured teardown, reach for
|
||||||
|
/// [`GenServerRef::shutdown`](crate::GenServerRef::shutdown) or a
|
||||||
|
/// [`supervisor`](crate::supervisor) instead of calling this directly.
|
||||||
///
|
///
|
||||||
/// Because it's cooperative, an actor stuck in a tight loop with no
|
/// Because it's cooperative, an actor stuck in a tight loop with no
|
||||||
/// blocking call, no [`check!`](crate::check), and no allocation cannot be
|
/// blocking call, no [`check!`](crate::check), and no allocation cannot be
|
||||||
@@ -635,6 +637,53 @@ pub(crate) fn request_stop_inner(inner: &RuntimeInner, pid: Pid) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Ask an actor to shut down gracefully — OTP's `exit(Pid, shutdown)`, where
|
||||||
|
/// [`request_stop`] is `exit(Pid, kill)`.
|
||||||
|
///
|
||||||
|
/// If the target has called [`trap_exit`](crate::trap_exit), it receives an
|
||||||
|
/// [`ExitSignal`](crate::ExitSignal) with reason
|
||||||
|
/// [`DownReason::Shutdown`](crate::DownReason::Shutdown) on its trap inbox and
|
||||||
|
/// keeps running: the request is advisory, and the target is expected to wind
|
||||||
|
/// down and exit normally in its own time (a supervisor bounds that time with
|
||||||
|
/// its child's [`Shutdown`](crate::supervisor::Shutdown) policy and falls back
|
||||||
|
/// to `request_stop`). A target that is not trapping is stopped exactly as by
|
||||||
|
/// `request_stop`. A dead pid is a no-op.
|
||||||
|
///
|
||||||
|
/// The signal's `from` is the calling actor, or `ROOT_PID` when driven from
|
||||||
|
/// outside the runtime (see [`RuntimeHandle::request_shutdown`](crate::RuntimeHandle::request_shutdown)).
|
||||||
|
pub fn request_shutdown<A>(pid: Pid<A>) {
|
||||||
|
let pid = pid.erase();
|
||||||
|
let from = current_pid().unwrap_or(crate::runtime::ROOT_PID);
|
||||||
|
let _ = try_with_runtime(|inner| request_shutdown_inner(inner, pid, from));
|
||||||
|
}
|
||||||
|
|
||||||
|
// The core of `request_shutdown`. Reads the target's trap sender under its
|
||||||
|
// cold lock (generation-verified), then acts outside the lock: a trap send
|
||||||
|
// may unpark the receiver, and `request_stop_inner` re-takes the lock.
|
||||||
|
pub(crate) fn request_shutdown_inner(inner: &RuntimeInner, pid: Pid, from: Pid) {
|
||||||
|
let trap = match inner.slot_at(pid) {
|
||||||
|
Some(slot) => {
|
||||||
|
let cold = slot.cold.lock();
|
||||||
|
if slot.generation() == pid.generation() {
|
||||||
|
cold.actor.as_ref().map(|a| a.trap.clone())
|
||||||
|
} else {
|
||||||
|
None // stale pid: nothing there to shut down
|
||||||
|
}
|
||||||
|
}
|
||||||
|
None => None,
|
||||||
|
};
|
||||||
|
match trap {
|
||||||
|
Some(Some(tx)) => {
|
||||||
|
let _ = tx.send(crate::link::ExitSignal {
|
||||||
|
from,
|
||||||
|
reason: crate::monitor::DownReason::Shutdown,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
Some(None) => request_stop_inner(inner, pid),
|
||||||
|
None => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
// NoPreempt
|
// NoPreempt
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|||||||
+187
-66
@@ -140,7 +140,8 @@ impl Signal {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
use crate::channel::channel;
|
use crate::channel::{channel, RecvTimeoutError};
|
||||||
|
use crate::monitor::DownReason;
|
||||||
use std::collections::{HashMap, VecDeque};
|
use std::collections::{HashMap, VecDeque};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
@@ -165,15 +166,53 @@ pub enum Restart {
|
|||||||
pub struct ChildSpec {
|
pub struct ChildSpec {
|
||||||
start: Arc<dyn Fn() + Send + Sync + 'static>,
|
start: Arc<dyn Fn() + Send + Sync + 'static>,
|
||||||
restart: Restart,
|
restart: Restart,
|
||||||
|
shutdown: Shutdown,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl ChildSpec {
|
impl ChildSpec {
|
||||||
|
/// A child with the given restart policy and the default
|
||||||
|
/// [`Shutdown::Timeout`] of 5 seconds.
|
||||||
pub fn new(restart: Restart, start: impl Fn() + Send + Sync + 'static) -> Self {
|
pub fn new(restart: Restart, start: impl Fn() + Send + Sync + 'static) -> Self {
|
||||||
Self {
|
Self {
|
||||||
start: Arc::new(start),
|
start: Arc::new(start),
|
||||||
restart,
|
restart,
|
||||||
|
shutdown: Shutdown::default(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Set how the supervisor stops this child (see [`Shutdown`]). A child
|
||||||
|
/// that is itself a supervisor should use [`Shutdown::Infinity`] so its
|
||||||
|
/// own subtree gets its full grace periods.
|
||||||
|
pub fn shutdown(mut self, shutdown: Shutdown) -> Self {
|
||||||
|
self.shutdown = shutdown;
|
||||||
|
self
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// How a supervisor stops a child it is taking down — the OTP child-spec
|
||||||
|
/// `shutdown` value. Applies to every supervisor-initiated stop: the ordered
|
||||||
|
/// shutdown of the whole set and the sibling cycling of
|
||||||
|
/// [`Strategy::OneForAll`] / [`Strategy::RestForOne`].
|
||||||
|
///
|
||||||
|
/// A graceful stop is a [`request_shutdown`](crate::request_shutdown): a child
|
||||||
|
/// that traps exits receives the request as a message and winds down in its
|
||||||
|
/// own time; one that does not is stopped outright.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum Shutdown {
|
||||||
|
/// `request_stop` immediately; no request, no grace period.
|
||||||
|
BrutalKill,
|
||||||
|
/// `request_shutdown`, wait up to the duration for the child to exit, then
|
||||||
|
/// `request_stop` it. The default, at 5 seconds.
|
||||||
|
Timeout(Duration),
|
||||||
|
/// `request_shutdown` and wait however long the child takes. Use for a
|
||||||
|
/// child supervisor, whose subtree has its own timeouts.
|
||||||
|
Infinity,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Default for Shutdown {
|
||||||
|
fn default() -> Self {
|
||||||
|
Shutdown::Timeout(Duration::from_secs(5))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// How a supervisor reacts when one child terminates and a restart is due.
|
/// How a supervisor reacts when one child terminates and a restart is due.
|
||||||
@@ -244,23 +283,35 @@ impl OneForOne {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Run the supervision loop on the current actor. Returns when every child
|
/// Run the supervision loop on the current actor. Returns when every child
|
||||||
/// has reached a terminal, non-restartable state, or when the restart
|
/// has reached a terminal, non-restartable state, when the restart
|
||||||
/// intensity cap is tripped.
|
/// intensity cap is tripped, or when the supervisor is asked to shut down
|
||||||
|
/// (a [`request_shutdown`](crate::request_shutdown) — from its own
|
||||||
|
/// supervisor, or from the app). On every one of those exits the survivors
|
||||||
|
/// are stopped in reverse start order, each per its
|
||||||
|
/// [`Shutdown`] policy, before this returns.
|
||||||
|
///
|
||||||
|
/// The supervisor traps exits for the length of the loop (that is how the
|
||||||
|
/// shutdown request reaches it as a message). Should the supervisor itself
|
||||||
|
/// be hard-stopped with [`request_stop`](crate::request_stop), it unwinds
|
||||||
|
/// without waiting for anything — but a drop guard hard-stops its live
|
||||||
|
/// children on the way out, so the subtree is not orphaned (a child
|
||||||
|
/// supervisor unwinds the same way, recursively).
|
||||||
pub fn run(self) {
|
pub fn run(self) {
|
||||||
let me = crate::scheduler::self_pid();
|
let me = crate::scheduler::self_pid();
|
||||||
let (tx, rx) = channel::<Signal>();
|
let (tx, rx) = channel::<Signal>();
|
||||||
crate::scheduler::register_supervisor_channel(me, tx);
|
crate::scheduler::register_supervisor_channel(me, tx);
|
||||||
|
let exits = crate::link::trap_exit();
|
||||||
|
|
||||||
// pid -> index into `self.children`, for the children currently alive.
|
// pid -> index into `self.children`, for the children currently alive.
|
||||||
let mut by_pid: HashMap<Pid, usize> = HashMap::new();
|
let mut live = Live::default();
|
||||||
let mut active: usize = 0;
|
let mut active: usize = 0;
|
||||||
// Sliding window of recent restart instants, for the intensity cap.
|
// Sliding window of recent restart instants, for the intensity cap.
|
||||||
let mut restarts: Vec<Instant> = Vec::new();
|
let mut restarts: Vec<Instant> = Vec::new();
|
||||||
|
|
||||||
let start_child = |idx: usize, by_pid: &mut HashMap<Pid, usize>| {
|
let start_child = |idx: usize, live: &mut Live| {
|
||||||
let start = self.children[idx].start.clone();
|
let start = self.children[idx].start.clone();
|
||||||
let h = crate::scheduler::spawn_under(me, move || (start)());
|
let h = crate::scheduler::spawn_under(me, move || (start)());
|
||||||
by_pid.insert(h.pid(), idx);
|
live.insert(h.pid(), idx);
|
||||||
// We supervise via the signal funnel, not by joining; drop the
|
// We supervise via the signal funnel, not by joining; drop the
|
||||||
// handle so the child's slot is reclaimed promptly on death (the
|
// handle so the child's slot is reclaimed promptly on death (the
|
||||||
// termination Signal is delivered before reclamation regardless).
|
// termination Signal is delivered before reclamation regardless).
|
||||||
@@ -268,28 +319,105 @@ impl OneForOne {
|
|||||||
};
|
};
|
||||||
|
|
||||||
for idx in 0..self.children.len() {
|
for idx in 0..self.children.len() {
|
||||||
start_child(idx, &mut by_pid);
|
start_child(idx, &mut live);
|
||||||
active += 1;
|
active += 1;
|
||||||
}
|
}
|
||||||
|
|
||||||
// A signal that arrives while we are awaiting stop-confirmations (for a
|
// A signal that arrives while we are awaiting stop-confirmations (for a
|
||||||
// child we are *not* currently stopping) is stashed here and processed
|
// child we are *not* currently stopping) is stashed here and processed
|
||||||
// by the main loop before it blocks on `recv` again.
|
// by the main loop before it blocks again.
|
||||||
let mut pending: VecDeque<Signal> = VecDeque::new();
|
let mut pending: VecDeque<Signal> = VecDeque::new();
|
||||||
let next_signal = |pending: &mut VecDeque<Signal>| -> Option<Signal> {
|
|
||||||
|
// Stop one child per its policy and wait for its termination signal.
|
||||||
|
// Signals for other pids that arrive meanwhile are stashed. Bounded by
|
||||||
|
// construction: `request_stop` (used directly, or as the fallback once
|
||||||
|
// the grace period lapses) always produces a signal.
|
||||||
|
let stop_child = |pid: Pid, idx: usize, pending: &mut VecDeque<Signal>| {
|
||||||
|
let await_one = |deadline: Option<Instant>, pending: &mut VecDeque<Signal>| -> bool {
|
||||||
|
loop {
|
||||||
|
let sig = match pending.iter().position(|s| s.pid() == pid) {
|
||||||
|
Some(i) => pending.remove(i),
|
||||||
|
None => match deadline {
|
||||||
|
None => rx.recv().ok(),
|
||||||
|
Some(dl) => {
|
||||||
|
match rx.recv_timeout(dl.saturating_duration_since(Instant::now()))
|
||||||
|
{
|
||||||
|
Ok(s) => Some(s),
|
||||||
|
Err(RecvTimeoutError::Timeout) => return false,
|
||||||
|
Err(RecvTimeoutError::Disconnected) => None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
};
|
||||||
|
match sig {
|
||||||
|
Some(s) if s.pid() == pid => return true,
|
||||||
|
Some(s) => pending.push_back(s),
|
||||||
|
None => return true, // funnel closed: nothing more can arrive
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
match self.children[idx].shutdown {
|
||||||
|
Shutdown::BrutalKill => {
|
||||||
|
crate::scheduler::request_stop(pid);
|
||||||
|
await_one(None, pending);
|
||||||
|
}
|
||||||
|
Shutdown::Timeout(grace) => {
|
||||||
|
crate::scheduler::request_shutdown(pid);
|
||||||
|
if !await_one(Some(Instant::now() + grace), pending) {
|
||||||
|
crate::scheduler::request_stop(pid);
|
||||||
|
await_one(None, pending);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Shutdown::Infinity => {
|
||||||
|
crate::scheduler::request_shutdown(pid);
|
||||||
|
await_one(None, pending);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
// Stop a set of children in reverse start order, one at a time.
|
||||||
|
let stop_set =
|
||||||
|
|set: &mut Vec<(Pid, usize)>, live: &mut Live, pending: &mut VecDeque<Signal>| {
|
||||||
|
set.sort_unstable_by_key(|x| std::cmp::Reverse(x.1));
|
||||||
|
for (pid, idx) in set.iter() {
|
||||||
|
live.remove(pid);
|
||||||
|
stop_child(*pid, *idx, pending);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
// Wait for the next event: a stashed signal, a child signal, or a
|
||||||
|
// shutdown request. `Ok(sig)`, or `Err(())` when we must wind down.
|
||||||
|
let next_event = |pending: &mut VecDeque<Signal>| -> Result<Signal, ()> {
|
||||||
|
loop {
|
||||||
if let Some(s) = pending.pop_front() {
|
if let Some(s) = pending.pop_front() {
|
||||||
Some(s)
|
return Ok(s);
|
||||||
} else {
|
}
|
||||||
rx.recv().ok()
|
// The trap inbox is arm 0: a shutdown request is noticed even
|
||||||
|
// under a flood of child signals.
|
||||||
|
match crate::channel::select(&[&exits, &rx]) {
|
||||||
|
0 => match exits.try_recv() {
|
||||||
|
Ok(Some(sig)) if sig.reason == DownReason::Shutdown => return Err(()),
|
||||||
|
// Any other exit signal (a linked peer's death — a
|
||||||
|
// supervisor links nothing itself, but may be linked
|
||||||
|
// to) is not ours to act on; a closed trap inbox is
|
||||||
|
// impossible while `exits` is held here.
|
||||||
|
_ => {}
|
||||||
|
},
|
||||||
|
_ => match rx.try_recv() {
|
||||||
|
Ok(Some(s)) => return Ok(s),
|
||||||
|
Ok(None) => {}
|
||||||
|
Err(_) => return Err(()), // funnel closed: nothing left to supervise
|
||||||
|
},
|
||||||
|
}
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
while active > 0 {
|
while active > 0 {
|
||||||
let sig = match next_signal(&mut pending) {
|
let sig = match next_event(&mut pending) {
|
||||||
Some(s) => s,
|
Ok(s) => s,
|
||||||
None => break, // mailbox closed: nothing left to supervise
|
Err(()) => break,
|
||||||
};
|
};
|
||||||
let idx = match by_pid.remove(&sig.pid()) {
|
let idx = match live.remove(&sig.pid()) {
|
||||||
Some(i) => i,
|
Some(i) => i,
|
||||||
None => continue, // stray/duplicate signal
|
None => continue, // stray/duplicate signal
|
||||||
};
|
};
|
||||||
@@ -321,76 +449,69 @@ impl OneForOne {
|
|||||||
restarts.push(now);
|
restarts.push(now);
|
||||||
|
|
||||||
// Which *live* siblings get cycled along with the failed child.
|
// Which *live* siblings get cycled along with the failed child.
|
||||||
// (The failed child is already gone — removed from `by_pid` above.)
|
// (The failed child is already gone — removed from `live` above.)
|
||||||
let mut to_stop: Vec<(Pid, usize)> = match self.strategy {
|
let mut to_stop: Vec<(Pid, usize)> = match self.strategy {
|
||||||
Strategy::OneForOne => Vec::new(),
|
Strategy::OneForOne => Vec::new(),
|
||||||
Strategy::OneForAll => by_pid.iter().map(|(p, i)| (*p, *i)).collect(),
|
Strategy::OneForAll => live.iter().map(|(p, i)| (*p, *i)).collect(),
|
||||||
Strategy::RestForOne => by_pid
|
Strategy::RestForOne => live
|
||||||
.iter()
|
.iter()
|
||||||
.filter(|(_, i)| **i > idx)
|
.filter(|(_, i)| **i > idx)
|
||||||
.map(|(p, i)| (*p, *i))
|
.map(|(p, i)| (*p, *i))
|
||||||
.collect(),
|
.collect(),
|
||||||
};
|
};
|
||||||
// Stop survivors in reverse start order (highest child index first).
|
|
||||||
to_stop.sort_unstable_by_key(|x| std::cmp::Reverse(x.1));
|
|
||||||
|
|
||||||
// The set we will restart: the failed child plus every sibling we
|
// The set we will restart: the failed child plus every sibling we
|
||||||
// are about to stop, restarted in start (ascending index) order.
|
// are about to stop, restarted in start (ascending index) order.
|
||||||
let mut restart_set: Vec<usize> = Vec::with_capacity(to_stop.len() + 1);
|
let mut restart_set: Vec<usize> = Vec::with_capacity(to_stop.len() + 1);
|
||||||
restart_set.push(idx);
|
restart_set.push(idx);
|
||||||
|
restart_set.extend(to_stop.iter().map(|(_, i)| *i));
|
||||||
|
|
||||||
// Request stops, then await each survivor's termination signal
|
// Stop the survivors (each per its policy, reverse start order),
|
||||||
// before restarting. `request_stop` on an already-dead pid is a
|
// then restart the whole set in start order. Net effect on
|
||||||
// no-op; in that case its (already-sent) Exit signal serves as the
|
// `active`: one child died (idx), `to_stop.len()` were stopped,
|
||||||
// confirmation. Any signal for a pid we are *not* awaiting is
|
// and `restart_set.len() == 1 + to_stop.len()` are started — so
|
||||||
// stashed for the main loop.
|
|
||||||
let mut awaiting: Vec<Pid> = Vec::with_capacity(to_stop.len());
|
|
||||||
for (pid, cidx) in &to_stop {
|
|
||||||
by_pid.remove(pid);
|
|
||||||
restart_set.push(*cidx);
|
|
||||||
crate::scheduler::request_stop(*pid);
|
|
||||||
awaiting.push(*pid);
|
|
||||||
}
|
|
||||||
while !awaiting.is_empty() {
|
|
||||||
let s = match next_signal(&mut pending) {
|
|
||||||
Some(s) => s,
|
|
||||||
None => break, // mailbox closed mid-await; stop waiting
|
|
||||||
};
|
|
||||||
if let Some(pos) = awaiting.iter().position(|p| *p == s.pid()) {
|
|
||||||
awaiting.swap_remove(pos);
|
|
||||||
} else {
|
|
||||||
pending.push_back(s);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Restart the whole set in start order. Net effect on `active`:
|
|
||||||
// one child died (idx), `to_stop.len()` were stopped, and
|
|
||||||
// `restart_set.len() == 1 + to_stop.len()` are started — so
|
|
||||||
// `active` is unchanged and needs no adjustment here.
|
// `active` is unchanged and needs no adjustment here.
|
||||||
|
stop_set(&mut to_stop, &mut live, &mut pending);
|
||||||
restart_set.sort_unstable();
|
restart_set.sort_unstable();
|
||||||
for cidx in restart_set {
|
for cidx in restart_set {
|
||||||
start_child(cidx, &mut by_pid);
|
start_child(cidx, &mut live);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Ordered shutdown: stop any survivors in reverse start order and await
|
// Ordered shutdown: stop any survivors in reverse start order, each per
|
||||||
// their termination. On the normal `active == 0` exit `by_pid` is empty
|
// its policy. On the normal `active == 0` exit `live` is empty and this
|
||||||
// and this is a no-op; on a cap-trip or mailbox-closed break it tears
|
// is a no-op; on a shutdown request, a cap-trip, or a closed funnel it
|
||||||
// the remaining children down deterministically instead of leaking them.
|
// tears the remaining children down deterministically.
|
||||||
let mut survivors: Vec<(Pid, usize)> = by_pid.iter().map(|(p, i)| (*p, *i)).collect();
|
let mut survivors: Vec<(Pid, usize)> = live.iter().map(|(p, i)| (*p, *i)).collect();
|
||||||
survivors.sort_unstable_by_key(|x| std::cmp::Reverse(x.1));
|
stop_set(&mut survivors, &mut live, &mut pending);
|
||||||
let mut awaiting: Vec<Pid> = Vec::with_capacity(survivors.len());
|
|
||||||
for (pid, _) in &survivors {
|
|
||||||
crate::scheduler::request_stop(*pid);
|
|
||||||
awaiting.push(*pid);
|
|
||||||
}
|
}
|
||||||
while !awaiting.is_empty() {
|
}
|
||||||
let s = match next_signal(&mut pending) {
|
|
||||||
Some(s) => s,
|
/// The live children of a supervisor, with a drop guard: if the supervisor is
|
||||||
None => break,
|
/// unwound (a hard `request_stop`, or a panic in the loop) its children are
|
||||||
};
|
/// hard-stopped rather than orphaned. Fire-and-forget by necessity — a guard
|
||||||
if let Some(pos) = awaiting.iter().position(|p| *p == s.pid()) {
|
/// running mid-unwind cannot park to await anything.
|
||||||
awaiting.swap_remove(pos);
|
#[derive(Default)]
|
||||||
|
struct Live(HashMap<Pid, usize>);
|
||||||
|
|
||||||
|
impl std::ops::Deref for Live {
|
||||||
|
type Target = HashMap<Pid, usize>;
|
||||||
|
fn deref(&self) -> &Self::Target {
|
||||||
|
&self.0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl std::ops::DerefMut for Live {
|
||||||
|
fn deref_mut(&mut self) -> &mut Self::Target {
|
||||||
|
&mut self.0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Drop for Live {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
if std::thread::panicking() {
|
||||||
|
for pid in self.0.keys() {
|
||||||
|
crate::scheduler::request_stop(*pid);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,216 @@
|
|||||||
|
//! gen_server graceful shutdown.
|
||||||
|
//!
|
||||||
|
//! - A server that does not opt in (`ctx.trap_exit()` in `init`) is stopped
|
||||||
|
//! outright by `request_shutdown`, exactly as by `request_stop`.
|
||||||
|
//! - A trapping server receives the request as `handle_shutdown`. The default
|
||||||
|
//! returns `ShutdownAction::Exit`: the loop breaks and `terminate` runs on
|
||||||
|
//! the normal (non-unwind) path, so it may block. `Continue` keeps the loop
|
||||||
|
//! dispatching; the state later ends itself with a `StopHandle` — the only
|
||||||
|
//! way for a gen_server to exit *normally* on its own (`request_stop` on
|
||||||
|
//! self is an abnormal `Stopped`, which `Transient` restarts).
|
||||||
|
//! - Other exit signals (linked peers dying) reach a trapping server via
|
||||||
|
//! `handle_exit`.
|
||||||
|
|
||||||
|
use smarm::gen_server::{
|
||||||
|
start, GenServer, GenServerBuilder, GenServerCtx, GenServerRef, ShutdownAction, StopHandle,
|
||||||
|
};
|
||||||
|
use smarm::supervisor::{ChildSpec, OneForOne, Restart};
|
||||||
|
use smarm::{link, monitor, request_shutdown, run, self_pid, sleep, spawn, DownReason, ExitSignal};
|
||||||
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
use std::time::Duration;
|
||||||
|
|
||||||
|
#[derive(Default, Clone)]
|
||||||
|
struct Log {
|
||||||
|
events: Arc<Mutex<Vec<&'static str>>>,
|
||||||
|
}
|
||||||
|
impl Log {
|
||||||
|
fn push(&self, e: &'static str) {
|
||||||
|
self.events.lock().unwrap().push(e);
|
||||||
|
}
|
||||||
|
fn get(&self) -> Vec<&'static str> {
|
||||||
|
self.events.lock().unwrap().clone()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A server with configurable shutdown behaviour.
|
||||||
|
struct Srv {
|
||||||
|
log: Log,
|
||||||
|
trap: bool,
|
||||||
|
action: ShutdownAction,
|
||||||
|
stop: Option<StopHandle<Srv>>,
|
||||||
|
exits: Arc<Mutex<Vec<ExitSignal>>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Srv {
|
||||||
|
fn new(log: &Log, trap: bool, action: ShutdownAction) -> Self {
|
||||||
|
Srv {
|
||||||
|
log: log.clone(),
|
||||||
|
trap,
|
||||||
|
action,
|
||||||
|
stop: None,
|
||||||
|
exits: Default::default(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
enum Cast {
|
||||||
|
Note(&'static str),
|
||||||
|
StopNow,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl GenServer for Srv {
|
||||||
|
type Call = ();
|
||||||
|
type Reply = ();
|
||||||
|
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());
|
||||||
|
}
|
||||||
|
fn handle_call(&mut self, _: ()) {}
|
||||||
|
fn handle_cast(&mut self, c: Cast) {
|
||||||
|
match c {
|
||||||
|
Cast::Note(s) => self.log.push(s),
|
||||||
|
Cast::StopNow => self.stop.as_ref().unwrap().stop(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
fn handle_shutdown(&mut self) -> ShutdownAction {
|
||||||
|
self.log.push("handle_shutdown");
|
||||||
|
self.action
|
||||||
|
}
|
||||||
|
fn handle_exit(&mut self, sig: ExitSignal) {
|
||||||
|
self.log.push("handle_exit");
|
||||||
|
self.exits.lock().unwrap().push(sig);
|
||||||
|
}
|
||||||
|
fn terminate(&mut self) {
|
||||||
|
// Allowed to block on the graceful path.
|
||||||
|
if self.trap {
|
||||||
|
sleep(Duration::from_millis(10));
|
||||||
|
}
|
||||||
|
self.log.push("terminate");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn spawn_settled<G: GenServer>(state: G) -> GenServerRef<G> {
|
||||||
|
let r = start(state);
|
||||||
|
sleep(Duration::from_millis(20)); // let init (trap_exit) run
|
||||||
|
r
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn non_trapping_server_is_stopped_outright() {
|
||||||
|
let log = Log::default();
|
||||||
|
let l = log.clone();
|
||||||
|
run(move || {
|
||||||
|
let r = spawn_settled(Srv::new(&l, false, ShutdownAction::Exit));
|
||||||
|
let mon = monitor(r.pid());
|
||||||
|
request_shutdown(r.pid());
|
||||||
|
let d = mon.rx.recv().unwrap();
|
||||||
|
assert_eq!(d.reason, DownReason::Stopped);
|
||||||
|
});
|
||||||
|
assert_eq!(log.get(), vec!["terminate"]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn trapping_server_exits_normally_via_handle_shutdown() {
|
||||||
|
let log = Log::default();
|
||||||
|
let l = log.clone();
|
||||||
|
run(move || {
|
||||||
|
let r = spawn_settled(Srv::new(&l, true, ShutdownAction::Exit));
|
||||||
|
let mon = monitor(r.pid());
|
||||||
|
request_shutdown(r.pid());
|
||||||
|
let d = mon.rx.recv().unwrap();
|
||||||
|
assert_eq!(d.reason, DownReason::Exit);
|
||||||
|
});
|
||||||
|
assert_eq!(log.get(), vec!["handle_shutdown", "terminate"]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn continue_keeps_dispatching_until_stop_handle() {
|
||||||
|
let log = Log::default();
|
||||||
|
let l = log.clone();
|
||||||
|
run(move || {
|
||||||
|
let r = spawn_settled(Srv::new(&l, true, ShutdownAction::Continue));
|
||||||
|
let mon = monitor(r.pid());
|
||||||
|
request_shutdown(r.pid());
|
||||||
|
sleep(Duration::from_millis(20));
|
||||||
|
r.cast(Cast::Note("after-shutdown-request")).unwrap();
|
||||||
|
r.cast(Cast::StopNow).unwrap();
|
||||||
|
let d = mon.rx.recv().unwrap();
|
||||||
|
assert_eq!(d.reason, DownReason::Exit);
|
||||||
|
});
|
||||||
|
assert_eq!(
|
||||||
|
log.get(),
|
||||||
|
vec!["handle_shutdown", "after-shutdown-request", "terminate"]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn stop_handle_is_a_normal_exit_that_transient_does_not_restart() {
|
||||||
|
let starts = Arc::new(AtomicUsize::new(0));
|
||||||
|
let s = starts.clone();
|
||||||
|
run(move || {
|
||||||
|
let s2 = s.clone();
|
||||||
|
let sup = spawn(move || {
|
||||||
|
let s3 = s2.clone();
|
||||||
|
OneForOne::new()
|
||||||
|
.child(ChildSpec::new(Restart::Transient, move || {
|
||||||
|
s3.fetch_add(1, Ordering::SeqCst);
|
||||||
|
let log = Log::default();
|
||||||
|
let r = GenServerBuilder::new(Srv::new(&log, false, ShutdownAction::Exit))
|
||||||
|
.under(self_pid())
|
||||||
|
.start();
|
||||||
|
r.cast(Cast::StopNow).unwrap();
|
||||||
|
// Block until the server is gone; a bare spawn parent
|
||||||
|
// returning would not itself end the server.
|
||||||
|
let mon = monitor(r.pid());
|
||||||
|
let _ = mon.rx.recv();
|
||||||
|
}))
|
||||||
|
.run();
|
||||||
|
});
|
||||||
|
sup.join().unwrap(); // returns only if the child was not restarted forever
|
||||||
|
});
|
||||||
|
assert_eq!(starts.load(Ordering::SeqCst), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn linked_peer_death_reaches_handle_exit() {
|
||||||
|
let log = Log::default();
|
||||||
|
let l = log.clone();
|
||||||
|
let alive = Arc::new(AtomicBool::new(false));
|
||||||
|
let a = alive.clone();
|
||||||
|
run(move || {
|
||||||
|
let r = spawn_settled(Srv::new(&l, true, ShutdownAction::Exit));
|
||||||
|
let pid = r.pid();
|
||||||
|
let peer = spawn(move || {
|
||||||
|
link(pid);
|
||||||
|
panic!("peer dies");
|
||||||
|
});
|
||||||
|
let _ = peer.join();
|
||||||
|
sleep(Duration::from_millis(20));
|
||||||
|
r.cast(Cast::Note("still-serving")).unwrap();
|
||||||
|
sleep(Duration::from_millis(20));
|
||||||
|
a.store(true, Ordering::SeqCst);
|
||||||
|
let mon = monitor(r.pid());
|
||||||
|
drop(r); // inbox closes → clean exit
|
||||||
|
let _ = mon.rx.recv(); // don't let the root-exit sweep race terminate
|
||||||
|
});
|
||||||
|
assert!(alive.load(Ordering::SeqCst));
|
||||||
|
assert_eq!(log.get(), vec!["handle_exit", "still-serving", "terminate"]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn gen_server_ref_shutdown_is_graceful_for_a_trapping_server() {
|
||||||
|
let log = Log::default();
|
||||||
|
let l = log.clone();
|
||||||
|
run(move || {
|
||||||
|
let r = spawn_settled(Srv::new(&l, true, ShutdownAction::Exit));
|
||||||
|
r.shutdown();
|
||||||
|
});
|
||||||
|
assert_eq!(log.get(), vec!["handle_shutdown", "terminate"]);
|
||||||
|
}
|
||||||
@@ -0,0 +1,122 @@
|
|||||||
|
//! Graceful shutdown — `request_shutdown` (OTP `exit(Pid, shutdown)`).
|
||||||
|
//!
|
||||||
|
//! `request_stop` is `exit(Pid, kill)`: an uncatchable unwind at the target's
|
||||||
|
//! next observation point. `request_shutdown` is the polite form:
|
||||||
|
//! - a target that is NOT trapping exits is stopped exactly as by
|
||||||
|
//! `request_stop` (OTP's rule: don't trap, you die);
|
||||||
|
//! - a target that IS trapping receives an `ExitSignal { reason: Shutdown }`
|
||||||
|
//! on its trap inbox and keeps running — it is expected to wind down and
|
||||||
|
//! exit normally on its own.
|
||||||
|
|
||||||
|
use smarm::{monitor, request_shutdown, run, self_pid, sleep, spawn, trap_exit, DownReason, Pid};
|
||||||
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
|
use std::sync::{mpsc, Arc};
|
||||||
|
use std::thread;
|
||||||
|
use std::time::Duration;
|
||||||
|
|
||||||
|
const WATCHDOG: Duration = Duration::from_secs(10);
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn request_shutdown_stops_a_non_trapping_actor() {
|
||||||
|
run(|| {
|
||||||
|
let h = spawn(|| sleep(Duration::from_secs(3600)));
|
||||||
|
let mon = monitor(h.pid());
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
let down = mon.rx.recv().expect("down");
|
||||||
|
assert_eq!(down.reason, DownReason::Stopped);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn request_shutdown_is_a_message_to_a_trapping_actor() {
|
||||||
|
let unwound = Arc::new(AtomicBool::new(false));
|
||||||
|
let u = unwound.clone();
|
||||||
|
run(move || {
|
||||||
|
struct Unwound(Arc<AtomicBool>);
|
||||||
|
impl Drop for Unwound {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
if std::thread::panicking() {
|
||||||
|
self.0.store(true, Ordering::SeqCst);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let (tx, rx) = smarm::channel::<(Pid, DownReason)>();
|
||||||
|
let (ready_tx, ready_rx) = smarm::channel::<()>();
|
||||||
|
let h = spawn(move || {
|
||||||
|
let _g = Unwound(u);
|
||||||
|
let inbox = trap_exit();
|
||||||
|
let _ = ready_tx.send(());
|
||||||
|
let sig = inbox.recv().expect("exit signal");
|
||||||
|
let _ = tx.send((sig.from, sig.reason));
|
||||||
|
// Keep doing work after the request: shutdown is advisory.
|
||||||
|
sleep(Duration::from_millis(20));
|
||||||
|
});
|
||||||
|
// Trapping is set by the target itself; a request that beats it is a
|
||||||
|
// plain stop (same window as OTP's exit-before-process_flag).
|
||||||
|
ready_rx.recv().expect("ready");
|
||||||
|
let me = self_pid();
|
||||||
|
let mon = monitor(h.pid());
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
let (from, reason) = rx.recv().expect("relayed");
|
||||||
|
assert_eq!(from, me);
|
||||||
|
assert_eq!(reason, DownReason::Shutdown);
|
||||||
|
let down = mon.rx.recv().expect("down");
|
||||||
|
assert_eq!(
|
||||||
|
down.reason,
|
||||||
|
DownReason::Exit,
|
||||||
|
"target exited normally, not stopped"
|
||||||
|
);
|
||||||
|
});
|
||||||
|
assert!(
|
||||||
|
!unwound.load(Ordering::SeqCst),
|
||||||
|
"trapping target must not be unwound"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn request_shutdown_on_dead_pid_is_a_no_op() {
|
||||||
|
run(|| {
|
||||||
|
let h = spawn(|| {});
|
||||||
|
let pid = h.pid();
|
||||||
|
let _ = h.join();
|
||||||
|
request_shutdown(pid); // must not panic
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn handle_request_shutdown_from_foreign_thread() {
|
||||||
|
let rt = smarm::init(smarm::Config::exact(2));
|
||||||
|
let handle = rt.handle();
|
||||||
|
|
||||||
|
let (pid_tx, pid_rx) = mpsc::channel::<Pid>();
|
||||||
|
let requester = thread::spawn(move || {
|
||||||
|
let pid = pid_rx.recv().expect("pid");
|
||||||
|
thread::sleep(Duration::from_millis(50));
|
||||||
|
handle.request_shutdown(pid);
|
||||||
|
});
|
||||||
|
|
||||||
|
let (done_tx, done_rx) = mpsc::channel();
|
||||||
|
thread::spawn(move || {
|
||||||
|
rt.run(move || {
|
||||||
|
let (tx, rx) = smarm::channel::<DownReason>();
|
||||||
|
let (ready_tx, ready_rx) = smarm::channel::<()>();
|
||||||
|
let h = spawn(move || {
|
||||||
|
let inbox = trap_exit();
|
||||||
|
let _ = ready_tx.send(());
|
||||||
|
let sig = inbox.recv().expect("exit signal");
|
||||||
|
let _ = tx.send(sig.reason);
|
||||||
|
});
|
||||||
|
ready_rx.recv().expect("ready");
|
||||||
|
pid_tx.send(h.pid()).expect("send pid");
|
||||||
|
let reason = rx.recv().expect("relayed");
|
||||||
|
assert_eq!(reason, DownReason::Shutdown);
|
||||||
|
let _ = h.join();
|
||||||
|
});
|
||||||
|
let _ = done_tx.send(());
|
||||||
|
});
|
||||||
|
|
||||||
|
done_rx
|
||||||
|
.recv_timeout(WATCHDOG)
|
||||||
|
.expect("run did not return: foreign-thread request_shutdown never reached the target");
|
||||||
|
requester.join().expect("requester thread");
|
||||||
|
}
|
||||||
@@ -0,0 +1,303 @@
|
|||||||
|
//! Supervisor shutdown — the OTP child-spec `shutdown` policy.
|
||||||
|
//!
|
||||||
|
//! A supervisor traps exits. A `request_shutdown` reaching it (from its parent
|
||||||
|
//! supervisor, or from the app via `request_shutdown`/`RuntimeHandle`) runs
|
||||||
|
//! the ordered shutdown: children are stopped in reverse start order, each
|
||||||
|
//! per its `Shutdown` policy — `request_shutdown`, wait up to the timeout for
|
||||||
|
//! its termination signal, `request_stop` if it overstays — and then `run()`
|
||||||
|
//! returns normally. Every supervisor-initiated child stop (ordered shutdown,
|
||||||
|
//! OneForAll/RestForOne sibling cycling) goes through the same policy.
|
||||||
|
//!
|
||||||
|
//! A *hard* `request_stop` on a supervisor unwinds it; a drop guard then
|
||||||
|
//! hard-stops its live children so the subtree is never orphaned.
|
||||||
|
|
||||||
|
use smarm::supervisor::{ChildSpec, OneForOne, Restart, Shutdown, Strategy};
|
||||||
|
use smarm::{
|
||||||
|
monitor, request_shutdown, request_stop, run, sleep, spawn, trap_exit, DownReason, JoinHandle,
|
||||||
|
};
|
||||||
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
|
/// A child that traps exits, records the order it was shut down in, and exits
|
||||||
|
/// normally on the request (after `delay`). Ignores the request if `comply`
|
||||||
|
/// is false — a straggler that must be hard-stopped.
|
||||||
|
fn polite_child(
|
||||||
|
tag: usize,
|
||||||
|
log: &Arc<Mutex<Vec<usize>>>,
|
||||||
|
delay: Duration,
|
||||||
|
comply: bool,
|
||||||
|
) -> impl Fn() + Send + Sync + 'static {
|
||||||
|
let log = log.clone();
|
||||||
|
move || {
|
||||||
|
let inbox = trap_exit();
|
||||||
|
loop {
|
||||||
|
let sig = match inbox.recv() {
|
||||||
|
Ok(s) => s,
|
||||||
|
Err(_) => return,
|
||||||
|
};
|
||||||
|
if sig.reason == DownReason::Shutdown {
|
||||||
|
log.lock().unwrap().push(tag);
|
||||||
|
if comply {
|
||||||
|
sleep(delay);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
// Not complying: keep running until hard-stopped.
|
||||||
|
loop {
|
||||||
|
sleep(Duration::from_millis(5));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Spawn `sup`, let its children reach `trap_exit`, return the handle.
|
||||||
|
fn spawn_settled(sup: OneForOne) -> JoinHandle {
|
||||||
|
let h = spawn(move || sup.run());
|
||||||
|
sleep(Duration::from_millis(30));
|
||||||
|
h
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn shutdown_stops_children_in_reverse_order_and_returns_normally() {
|
||||||
|
let log = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let l = log.clone();
|
||||||
|
run(move || {
|
||||||
|
let sup = OneForOne::new()
|
||||||
|
.child(ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(1, &l, Duration::ZERO, true),
|
||||||
|
))
|
||||||
|
.child(ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(2, &l, Duration::ZERO, true),
|
||||||
|
))
|
||||||
|
.child(ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(3, &l, Duration::ZERO, true),
|
||||||
|
));
|
||||||
|
let h = spawn_settled(sup);
|
||||||
|
let mon = monitor(h.pid());
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
let down = mon.rx.recv().expect("down");
|
||||||
|
assert_eq!(
|
||||||
|
down.reason,
|
||||||
|
DownReason::Exit,
|
||||||
|
"supervisor exits normally after shutdown"
|
||||||
|
);
|
||||||
|
});
|
||||||
|
assert_eq!(*log.lock().unwrap(), vec![3, 2, 1]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn non_trapping_child_is_simply_stopped() {
|
||||||
|
let dropped = Arc::new(AtomicBool::new(false));
|
||||||
|
let d = dropped.clone();
|
||||||
|
run(move || {
|
||||||
|
struct G(Arc<AtomicBool>);
|
||||||
|
impl Drop for G {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
self.0.store(true, Ordering::SeqCst);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let sup = OneForOne::new().child(ChildSpec::new(Restart::Permanent, move || {
|
||||||
|
let _g = G(d.clone());
|
||||||
|
loop {
|
||||||
|
sleep(Duration::from_millis(5));
|
||||||
|
}
|
||||||
|
}));
|
||||||
|
let h = spawn_settled(sup);
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
h.join().expect("sup");
|
||||||
|
});
|
||||||
|
assert!(dropped.load(Ordering::SeqCst));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn straggler_is_hard_stopped_after_timeout() {
|
||||||
|
let log = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let l = log.clone();
|
||||||
|
run(move || {
|
||||||
|
let sup = OneForOne::new().child(
|
||||||
|
ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(1, &l, Duration::ZERO, false),
|
||||||
|
)
|
||||||
|
.shutdown(Shutdown::Timeout(Duration::from_millis(50))),
|
||||||
|
);
|
||||||
|
let h = spawn_settled(sup);
|
||||||
|
let t0 = Instant::now();
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
h.join().expect("sup");
|
||||||
|
let took = t0.elapsed();
|
||||||
|
assert!(
|
||||||
|
took >= Duration::from_millis(50),
|
||||||
|
"returned before the grace period: {took:?}"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
took < Duration::from_secs(2),
|
||||||
|
"did not fall back to a hard stop: {took:?}"
|
||||||
|
);
|
||||||
|
});
|
||||||
|
assert_eq!(
|
||||||
|
*log.lock().unwrap(),
|
||||||
|
vec![1],
|
||||||
|
"the straggler did receive the request"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn infinity_waits_for_a_slow_but_compliant_child() {
|
||||||
|
let log = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let l = log.clone();
|
||||||
|
let finished = Arc::new(AtomicBool::new(false));
|
||||||
|
let f = finished.clone();
|
||||||
|
run(move || {
|
||||||
|
let f2 = f.clone();
|
||||||
|
let l2 = l.clone();
|
||||||
|
let sup = OneForOne::new().child(
|
||||||
|
ChildSpec::new(Restart::Permanent, move || {
|
||||||
|
let inbox = trap_exit();
|
||||||
|
let _ = inbox.recv();
|
||||||
|
l2.lock().unwrap().push(1);
|
||||||
|
sleep(Duration::from_millis(150));
|
||||||
|
f2.store(true, Ordering::SeqCst); // only reached if not hard-stopped
|
||||||
|
})
|
||||||
|
.shutdown(Shutdown::Infinity),
|
||||||
|
);
|
||||||
|
let h = spawn_settled(sup);
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
h.join().expect("sup");
|
||||||
|
});
|
||||||
|
assert!(
|
||||||
|
finished.load(Ordering::SeqCst),
|
||||||
|
"Infinity must not hard-stop a compliant child"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn brutal_kill_skips_the_request() {
|
||||||
|
let log = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let l = log.clone();
|
||||||
|
run(move || {
|
||||||
|
let sup = OneForOne::new().child(
|
||||||
|
ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(1, &l, Duration::ZERO, true),
|
||||||
|
)
|
||||||
|
.shutdown(Shutdown::BrutalKill),
|
||||||
|
);
|
||||||
|
let h = spawn_settled(sup);
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
h.join().expect("sup");
|
||||||
|
});
|
||||||
|
assert!(
|
||||||
|
log.lock().unwrap().is_empty(),
|
||||||
|
"a BrutalKill child never sees the request"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn hard_stop_of_supervisor_does_not_orphan_children() {
|
||||||
|
let alive = Arc::new(AtomicUsize::new(0));
|
||||||
|
let a = alive.clone();
|
||||||
|
run(move || {
|
||||||
|
struct Alive(Arc<AtomicUsize>);
|
||||||
|
impl Drop for Alive {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
self.0.fetch_sub(1, Ordering::SeqCst);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let mk = |a: Arc<AtomicUsize>| {
|
||||||
|
move || {
|
||||||
|
a.fetch_add(1, Ordering::SeqCst);
|
||||||
|
let _g = Alive(a.clone());
|
||||||
|
loop {
|
||||||
|
sleep(Duration::from_millis(5));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let sup = OneForOne::new()
|
||||||
|
.child(ChildSpec::new(Restart::Permanent, mk(a.clone())))
|
||||||
|
.child(ChildSpec::new(Restart::Permanent, mk(a.clone())));
|
||||||
|
let h = spawn_settled(sup);
|
||||||
|
assert_eq!(a.load(Ordering::SeqCst), 2);
|
||||||
|
let mon = monitor(h.pid());
|
||||||
|
request_stop(h.pid());
|
||||||
|
let _ = mon.rx.recv();
|
||||||
|
sleep(Duration::from_millis(50));
|
||||||
|
assert_eq!(
|
||||||
|
a.load(Ordering::SeqCst),
|
||||||
|
0,
|
||||||
|
"children orphaned by a hard supervisor stop"
|
||||||
|
);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn nested_shutdown_reaches_grandchildren() {
|
||||||
|
let log = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let l = log.clone();
|
||||||
|
run(move || {
|
||||||
|
let l_inner = l.clone();
|
||||||
|
let inner = move || {
|
||||||
|
OneForOne::new()
|
||||||
|
.child(ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(10, &l_inner, Duration::ZERO, true),
|
||||||
|
))
|
||||||
|
.child(ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(11, &l_inner, Duration::ZERO, true),
|
||||||
|
))
|
||||||
|
.run()
|
||||||
|
};
|
||||||
|
let sup = OneForOne::new()
|
||||||
|
.child(ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(1, &l, Duration::ZERO, true),
|
||||||
|
))
|
||||||
|
.child(ChildSpec::new(Restart::Permanent, inner).shutdown(Shutdown::Infinity));
|
||||||
|
let h = spawn_settled(sup);
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
h.join().expect("sup");
|
||||||
|
});
|
||||||
|
assert_eq!(*log.lock().unwrap(), vec![11, 10, 1]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn sibling_cycling_uses_graceful_shutdown() {
|
||||||
|
// OneForAll: when child A dies, sibling B (trapping) must receive a
|
||||||
|
// Shutdown request rather than a bare stop.
|
||||||
|
let log = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let l = log.clone();
|
||||||
|
let a_runs = Arc::new(AtomicUsize::new(0));
|
||||||
|
let ar = a_runs.clone();
|
||||||
|
run(move || {
|
||||||
|
let ar2 = ar.clone();
|
||||||
|
let sup = OneForOne::new()
|
||||||
|
.strategy(Strategy::OneForAll)
|
||||||
|
.intensity(5, Duration::from_secs(60))
|
||||||
|
.child(ChildSpec::new(Restart::Transient, move || {
|
||||||
|
let n = ar2.fetch_add(1, Ordering::SeqCst) + 1;
|
||||||
|
sleep(Duration::from_millis(30));
|
||||||
|
if n == 1 {
|
||||||
|
panic!("first run dies");
|
||||||
|
}
|
||||||
|
// Second run: park until shut down.
|
||||||
|
let inbox = trap_exit();
|
||||||
|
let _ = inbox.recv();
|
||||||
|
}))
|
||||||
|
.child(ChildSpec::new(
|
||||||
|
Restart::Permanent,
|
||||||
|
polite_child(2, &l, Duration::ZERO, true),
|
||||||
|
));
|
||||||
|
let h = spawn(move || sup.run());
|
||||||
|
sleep(Duration::from_millis(150));
|
||||||
|
request_shutdown(h.pid());
|
||||||
|
h.join().expect("sup");
|
||||||
|
});
|
||||||
|
// B was shut down once by the cycle and once by the final shutdown.
|
||||||
|
assert_eq!(*log.lock().unwrap(), vec![2, 2]);
|
||||||
|
assert_eq!(a_runs.load(Ordering::SeqCst), 2);
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user