From 250f31265b8e53ce266b6c3ffa61fd4c04f9de6e Mon Sep 17 00:00:00 2001 From: "Claude (sandbox)" Date: Wed, 19 Aug 2026 07:11:46 +0000 Subject: [PATCH] feat(runtime): root exit is graceful shutdown of the forest roots MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- ROADMAP.md | 18 ++- src/runtime.rs | 108 ++++++++--------- src/scheduler.rs | 3 +- tests/root_exit.rs | 288 +++++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 359 insertions(+), 58 deletions(-) create mode 100644 tests/root_exit.rs diff --git a/ROADMAP.md b/ROADMAP.md index e2bea06..aa75071 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -77,10 +77,10 @@ Delivered surface: `ServerBuilder::start` untouched; free `call` / `cast` / `whereis_server`; `ServerRef::shutdown` + free `shutdown` as the sys-style synchronous stop. - **Root-exit teardown** (final phase): the run's initial actor is the root; - when it exits, the scheduler's idle verdict stops the parked-forever remainder - (deferred past the queue drain, so actors with in-flight work finish rather - than unwinding on the stop). Closes the "app actor blocks AllDone" stall — see - Look into, below. + when it exits the run winds down. *(Reworked with the graceful-shutdown work: + root exit now delivers `request_shutdown` to every forest root — see + "Root exit" below and `tests/root_exit.rs`.)* Closes the "app actor blocks + AllDone" stall — see Look into, below. Extends — does not retire — the "select exists; a unified per-process mailbox still does not" invariant: 014 adds addressable *delivery*, not a unified inbox; @@ -260,6 +260,16 @@ path the atomic-bool workaround stood in for. Re-check the urus crud repro to confirm the workaround can be retired (the teardown is cooperative — an actor in a tight loop with no observation point still can't be stopped). +**Update (graceful shutdown):** the RFC 014 sweep was a hard `request_stop` of +every live slot, deferred until nothing was runnable — which killed a sleeping +actor (timer pending) but drained a queued one, for no principled reason. It is +now the OTP semantics: root exit = "the program is done" = `request_shutdown` +to every **forest root** (live actor whose parent is the run or is dead), run +synchronously on the root's finalize path. Supervisors cascade with their child +`Shutdown` policies; trapping actors may `Continue`/drain (timers keep working) +and end the run when they stop themselves; non-trapping actors are stopped +outright — `join` what you need finished. No forcing sweep follows. + --- ## Invariants & gotchas (respect these across all cycles) diff --git a/src/runtime.rs b/src/runtime.rs index 9448722..16f45a3 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -904,17 +904,10 @@ pub(crate) struct RuntimeInner { pub(crate) live_actors: AtomicU32, /// Packed `(index << 32 | generation)` of the run's root (initial) actor, /// or `u64::MAX` (the ROOT_PID sentinel) before one is set. When this actor - /// finalizes it flags `root_exited`; the scheduler's idle verdict then - /// stops the remaining (parked-forever) actors. Set once per `run()`, right - /// after the initial spawn. + /// finalizes, `finalize_actor` runs the root-exit shutdown (see + /// `shutdown_forest_roots`). Set once per `run()`, right after the initial + /// spawn. pub(crate) root_bits: AtomicU64, - /// Set when the root actor finalizes; read by the scheduler's idle verdict - /// to trigger the one-shot teardown sweep. Reset per `run()`. - pub(crate) root_exited: AtomicBool, - /// Guards the teardown sweep to fire at most once per run (a parked-forever - /// remainder that survives the sweep falls through to the normal idle wait - /// rather than busy-spinning). Reset per `run()`. - pub(crate) root_swept: AtomicBool, /// Timer heap. Independent lock: never nested with any other. pub(crate) timers: Mutex, /// IO subsystem. `None` between runs. Lock order: io before everything. @@ -1004,8 +997,6 @@ impl RuntimeInner { free: RawMutex::new(free), live_actors: AtomicU32::new(0), root_bits: AtomicU64::new(u64::MAX), - root_exited: AtomicBool::new(false), - root_swept: AtomicBool::new(false), timers: Mutex::new(timers), io: Mutex::new(None), next_monitor_id: AtomicU64::new(0), @@ -1338,10 +1329,9 @@ impl Runtime { // requires a running runtime in the thread-local). RUNTIME.with(|r| *r.borrow_mut() = Some(self.inner.clone())); let initial_handle = crate::scheduler::spawn(f); - // The initial actor is the run's root: when it exits, remaining actors - // are stopped so the run winds down (see finalize_actor / schedule_loop). - self.inner.root_exited.store(false, Ordering::Relaxed); - self.inner.root_swept.store(false, Ordering::Relaxed); + // The initial actor is the run's root: its exit means "the program is + // done" — every remaining top-level actor is asked to shut down (see + // `finalize_actor` / `shutdown_forest_roots`). self.inner.set_root(initial_handle.pid()); // Launch N-1 extra scheduler threads, named `smarm-sched-{slot}` so @@ -1916,14 +1906,12 @@ fn finalize_actor(inner: &Arc, pid: Pid, outcome: Outcome) { // Reclaim if no outstanding handles (re-verified inside). reclaim_slot(inner, pid); - // Root-exit teardown is DEFERRED to the scheduler's idle verdict, not done - // here: stopping eagerly would cut off actors that still have queued work - // (they'd unwind on the stop before draining their mailbox). Flagging it - // instead lets the run queue drain naturally first; only the parked-forever - // remainder (e.g. a server pinned alive by a registered name) is then - // stopped, once nothing runnable is left. See `schedule_loop`. + // Root exit = the program is done. Ask every top-level survivor to shut + // down, right here, before the live-count decrement below: any wake this + // produces is then ordered before `live_actors` can be observed at its + // decremented value, same as every other wakeup finalize issues. if inner.is_root(pid) { - inner.root_exited.store(true, Ordering::Release); + shutdown_forest_roots(inner, pid); } // The decrement is LAST: every wakeup this finalize produced (joiners, @@ -1934,16 +1922,51 @@ fn finalize_actor(inner: &Arc, pid: Pid, outcome: Outcome) { debug_assert!(prev >= 1, "live_actors underflow — double finalize"); } -/// Cooperatively stop every live actor — the root-exit teardown sweep, run from -/// `schedule_loop` once the run queue is empty after the root has exited. Each -/// [`request_stop_inner`](crate::scheduler::request_stop_inner) re-verifies the -/// target under its cold lock, so the racy per-slot generation read is safe: a -/// vacant, dead, or reused slot no-ops. The swept actors unpark, unwind at their -/// next observation point, and finalize, dropping `live_actors` to zero. -fn stop_live_actors(inner: &Arc) { +/// The root-exit shutdown. Delivers [`request_shutdown`](crate::request_shutdown) +/// to every **forest root**: each live actor whose recorded parent +/// (`Actor::supervisor` — the spawner for a plain `spawn`, the supervisor for +/// `spawn_under`) is the run itself (`ROOT_PID`) or is no longer live. Actors +/// under a live parent are not addressed — that parent is responsible for +/// them: a supervisor traps and runs its ordered, policy-driven shutdown; a +/// bare parent that dies takes non-trapping children with it via the next +/// pass of this same rule only if it dies *now*, so a parent that outlives +/// this scan and later dies leaves its subtree to itself (Erlang semantics: an +/// unlinked spawn is nobody's child). +/// +/// Semantics per target follow `request_shutdown`: a trapping actor receives +/// `ExitSignal { from: root, reason: Shutdown }` and may finish work — drain, +/// keep its timers ticking, then stop itself; a non-trapping one is stopped +/// outright. There is no second, forcing sweep: an actor that traps and never +/// stops keeps the run alive by design (put it under a supervisor with a +/// `Shutdown::Timeout` policy if that is not wanted). Runs once, on the root's +/// finalize path, so it races only against actors that are still running — +/// each `request_shutdown_inner` re-verifies its target under the cold lock, +/// so a slot that dies or is reused mid-scan is a no-op. +fn shutdown_forest_roots(inner: &Arc, root: Pid) { for idx in 0..inner.slots.len() as u32 { - let pid = Pid::new(idx, inner.slots[idx as usize].generation()); - crate::scheduler::request_stop_inner(inner, pid); + let slot = &inner.slots[idx as usize]; + let pid = Pid::new(idx, slot.generation()); + if pid == root { + continue; + } + // Read the parent under the cold lock (generation-verified); act + // outside it — `request_shutdown_inner` sends and may unpark. + let parent = { + let cold = slot.cold.lock(); + if slot.generation() != pid.generation() { + continue; + } + match cold.actor.as_ref() { + Some(a) => a.supervisor, + None => continue, + } + }; + let parent_live = inner + .slot_at(parent) + .is_some_and(|ps| ps.is_live_for(parent)); + if !parent_live { + crate::scheduler::request_shutdown_inner(inner, pid, root); + } } } @@ -2027,9 +2050,6 @@ fn schedule_loop(inner: &Arc, slot_idx: usize) { Got(Pid), Idle, AllDone, - /// Root has exited and nothing is runnable: stop the parked-forever - /// remainder, then re-pop. Fires at most once per run. - RootDrain, } // 2a. RFC 005: drain this thread's wake slot before touching the @@ -2079,15 +2099,6 @@ fn schedule_loop(inner: &Arc, slot_idx: usize) { let live = inner.live_actors.load(Ordering::Acquire); if live == 0 && io_out == 0 { Pop::AllDone - } else if inner.root_exited.load(Ordering::Acquire) - && !inner.root_swept.swap(true, Ordering::AcqRel) - { - // Root gone and nothing runnable — the live remainder - // are parked-forever daemons (Queued actors with pending - // work drained before the queue emptied). Stop them so - // the run can end. One-shot: a survivor falls through to - // the idle wait below on the next pass. - Pop::RootDrain } else { Pop::Idle } @@ -2115,13 +2126,6 @@ fn schedule_loop(inner: &Arc, slot_idx: usize) { inner.coord.wake_all(); return; } - Pop::RootDrain => { - // Root has exited and nothing is runnable: stop the - // parked-forever remainder, then loop back to re-pop the - // now-runnable (stopping) actors. - stop_live_actors(inner); - continue; - } Pop::Idle => { // Something is still in flight. Park on our own futex // until a producer wakes us (enqueue tail), a deadline @@ -2152,8 +2156,6 @@ fn schedule_loop(inner: &Arc, slot_idx: usize) { || (inner.live_actors.load(Ordering::Acquire) == 0 && inner.io_outstanding.load(Ordering::Acquire) == 0 && inner.io_fd_waiters.load(Ordering::Acquire) == 0) - || (inner.root_exited.load(Ordering::Acquire) - && !inner.root_swept.load(Ordering::Acquire)) || inner.coord.deadline_due() }); if tk_deadline.is_some() { diff --git a/src/scheduler.rs b/src/scheduler.rs index a8b5f47..00d9125 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -617,7 +617,8 @@ pub fn request_stop(pid: Pid) { } // The core of `request_stop`, taking the runtime directly so it can also be -// driven from inside the runtime itself (the root-exit sweep) without +// driven from inside the runtime itself (the RuntimeHandle path, supervisor +// sweeps) without // re-borrowing the thread-local. Sets the stop flag under the target's lock // (a generation mismatch, or no live actor there, makes it a no-op) and // wakes the target. diff --git a/tests/root_exit.rs b/tests/root_exit.rs new file mode 100644 index 0000000..15ec16b --- /dev/null +++ b/tests/root_exit.rs @@ -0,0 +1,288 @@ +//! 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, +} +impl GenServer for Ticker { + type Call = (); + type Reply = (); + type Cast = (); + type Info = (); + type Timer = (); + fn init(&mut self, ctx: &GenServerCtx) { + 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>>, + shutdowns: Arc, + ticks_after_shutdown: usize, + stop: Option>, + timer: Option>, + draining: bool, +} +impl GenServer for Drainer { + type Call = (); + type Reply = (); + type Cast = (); + type Info = (); + type Timer = (); + fn init(&mut self, ctx: &GenServerCtx) { + 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>>, shutdowns: &Arc) -> 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"); +}