feat(runtime): root exit is graceful shutdown of the forest roots
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.
This commit is contained in:
+14
-4
@@ -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)
|
||||
|
||||
+55
-53
@@ -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<Timers>,
|
||||
/// 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<RuntimeInner>, 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<RuntimeInner>, 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<RuntimeInner>) {
|
||||
/// 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<RuntimeInner>, 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<RuntimeInner>, 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<RuntimeInner>, 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<RuntimeInner>, 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<RuntimeInner>, 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() {
|
||||
|
||||
+2
-1
@@ -617,7 +617,8 @@ pub fn request_stop<A>(pid: Pid<A>) {
|
||||
}
|
||||
|
||||
// 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.
|
||||
|
||||
@@ -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<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");
|
||||
}
|
||||
Reference in New Issue
Block a user