scheduler: spawn_monitor / spawn_monitor_with — monitor registered on the child's slot before publish, so the Down always carries the real reason (spawn-then-monitor could race to NoProc); tests; channel test uses it
This commit is contained in:
+3
-3
@@ -89,9 +89,9 @@ pub use runtime::{init, Config, Runtime, RuntimeHandle};
|
||||
pub use scheduler::{
|
||||
block_on_io, cancel_timer, request_shutdown, request_stop, run, self_pid, send_after,
|
||||
send_after_named, send_after_named_wall, send_after_wall, sleep, sleep_wall, spawn, spawn_addr,
|
||||
spawn_addr_with, spawn_under, spawn_under_with, spawn_with, try_spawn, try_spawn_under_with,
|
||||
wait_readable, wait_readable_timeout, wait_writable, wait_writable_timeout, yield_now, FdArm,
|
||||
JoinError, JoinHandle, SpawnError, SpawnOpts,
|
||||
spawn_addr_with, spawn_monitor, spawn_monitor_with, spawn_under, spawn_under_with, spawn_with,
|
||||
try_spawn, try_spawn_under_with, wait_readable, wait_readable_timeout, wait_writable,
|
||||
wait_writable_timeout, yield_now, FdArm, JoinError, JoinHandle, SpawnError, SpawnOpts,
|
||||
};
|
||||
pub use supervisor::{ChildSpec, OneForOne, Restart, Shutdown, Signal, Strategy};
|
||||
pub use timer::TimerId;
|
||||
|
||||
@@ -145,6 +145,11 @@ pub struct Monitor {
|
||||
/// Monitor `target`. Returns a [`Monitor`] whose `rx` receives exactly one
|
||||
/// [`Down`].
|
||||
///
|
||||
/// To monitor a child you are spawning yourself, prefer
|
||||
/// [`spawn_monitor`](crate::spawn_monitor): `spawn` followed by `monitor` on
|
||||
/// the returned pid can race the child's death and observe `NoProc` instead
|
||||
/// of its real reason.
|
||||
///
|
||||
/// If `target` is still live, the `Down` arrives when it terminates. If
|
||||
/// `target` is already gone, a [`DownReason::NoProc`] `Down` is queued
|
||||
/// immediately so the caller's `rx.recv()` returns without parking.
|
||||
|
||||
@@ -1764,6 +1764,7 @@ pub(crate) fn install_actor(
|
||||
stack: crate::stack::Stack,
|
||||
supervisor: Pid,
|
||||
closure: Closure,
|
||||
monitor: Option<(crate::monitor::MonitorId, crate::channel::Sender<crate::monitor::Down>)>,
|
||||
) -> Pid {
|
||||
let slot = &inner.slots[idx as usize];
|
||||
let gen = slot.generation(); // stable: we own the vacant slot via the free list
|
||||
@@ -1791,6 +1792,12 @@ pub(crate) fn install_actor(
|
||||
cold.outstanding_handles = 1;
|
||||
cold.outcome = None;
|
||||
cold.pending_io_result = None;
|
||||
// `spawn_monitor`: register before publish, so no scheduler can run
|
||||
// (and finalize) the child before its monitor exists. Same slot as a
|
||||
// `monitor()` registration; the send-from-finalize path is unchanged.
|
||||
if let Some(m) = monitor {
|
||||
cold.monitors.push(m);
|
||||
}
|
||||
}
|
||||
slot.sp.store(sp, Ordering::Relaxed);
|
||||
// RFC 019: a fresh incarnation starts with its high-water at the fresh
|
||||
|
||||
+47
-2
@@ -362,7 +362,7 @@ pub fn spawn_under_with<A>(
|
||||
|
||||
let pid = with_runtime(|inner| {
|
||||
let idx = inner.allocate_slot(); // panics loudly on slab exhaustion
|
||||
crate::runtime::install_actor(inner, idx, sp, stack, supervisor, closure)
|
||||
crate::runtime::install_actor(inner, idx, sp, stack, supervisor, closure, None)
|
||||
});
|
||||
|
||||
JoinHandle {
|
||||
@@ -371,6 +371,51 @@ pub fn spawn_under_with<A>(
|
||||
}
|
||||
}
|
||||
|
||||
/// [`spawn`] and [`monitor`](crate::monitor) the child in one step, with no
|
||||
/// window in which the child can die unobserved.
|
||||
///
|
||||
/// `spawn` followed by `monitor(h.pid())` races: on a multi-scheduler
|
||||
/// runtime the child can run to completion before the monitor registers, and
|
||||
/// a monitor on a dead pid delivers [`DownReason::NoProc`](crate::DownReason::NoProc)
|
||||
/// — the real reason (Exit vs Panic) is lost. Here the monitor is registered
|
||||
/// on the child's slot *before* the child is published to any run queue, so
|
||||
/// the `Down` always carries the child's actual termination reason. This is
|
||||
/// Erlang's `spawn_monitor/1`.
|
||||
pub fn spawn_monitor(f: impl FnOnce() + Send + 'static) -> (JoinHandle, crate::monitor::Monitor) {
|
||||
spawn_monitor_with(SpawnOpts::default(), f)
|
||||
}
|
||||
|
||||
/// [`spawn_monitor`] with per-actor stack shape overrides (RFC 019).
|
||||
pub fn spawn_monitor_with(
|
||||
opts: SpawnOpts,
|
||||
f: impl FnOnce() + Send + 'static,
|
||||
) -> (JoinHandle, crate::monitor::Monitor) {
|
||||
let parent = current_pid().unwrap_or_else(|| with_runtime(|_| crate::runtime::ROOT_PID));
|
||||
let (tx, rx) = crate::channel::channel::<crate::monitor::Down>();
|
||||
let stack = with_runtime(|inner| crate::runtime::acquire_stack(inner, opts));
|
||||
let sp = init_actor_stack(stack.top(), crate::actor::trampoline);
|
||||
let closure: crate::runtime::Closure = Box::new(f);
|
||||
|
||||
let (pid, id) = with_runtime(|inner| {
|
||||
let idx = inner.allocate_slot(); // panics loudly on slab exhaustion
|
||||
let id = inner.alloc_monitor_id();
|
||||
let pid = crate::runtime::install_actor(inner, idx, sp, stack, parent, closure, Some((id, tx)));
|
||||
(pid, id)
|
||||
});
|
||||
|
||||
(
|
||||
JoinHandle {
|
||||
pid,
|
||||
consumed: false,
|
||||
},
|
||||
crate::monitor::Monitor {
|
||||
id,
|
||||
target: pid,
|
||||
rx,
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
/// [`spawn`] that reports a full actor slab instead of panicking.
|
||||
///
|
||||
/// Behaviour parity with [`spawn`] in every case except one: when the fixed
|
||||
@@ -429,7 +474,7 @@ pub fn try_spawn_under_with<A>(
|
||||
|
||||
claimed.0 = None; // install_actor takes ownership of the slot from here
|
||||
let pid = with_runtime(|inner| {
|
||||
crate::runtime::install_actor(inner, idx, sp, stack, supervisor, closure)
|
||||
crate::runtime::install_actor(inner, idx, sp, stack, supervisor, closure, None)
|
||||
});
|
||||
|
||||
Ok(JoinHandle {
|
||||
|
||||
+5
-11
@@ -137,19 +137,13 @@ fn channel_ops_interleaved_with_monitor_churn_multi_thread() {
|
||||
for i in 0..32i64 {
|
||||
let tx = tx.clone();
|
||||
handles.push(spawn(move || {
|
||||
// Short-lived target whose death fires the monitor below. It
|
||||
// waits for `go` so the monitor is registered before it can
|
||||
// die — otherwise a fast target yields an immediate NoProc
|
||||
// Down instead of the finalize-sent Exit this test is about
|
||||
// (pre-existing ~8% flake at 4 threads, independent of the
|
||||
// wake slot).
|
||||
let (go_tx, go_rx) = channel::<()>();
|
||||
let t = spawn(move || {
|
||||
let _ = go_rx.recv();
|
||||
// Short-lived target whose death fires the monitor below.
|
||||
// spawn_monitor: registered before publish, so the Down is
|
||||
// the finalize-sent Exit this test is about, never NoProc
|
||||
// (spawn-then-monitor raced ~8% at 4 threads).
|
||||
let (t, m) = smarm::spawn_monitor(move || {
|
||||
tx.send(i).unwrap();
|
||||
});
|
||||
let m = smarm::monitor(t.pid());
|
||||
go_tx.send(()).unwrap();
|
||||
t.join().unwrap();
|
||||
// Down delivery exercises send-from-finalize.
|
||||
let d = m.rx.recv().unwrap();
|
||||
|
||||
@@ -161,3 +161,63 @@ fn demonitor_after_fire_is_none() {
|
||||
let _ = h.join();
|
||||
});
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// spawn_monitor: registration precedes publish, so a child that dies before
|
||||
// the parent gets another instruction in still reports its real reason.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
#[test]
|
||||
fn spawn_monitor_never_reports_noproc_multi_thread() {
|
||||
// 4 schedulers, 200 instantly-dying children. With spawn+monitor this
|
||||
// observes NoProc a few percent of the time (the child finishes on
|
||||
// another scheduler before monitor() registers); with spawn_monitor it
|
||||
// must be Exit every time.
|
||||
let noproc = Arc::new(AtomicUsize::new(0));
|
||||
let exit = Arc::new(AtomicUsize::new(0));
|
||||
let (n, e) = (noproc.clone(), exit.clone());
|
||||
smarm::init(smarm::Config::exact(4)).run(move || {
|
||||
let mut hs = Vec::new();
|
||||
for _ in 0..200 {
|
||||
let (h, m) = smarm::spawn_monitor(|| {});
|
||||
let d = m.rx.recv().expect("Down");
|
||||
assert_eq!(d.pid, h.pid());
|
||||
match d.reason {
|
||||
DownReason::Exit => e.fetch_add(1, Ordering::Relaxed),
|
||||
DownReason::NoProc => n.fetch_add(1, Ordering::Relaxed),
|
||||
other => panic!("unexpected {other:?}"),
|
||||
};
|
||||
hs.push(h);
|
||||
}
|
||||
for h in hs {
|
||||
h.join().unwrap();
|
||||
}
|
||||
});
|
||||
assert_eq!(noproc.load(Ordering::Relaxed), 0);
|
||||
assert_eq!(exit.load(Ordering::Relaxed), 200);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn spawn_monitor_sees_panic_and_demonitor_works() {
|
||||
let ok = Arc::new(AtomicBool::new(false));
|
||||
let o = ok.clone();
|
||||
run(move || {
|
||||
let (h, m) = smarm::spawn_monitor(|| panic!("boom"));
|
||||
let d = m.rx.recv().expect("Down");
|
||||
assert_eq!(d.pid, h.pid());
|
||||
assert!(matches!(d.reason, DownReason::Panic));
|
||||
let _ = h.join();
|
||||
|
||||
// demonitor before the child runs: no Down ever arrives.
|
||||
let (h2, m2) = smarm::spawn_monitor(|| {});
|
||||
demonitor(&m2);
|
||||
h2.join().unwrap();
|
||||
// Last sender gone with nothing sent: the channel is closed and empty.
|
||||
assert!(
|
||||
!matches!(m2.rx.try_recv(), Ok(Some(_))),
|
||||
"demonitored spawn_monitor still delivered"
|
||||
);
|
||||
o.store(true, Ordering::SeqCst);
|
||||
});
|
||||
assert!(ok.load(Ordering::SeqCst));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user