396 lines
16 KiB
Rust
396 lines
16 KiB
Rust
//! Supervision: keep a set of actors alive.
|
|
//!
|
|
//! A *supervisor* is an actor whose only job is to start a fixed set of child
|
|
//! actors and react when one of them terminates — restarting it (and, depending
|
|
//! on the strategy, some of its siblings) according to a policy, or giving up
|
|
//! when failures arrive too fast. It is how you turn "an actor that might crash"
|
|
//! into "a service that stays up": a crash becomes a restart instead of a hole
|
|
//! in the process tree.
|
|
//!
|
|
//! The supervisor type is [`OneForOne`]. The name is historical — the restart
|
|
//! *strategy* is selectable via [`OneForOne::strategy`], and
|
|
//! [`Strategy::OneForOne`] is merely the default. You declare the children up
|
|
//! front as [`ChildSpec`]s, each carrying a [`Restart`] policy, then hand the
|
|
//! supervision loop an actor of its own with [`OneForOne::run`].
|
|
//!
|
|
//! ## A child that crashes and recovers
|
|
//!
|
|
//! ```
|
|
//! use smarm::{run, spawn, ChildSpec, OneForOne, Restart};
|
|
//! use std::sync::Arc;
|
|
//! use std::sync::atomic::{AtomicUsize, Ordering};
|
|
//! use std::time::Duration;
|
|
//!
|
|
//! run(|| {
|
|
//! // A flaky child: it panics on its first two starts, then settles.
|
|
//! let starts = Arc::new(AtomicUsize::new(0));
|
|
//! let s = starts.clone();
|
|
//! let child = move || {
|
|
//! let n = s.fetch_add(1, Ordering::SeqCst) + 1;
|
|
//! if n < 3 {
|
|
//! panic!("boom {n}");
|
|
//! }
|
|
//! // The third start returns normally.
|
|
//! };
|
|
//!
|
|
//! // The supervisor runs on its own actor. `Transient` restarts a child
|
|
//! // that panics but treats a clean return as "done", so once the child
|
|
//! // finally succeeds the supervisor has nothing left to do and `run()`
|
|
//! // returns. smarm catches the child's panic and turns it into a restart;
|
|
//! // it never reaches the process as a real crash.
|
|
//! let sup = spawn(move || {
|
|
//! OneForOne::new()
|
|
//! .intensity(5, Duration::from_secs(60))
|
|
//! .child(ChildSpec::new(Restart::Transient, child))
|
|
//! .run();
|
|
//! });
|
|
//! sup.join().unwrap();
|
|
//!
|
|
//! assert_eq!(starts.load(Ordering::SeqCst), 3); // one start, two restarts
|
|
//! });
|
|
//! ```
|
|
//!
|
|
//! ## Restart policies
|
|
//!
|
|
//! Each child carries a [`Restart`] policy that decides whether *that child*
|
|
//! comes back when it terminates:
|
|
//!
|
|
//! - [`Restart::Permanent`] restarts on any termination, normal or panic —
|
|
//! for a service that should never be down.
|
|
//! - [`Restart::Transient`] restarts only on an abnormal exit (a panic or a
|
|
//! cooperative stop); a clean return means "done" — for work that runs to
|
|
//! completion but should be retried if it crashes.
|
|
//! - [`Restart::Temporary`] never restarts; the death is simply noted.
|
|
//!
|
|
//! ## Strategies: which siblings get cycled
|
|
//!
|
|
//! When a restart is due, the [`Strategy`] decides which *other* children are
|
|
//! cycled along with the one that died. The triggering child's own policy still
|
|
//! decides whether anything restarts at all.
|
|
//!
|
|
//! - [`Strategy::OneForOne`] restarts only the child that died — the default.
|
|
//! - [`Strategy::OneForAll`] restarts every child: the survivors are stopped,
|
|
//! then the whole set is restarted.
|
|
//! - [`Strategy::RestForOne`] restarts the dead child and every child started
|
|
//! after it, leaving earlier children untouched.
|
|
//!
|
|
//! ## Stopping a sibling is cooperative
|
|
//!
|
|
//! Cycling a sibling means stopping it first, and a supervisor never tears a
|
|
//! running actor down from outside: smarm actors share a heap and rely on
|
|
//! Drop/RAII, so unwinding a peer's stack from elsewhere would be unsound.
|
|
//! Instead the supervisor *requests* the stop and the child unwinds at its next
|
|
//! observation point — a `check!()`, an allocation, or a blocking call. A child
|
|
//! wedged in a tight loop with no observation point cannot be stopped, for the
|
|
//! same reason it cannot be preempted.
|
|
//!
|
|
//! ## Giving up: the restart-intensity cap
|
|
//!
|
|
//! A child that crashes the instant it starts would otherwise restart forever.
|
|
//! [`OneForOne::intensity`] bounds that: at most `max` restarts within any
|
|
//! `period`-long sliding window. One terminating child counts as a single
|
|
//! restart event even when the strategy cycles several siblings. When the cap
|
|
//! trips, the supervisor stops restarting, cooperatively stops any survivors in
|
|
//! reverse start order, and `run()` returns.
|
|
//!
|
|
//! ## Running context
|
|
//!
|
|
//! [`OneForOne::run`] takes over the calling actor as the supervision loop, so
|
|
//! a supervisor gets an actor of its own — typically
|
|
//! `spawn(|| OneForOne::new()/* … */.run())`, all from inside
|
|
//! [`run`](crate::run). Each child is spawned beneath the supervisor's pid, so
|
|
//! every child termination funnels back to it as a [`Signal`].
|
|
|
|
use crate::pid::Pid;
|
|
use std::any::Any;
|
|
|
|
/// A child-termination notice delivered to its supervisor.
|
|
///
|
|
/// Every child a supervisor starts is spawned beneath the supervisor's pid, so
|
|
/// each child's termination funnels back to it as one of these. The variant
|
|
/// records *how* the child went — which is what its [`Restart`] policy keys off.
|
|
pub enum Signal {
|
|
/// The child exited normally.
|
|
Exit(Pid),
|
|
/// The child panicked. Payload is whatever `panic!` was called with.
|
|
Panic(Pid, Box<dyn Any + Send>),
|
|
/// The child was cooperatively cancelled via `request_stop`. Carries no
|
|
/// payload — there is nothing to propagate. Kept distinct from `Exit` so a
|
|
/// supervisor can distinguish a stop it requested from a self-termination.
|
|
Stopped(Pid),
|
|
}
|
|
|
|
impl std::fmt::Debug for Signal {
|
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
|
match self {
|
|
Signal::Exit(pid) => write!(f, "Signal::Exit({:?})", pid),
|
|
Signal::Panic(pid, _) => write!(f, "Signal::Panic({:?}, ..)", pid),
|
|
Signal::Stopped(pid) => write!(f, "Signal::Stopped({:?})", pid),
|
|
}
|
|
}
|
|
}
|
|
|
|
impl Signal {
|
|
pub fn pid(&self) -> Pid {
|
|
match self {
|
|
Signal::Exit(p) => *p,
|
|
Signal::Panic(p, _) => *p,
|
|
Signal::Stopped(p) => *p,
|
|
}
|
|
}
|
|
}
|
|
|
|
use crate::channel::channel;
|
|
use std::collections::{HashMap, VecDeque};
|
|
use std::sync::Arc;
|
|
use std::time::{Duration, Instant};
|
|
|
|
/// When a terminated child should be restarted.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum Restart {
|
|
/// Always restart, on normal exit or panic.
|
|
Permanent,
|
|
/// Restart only on panic; a normal exit is treated as "done".
|
|
Transient,
|
|
/// Never restart.
|
|
Temporary,
|
|
}
|
|
|
|
/// A child managed by a supervisor: a restart policy plus a *factory* that
|
|
/// produces a fresh instance of the child's work on each (re)start.
|
|
///
|
|
/// The factory is `Fn` (not `FnOnce`) precisely so it can be called again to
|
|
/// restart; it is shared via `Arc` so the spec stays cheap to clone.
|
|
#[derive(Clone)]
|
|
pub struct ChildSpec {
|
|
start: Arc<dyn Fn() + Send + Sync + 'static>,
|
|
restart: Restart,
|
|
}
|
|
|
|
impl ChildSpec {
|
|
pub fn new(restart: Restart, start: impl Fn() + Send + Sync + 'static) -> Self {
|
|
Self { start: Arc::new(start), restart }
|
|
}
|
|
}
|
|
|
|
/// How a supervisor reacts when one child terminates and a restart is due.
|
|
///
|
|
/// The *triggering* child's [`Restart`] policy still decides whether a restart
|
|
/// happens at all; the strategy only decides *which other children* are cycled
|
|
/// along with it.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum Strategy {
|
|
/// Restart only the child that terminated. Siblings are untouched.
|
|
OneForOne,
|
|
/// Restart every child: the survivors are cooperatively stopped (in
|
|
/// reverse start order) and the whole set is restarted in start order.
|
|
OneForAll,
|
|
/// Restart the terminated child and every child started *after* it; those
|
|
/// started before it are untouched.
|
|
RestForOne,
|
|
}
|
|
|
|
/// A supervisor over a fixed set of children.
|
|
///
|
|
/// Build it with [`new`](Self::new), add children with [`child`](Self::child),
|
|
/// pick a [`strategy`](Self::strategy) and an [`intensity`](Self::intensity)
|
|
/// cap, then drive the loop with [`run`](Self::run) on an actor of its own. See
|
|
/// the [module docs](self) for the full picture.
|
|
pub struct OneForOne {
|
|
children: Vec<ChildSpec>,
|
|
strategy: Strategy,
|
|
intensity: u32,
|
|
period: Duration,
|
|
}
|
|
|
|
impl Default for OneForOne {
|
|
fn default() -> Self {
|
|
// Erlang's default supervisor intensity: 1 restart per 5 seconds. We
|
|
// start a little more permissive but in the same spirit.
|
|
Self {
|
|
children: Vec::new(),
|
|
strategy: Strategy::OneForOne,
|
|
intensity: 3,
|
|
period: Duration::from_secs(5),
|
|
}
|
|
}
|
|
}
|
|
|
|
impl OneForOne {
|
|
pub fn new() -> Self {
|
|
Self::default()
|
|
}
|
|
|
|
/// Select the restart strategy (default [`Strategy::OneForOne`]).
|
|
pub fn strategy(mut self, strategy: Strategy) -> Self {
|
|
self.strategy = strategy;
|
|
self
|
|
}
|
|
|
|
/// Allow at most `max` restarts within any `period`-long window before the
|
|
/// supervisor gives up and `run()` returns.
|
|
pub fn intensity(mut self, max: u32, period: Duration) -> Self {
|
|
self.intensity = max;
|
|
self.period = period;
|
|
self
|
|
}
|
|
|
|
pub fn child(mut self, spec: ChildSpec) -> Self {
|
|
self.children.push(spec);
|
|
self
|
|
}
|
|
|
|
/// Run the supervision loop on the current actor. Returns when every child
|
|
/// has reached a terminal, non-restartable state, or when the restart
|
|
/// intensity cap is tripped.
|
|
pub fn run(self) {
|
|
let me = crate::scheduler::self_pid();
|
|
let (tx, rx) = channel::<Signal>();
|
|
crate::scheduler::register_supervisor_channel(me, tx);
|
|
|
|
// pid -> index into `self.children`, for the children currently alive.
|
|
let mut by_pid: HashMap<Pid, usize> = HashMap::new();
|
|
let mut active: usize = 0;
|
|
// Sliding window of recent restart instants, for the intensity cap.
|
|
let mut restarts: Vec<Instant> = Vec::new();
|
|
|
|
let start_child = |idx: usize, by_pid: &mut HashMap<Pid, usize>| {
|
|
let start = self.children[idx].start.clone();
|
|
let h = crate::scheduler::spawn_under(me, move || (start)());
|
|
by_pid.insert(h.pid(), idx);
|
|
// We supervise via the signal funnel, not by joining; drop the
|
|
// handle so the child's slot is reclaimed promptly on death (the
|
|
// termination Signal is delivered before reclamation regardless).
|
|
drop(h);
|
|
};
|
|
|
|
for idx in 0..self.children.len() {
|
|
start_child(idx, &mut by_pid);
|
|
active += 1;
|
|
}
|
|
|
|
// A signal that arrives while we are awaiting stop-confirmations (for a
|
|
// child we are *not* currently stopping) is stashed here and processed
|
|
// by the main loop before it blocks on `recv` again.
|
|
let mut pending: VecDeque<Signal> = VecDeque::new();
|
|
let next_signal = |pending: &mut VecDeque<Signal>| -> Option<Signal> {
|
|
if let Some(s) = pending.pop_front() {
|
|
Some(s)
|
|
} else {
|
|
rx.recv().ok()
|
|
}
|
|
};
|
|
|
|
while active > 0 {
|
|
let sig = match next_signal(&mut pending) {
|
|
Some(s) => s,
|
|
None => break, // mailbox closed: nothing left to supervise
|
|
};
|
|
let idx = match by_pid.remove(&sig.pid()) {
|
|
Some(i) => i,
|
|
None => continue, // stray/duplicate signal
|
|
};
|
|
|
|
// The triggering child's own policy decides whether *anything*
|
|
// restarts. A cooperative stop counts as abnormal alongside a panic.
|
|
let abnormal = matches!(sig, Signal::Panic(..) | Signal::Stopped(..));
|
|
let should_restart = match self.children[idx].restart {
|
|
Restart::Permanent => true,
|
|
Restart::Transient => abnormal,
|
|
Restart::Temporary => false,
|
|
};
|
|
|
|
if !should_restart {
|
|
active -= 1;
|
|
continue;
|
|
}
|
|
|
|
// Intensity cap: prune the window, then give up if we are already
|
|
// at the limit. One triggering failure counts as one restart event,
|
|
// even when the strategy cycles several children. On giving up we
|
|
// fall out of the loop and the ordered shutdown below stops any
|
|
// survivors.
|
|
let now = Instant::now();
|
|
restarts.retain(|t| now.duration_since(*t) <= self.period);
|
|
if restarts.len() as u32 >= self.intensity {
|
|
break;
|
|
}
|
|
restarts.push(now);
|
|
|
|
// Which *live* siblings get cycled along with the failed child.
|
|
// (The failed child is already gone — removed from `by_pid` above.)
|
|
let mut to_stop: Vec<(Pid, usize)> = match self.strategy {
|
|
Strategy::OneForOne => Vec::new(),
|
|
Strategy::OneForAll => by_pid.iter().map(|(p, i)| (*p, *i)).collect(),
|
|
Strategy::RestForOne => by_pid
|
|
.iter()
|
|
.filter(|(_, i)| **i > idx)
|
|
.map(|(p, i)| (*p, *i))
|
|
.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
|
|
// are about to stop, restarted in start (ascending index) order.
|
|
let mut restart_set: Vec<usize> = Vec::with_capacity(to_stop.len() + 1);
|
|
restart_set.push(idx);
|
|
|
|
// Request stops, then await each survivor's termination signal
|
|
// before restarting. `request_stop` on an already-dead pid is a
|
|
// no-op; in that case its (already-sent) Exit signal serves as the
|
|
// confirmation. Any signal for a pid we are *not* awaiting is
|
|
// 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.
|
|
restart_set.sort_unstable();
|
|
for cidx in restart_set {
|
|
start_child(cidx, &mut by_pid);
|
|
}
|
|
}
|
|
|
|
// Ordered shutdown: stop any survivors in reverse start order and await
|
|
// their termination. On the normal `active == 0` exit `by_pid` is empty
|
|
// and this is a no-op; on a cap-trip or mailbox-closed break it tears
|
|
// the remaining children down deterministically instead of leaking them.
|
|
let mut survivors: Vec<(Pid, usize)> = by_pid.iter().map(|(p, i)| (*p, *i)).collect();
|
|
survivors.sort_unstable_by_key(|x| std::cmp::Reverse(x.1));
|
|
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,
|
|
None => break,
|
|
};
|
|
if let Some(pos) = awaiting.iter().position(|p| *p == s.pid()) {
|
|
awaiting.swap_remove(pos);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|