Wakes issued from a non-scheduler OS thread were silent no-ops. Every off-runtime wake primitive (unpark, unpark_at, request_stop) reaches the runtime through the RUNTIME thread-local, which is unset on any foreign thread — so a send from a plain std::thread enqueued its message but never woke the parked receiver, and there was no way to drive a stop into a runtime from an application thread (e.g. an OS-signal handler). The former strands a parked recv forever; the latter is why a downstream server must poll a shutdown flag instead of parking on it. Generalize RFC 018's rule — a producer reaches the runtime through a Weak it holds — from the IO backend to channel senders and to a new handle: - A receiver captures a Weak<RuntimeInner> (provably live at that moment) alongside its (pid, epoch) when it parks. send() and the last-sender drop wake through scheduler::unpark_at_via, which takes the thread-local path when on a scheduler thread (preemption-gated, slot-eligible) and the captured Weak otherwise — the same cross-context wake the IO threads do. - Runtime::handle() returns a Send + Sync RuntimeHandle carrying that Weak; RuntimeHandle::request_stop drives a cooperative stop from any thread and is a no-op once the runtime is dropped. The in-runtime wake paths (recv/select timers) are unchanged; only the sites reachable from a foreign thread route through the Weak. RuntimeHandle exposes request_stop only: send-wake needs no user-facing handle, and off-runtime unpark is covered because request_stop drives unpark on the upgraded inner. tests/cross_thread_wake.rs: a foreign-thread send wakes a parked receiver; a foreign-thread request_stop wakes and stops a parked actor; a RuntimeHandle held across and beyond run never blocks all-done and degrades to a no-op.
126 lines
4.9 KiB
Rust
126 lines
4.9 KiB
Rust
//! Cross-thread wake: a thread that is *not* a smarm scheduler thread must be
|
|
//! able to wake (and stop) a parked actor.
|
|
//!
|
|
//! The gap this pins down: every off-runtime wake primitive (`unpark`,
|
|
//! `unpark_at`, `request_stop`) reaches the runtime through the `RUNTIME`
|
|
//! thread-local, which is `None` on any non-scheduler thread — so a wake
|
|
//! issued from a foreign OS thread is a silent no-op and the parked actor
|
|
//! sleeps forever. Both failure modes below manifest as `Runtime::run` never
|
|
//! returning, so each test is wrapped in a watchdog: a timeout is the failure.
|
|
//!
|
|
//! The fix mirrors RFC 018's IO backend — the waker reaches the runtime
|
|
//! through a `Weak<RuntimeInner>` it already holds (the receiver captures one
|
|
//! when it parks; `Runtime::handle()` hands one to an app thread).
|
|
|
|
use std::sync::mpsc;
|
|
use std::thread;
|
|
use std::time::Duration;
|
|
|
|
const WATCHDOG: Duration = Duration::from_secs(10);
|
|
/// Give the target actor time to actually park before the foreign thread pokes
|
|
/// it, so we exercise the *wake* of a parked actor rather than the entry-side
|
|
/// stop check.
|
|
const SETTLE: Duration = Duration::from_millis(200);
|
|
|
|
fn assert_send_sync<T: Send + Sync>() {}
|
|
|
|
/// A cross-thread `send` from a plain OS thread must wake a receiver parked in
|
|
/// `recv`. Under the thread-local-only wake path the send enqueues the message
|
|
/// but never wakes the receiver, so `recv` — and therefore `run` — hangs.
|
|
#[test]
|
|
fn foreign_thread_send_wakes_parked_receiver() {
|
|
let (done_tx, done_rx) = mpsc::channel();
|
|
thread::spawn(move || {
|
|
let rt = smarm::init(smarm::Config::exact(2));
|
|
rt.run(|| {
|
|
let (tx, rx) = smarm::channel::<u32>();
|
|
// Receiver actor: parks on recv until the foreign thread sends.
|
|
let h = smarm::spawn(move || {
|
|
assert_eq!(rx.recv().expect("recv"), 42);
|
|
});
|
|
// Foreign (non-scheduler) OS thread owns the Sender and sends
|
|
// after the receiver has parked.
|
|
let sender = thread::spawn(move || {
|
|
thread::sleep(SETTLE);
|
|
tx.send(42).expect("send");
|
|
});
|
|
let _ = h.join();
|
|
sender.join().expect("sender thread");
|
|
});
|
|
let _ = done_tx.send(());
|
|
});
|
|
done_rx
|
|
.recv_timeout(WATCHDOG)
|
|
.expect("run did not return: a foreign-thread send never woke the parked receiver");
|
|
}
|
|
|
|
/// A cross-thread `request_stop` through a `RuntimeHandle` must wake and stop a
|
|
/// parked actor. The actor parks on a long sleep (only a stop can end it); the
|
|
/// handle is grabbed before `run` and driven from a foreign thread.
|
|
#[test]
|
|
fn foreign_thread_request_stop_wakes_parked_actor() {
|
|
assert_send_sync::<smarm::RuntimeHandle>();
|
|
|
|
let rt = smarm::init(smarm::Config::exact(2));
|
|
let handle = rt.handle();
|
|
|
|
// Foreign thread: learn the target pid from inside the run, let it park,
|
|
// then stop it through the handle.
|
|
let (pid_tx, pid_rx) = mpsc::channel::<smarm::Pid>();
|
|
let stopper = thread::spawn(move || {
|
|
let pid = pid_rx.recv().expect("pid");
|
|
thread::sleep(SETTLE);
|
|
handle.request_stop(pid);
|
|
});
|
|
|
|
let (done_tx, done_rx) = mpsc::channel();
|
|
thread::spawn(move || {
|
|
rt.run(move || {
|
|
let h = smarm::spawn(|| {
|
|
// Parks indefinitely; only a cooperative stop unwinds it.
|
|
smarm::sleep(Duration::from_secs(3600));
|
|
});
|
|
pid_tx.send(h.pid()).expect("send pid");
|
|
let _ = h.join();
|
|
});
|
|
let _ = done_tx.send(());
|
|
});
|
|
|
|
done_rx
|
|
.recv_timeout(WATCHDOG)
|
|
.expect("run did not return: a foreign-thread request_stop never woke the parked actor");
|
|
stopper.join().expect("stopper thread");
|
|
}
|
|
|
|
/// A `RuntimeHandle` held across (and beyond) a run must not keep the runtime
|
|
/// alive or block all-done: `run` still returns, and once the `Runtime` is
|
|
/// dropped the handle degrades to a harmless no-op (Weak lifecycle) rather than
|
|
/// panicking or touching freed memory.
|
|
#[test]
|
|
fn lingering_handle_does_not_block_all_done() {
|
|
let rt = smarm::init(smarm::Config::exact(1));
|
|
let handle = rt.handle(); // outlives the run below
|
|
|
|
let (pid_tx, pid_rx) = mpsc::channel::<smarm::Pid>();
|
|
let (done_tx, done_rx) = mpsc::channel();
|
|
let runner = thread::spawn(move || {
|
|
rt.run(move || {
|
|
let h = smarm::spawn(|| {});
|
|
pid_tx.send(h.pid()).expect("send pid");
|
|
let _ = h.join();
|
|
});
|
|
// `rt` is dropped here, at the end of this thread.
|
|
let _ = done_tx.send(());
|
|
});
|
|
|
|
done_rx
|
|
.recv_timeout(WATCHDOG)
|
|
.expect("run did not return while a RuntimeHandle was held live");
|
|
runner.join().expect("runner thread");
|
|
|
|
// Runtime is now dropped. A stop through the lingering handle must be a
|
|
// silent no-op, not a panic or use-after-free.
|
|
let dead_pid = pid_rx.recv().expect("pid");
|
|
handle.request_stop(dead_pid);
|
|
}
|