The RFC 014 root-exit sweep hard-stopped every live slot once nothing was runnable. That deferral privileged queued work over parked-with-a-pending- wake work (a sleeper was killed, a queued cast was drained) and any attempt to widen the notion of pending wake (timers, fd readiness) re-wedges the run on the periodic-timer daemon the sweep exists to end. Root exit now means "the program is done": finalize_actor delivers request_shutdown to every forest root — each live actor whose parent is the run (ROOT_PID) or is dead — synchronously, before the live-count decrement. Supervisors cascade per child Shutdown policy; trapping actors may Continue/drain with working timers and end the run when they stop themselves; non-trapping actors are stopped outright. No forcing sweep. Removes root_exited/root_swept, Pop::RootDrain and the idle-verdict condition; adds tests/root_exit.rs.
289 lines
10 KiB
Rust
289 lines
10 KiB
Rust
//! Root exit — the run's initial actor returning means "the program is done".
|
|
//!
|
|
//! When the root finalizes, the runtime delivers `request_shutdown` to every
|
|
//! **forest root**: each live actor whose parent is the run itself (a plain
|
|
//! `spawn` from the root closure) or is already dead. Nothing below a live
|
|
//! parent is touched directly — a supervisor gets one request and runs its
|
|
//! own ordered shutdown per child `Shutdown` policy.
|
|
//!
|
|
//! - Non-trapping actors are stopped outright, exactly as by
|
|
//! `request_shutdown` — a `spawn(|| { sleep(..); work() })` the root did
|
|
//! not `join` does NOT get to finish. Join it, supervise it, or trap.
|
|
//! - Trapping actors get `handle_shutdown` / an `ExitSignal{Shutdown}` and
|
|
//! may keep running (`Continue`, drain, then stop themselves) — timers and
|
|
//! all; the run ends when they do. There is no second, forcing sweep.
|
|
//! - A periodic-timer daemon (the classic wedge) never blocks `run()`.
|
|
|
|
use smarm::gen_server::{start, GenServer, GenServerCtx, ShutdownAction, StopHandle, TimerHandle};
|
|
use smarm::supervisor::{ChildSpec, OneForOne, Restart, Shutdown};
|
|
use smarm::{run, sleep, spawn, trap_exit, DownReason};
|
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
|
use std::sync::{Arc, Mutex};
|
|
use std::time::{Duration, Instant};
|
|
|
|
fn assert_prompt(start: Instant, what: &str) {
|
|
assert!(
|
|
start.elapsed() < Duration::from_secs(2),
|
|
"{what}: run() took {:?}",
|
|
start.elapsed()
|
|
);
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Bare (non-gen_server) actors
|
|
// ---------------------------------------------------------------------------
|
|
|
|
/// A non-trapping sleeper the root did not join is stopped, not waited for.
|
|
#[test]
|
|
fn unjoined_non_trapping_sleeper_is_stopped() {
|
|
let finished = Arc::new(AtomicBool::new(false));
|
|
let f = finished.clone();
|
|
let t = Instant::now();
|
|
run(move || {
|
|
spawn(move || {
|
|
sleep(Duration::from_secs(5));
|
|
f.store(true, Ordering::SeqCst);
|
|
});
|
|
});
|
|
assert_prompt(t, "sleeper");
|
|
assert!(
|
|
!finished.load(Ordering::SeqCst),
|
|
"sleeper should have been stopped"
|
|
);
|
|
}
|
|
|
|
/// A trapping bare actor sees `Shutdown` from the root's exit and may keep
|
|
/// working — here it sleeps (a timer!) after the signal, then returns. The
|
|
/// run waits for it: no forcing sweep.
|
|
#[test]
|
|
fn trapping_actor_may_finish_after_shutdown_signal() {
|
|
let finished = Arc::new(AtomicBool::new(false));
|
|
let f = finished.clone();
|
|
run(move || {
|
|
spawn(move || {
|
|
let inbox = trap_exit();
|
|
let sig = inbox.recv().expect("shutdown signal");
|
|
assert_eq!(sig.reason, DownReason::Shutdown);
|
|
sleep(Duration::from_millis(100));
|
|
f.store(true, Ordering::SeqCst);
|
|
});
|
|
// Trapping is a runtime opt-in: give the actor a chance to run
|
|
// `trap_exit()`; a not-yet-run actor is non-trapping and is stopped.
|
|
sleep(Duration::from_millis(20));
|
|
});
|
|
assert!(
|
|
finished.load(Ordering::SeqCst),
|
|
"trapping actor must be allowed to finish"
|
|
);
|
|
}
|
|
|
|
/// The classic wedge: a lazily spawned daemon that never returns on its own.
|
|
#[test]
|
|
fn parked_forever_daemon_does_not_block_run() {
|
|
let t = Instant::now();
|
|
run(|| {
|
|
let (_tx, rx) = smarm::channel::<()>();
|
|
spawn(move || {
|
|
let _ = rx.recv(); // parked forever: sender is held by the root, which returns
|
|
});
|
|
// Leak the sender into the daemon's own scope so nothing else drops it.
|
|
std::mem::forget(_tx);
|
|
});
|
|
assert_prompt(t, "daemon");
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// gen_servers
|
|
// ---------------------------------------------------------------------------
|
|
|
|
/// A non-trapping ticker with a periodic timer: the timer wheel is never empty,
|
|
/// and root exit must still end the run.
|
|
struct Ticker {
|
|
ticks: Arc<AtomicUsize>,
|
|
}
|
|
impl GenServer for Ticker {
|
|
type Call = ();
|
|
type Reply = ();
|
|
type Cast = ();
|
|
type Info = ();
|
|
type Timer = ();
|
|
fn init(&mut self, ctx: &GenServerCtx<Self>) {
|
|
ctx.timer().tick_every(Duration::from_millis(5), ());
|
|
}
|
|
fn handle_call(&mut self, _: ()) {}
|
|
fn handle_cast(&mut self, _: ()) {}
|
|
fn handle_timer(&mut self, _: ()) {
|
|
self.ticks.fetch_add(1, Ordering::SeqCst);
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn periodic_timer_daemon_does_not_block_run() {
|
|
let ticks = Arc::new(AtomicUsize::new(0));
|
|
let tk = ticks.clone();
|
|
let t = Instant::now();
|
|
run(move || {
|
|
let _r = start(Ticker { ticks: tk });
|
|
sleep(Duration::from_millis(50));
|
|
});
|
|
assert_prompt(t, "ticker");
|
|
assert!(
|
|
ticks.load(Ordering::SeqCst) >= 3,
|
|
"ticker should have ticked"
|
|
);
|
|
}
|
|
|
|
/// A trapping server that answers `Continue`, keeps ticking on its own timer
|
|
/// (draining), and stops itself later. Root exit must not cut it short.
|
|
struct Drainer {
|
|
log: Arc<Mutex<Vec<&'static str>>>,
|
|
shutdowns: Arc<AtomicUsize>,
|
|
ticks_after_shutdown: usize,
|
|
stop: Option<StopHandle<Drainer>>,
|
|
timer: Option<TimerHandle<Drainer>>,
|
|
draining: bool,
|
|
}
|
|
impl GenServer for Drainer {
|
|
type Call = ();
|
|
type Reply = ();
|
|
type Cast = ();
|
|
type Info = ();
|
|
type Timer = ();
|
|
fn init(&mut self, ctx: &GenServerCtx<Self>) {
|
|
ctx.trap_exit();
|
|
self.stop = Some(ctx.stop_handle());
|
|
self.timer = Some(ctx.timer());
|
|
}
|
|
fn handle_call(&mut self, _: ()) {}
|
|
fn handle_cast(&mut self, _: ()) {}
|
|
fn handle_shutdown(&mut self) -> ShutdownAction {
|
|
self.shutdowns.fetch_add(1, Ordering::SeqCst);
|
|
self.log.lock().unwrap().push("handle_shutdown");
|
|
self.draining = true;
|
|
self.timer
|
|
.as_ref()
|
|
.unwrap()
|
|
.tick_every(Duration::from_millis(10), ());
|
|
ShutdownAction::Continue
|
|
}
|
|
fn handle_timer(&mut self, _: ()) {
|
|
if !self.draining {
|
|
return;
|
|
}
|
|
self.ticks_after_shutdown += 1;
|
|
if self.ticks_after_shutdown == 3 {
|
|
self.log.lock().unwrap().push("drained");
|
|
self.stop.as_ref().unwrap().stop();
|
|
}
|
|
}
|
|
fn terminate(&mut self) {
|
|
self.log.lock().unwrap().push("terminate");
|
|
}
|
|
}
|
|
|
|
fn drainer(log: &Arc<Mutex<Vec<&'static str>>>, shutdowns: &Arc<AtomicUsize>) -> Drainer {
|
|
Drainer {
|
|
log: log.clone(),
|
|
shutdowns: shutdowns.clone(),
|
|
ticks_after_shutdown: 0,
|
|
stop: None,
|
|
timer: None,
|
|
draining: false,
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn trapping_server_drains_with_timers_after_root_exit() {
|
|
let log = Arc::new(Mutex::new(Vec::new()));
|
|
let shutdowns = Arc::new(AtomicUsize::new(0));
|
|
let (l, s) = (log.clone(), shutdowns.clone());
|
|
run(move || {
|
|
let r = start(drainer(&l, &s));
|
|
// A gen_server's lifetime is governed by its refs: dropping the last
|
|
// one closes the inbox and ends the loop cleanly, which would cut the
|
|
// drain short for a reason unrelated to root exit. Pin it the way a
|
|
// registered name would.
|
|
std::mem::forget(r);
|
|
sleep(Duration::from_millis(20)); // let init (trap_exit) run
|
|
});
|
|
assert_eq!(
|
|
*log.lock().unwrap(),
|
|
vec!["handle_shutdown", "drained", "terminate"]
|
|
);
|
|
assert_eq!(shutdowns.load(Ordering::SeqCst), 1);
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Supervision trees: only forest roots are addressed
|
|
// ---------------------------------------------------------------------------
|
|
|
|
/// The supervisor gets ONE request and runs its ordered shutdown; a trapping
|
|
/// child under it sees exactly one `Shutdown` — from the supervisor, not a
|
|
/// second one from the runtime — and is allowed to finish its drain (a sleep,
|
|
/// i.e. a timer) under `Shutdown::Infinity`.
|
|
#[test]
|
|
fn supervised_children_are_shut_down_only_via_their_supervisor() {
|
|
let signals = Arc::new(AtomicUsize::new(0));
|
|
let drained = Arc::new(AtomicBool::new(false));
|
|
let sup_returned = Arc::new(AtomicBool::new(false));
|
|
let (sg, dr, sr) = (signals.clone(), drained.clone(), sup_returned.clone());
|
|
run(move || {
|
|
spawn(move || {
|
|
OneForOne::new()
|
|
.child(
|
|
ChildSpec::new(Restart::Permanent, move || {
|
|
let inbox = trap_exit();
|
|
while let Ok(sig) = inbox.recv() {
|
|
if sig.reason == DownReason::Shutdown {
|
|
sg.fetch_add(1, Ordering::SeqCst);
|
|
sleep(Duration::from_millis(100));
|
|
// A late second signal would land here.
|
|
while let Ok(Some(sig)) = inbox.try_recv() {
|
|
if sig.reason == DownReason::Shutdown {
|
|
sg.fetch_add(1, Ordering::SeqCst);
|
|
}
|
|
}
|
|
dr.store(true, Ordering::SeqCst);
|
|
return;
|
|
}
|
|
}
|
|
})
|
|
.shutdown(Shutdown::Infinity),
|
|
)
|
|
.run();
|
|
sr.store(true, Ordering::SeqCst);
|
|
});
|
|
sleep(Duration::from_millis(30)); // let the tree settle
|
|
});
|
|
assert!(
|
|
sup_returned.load(Ordering::SeqCst),
|
|
"supervisor should return normally"
|
|
);
|
|
assert!(
|
|
drained.load(Ordering::SeqCst),
|
|
"child should finish its drain"
|
|
);
|
|
assert_eq!(signals.load(Ordering::SeqCst), 1);
|
|
}
|
|
|
|
/// A supervised non-trapping child under `Shutdown::Timeout` is stopped by
|
|
/// the supervisor's policy, and the run ends promptly.
|
|
#[test]
|
|
fn supervisor_tree_is_torn_down_promptly_on_root_exit() {
|
|
let t = Instant::now();
|
|
run(|| {
|
|
spawn(|| {
|
|
OneForOne::new()
|
|
.child(
|
|
ChildSpec::new(Restart::Permanent, || loop {
|
|
sleep(Duration::from_millis(5));
|
|
})
|
|
.shutdown(Shutdown::Timeout(Duration::from_millis(50))),
|
|
)
|
|
.run();
|
|
});
|
|
sleep(Duration::from_millis(30));
|
|
});
|
|
assert_prompt(t, "tree");
|
|
}
|