From 1002777ef37eb8fa19c6ab4bfd67352816ec6d0c Mon Sep 17 00:00:00 2001 From: "Claude (sandbox)" Date: Wed, 19 Aug 2026 05:50:54 +0000 Subject: [PATCH] feat(channel,runtime): off-runtime cross-thread wake for parked receivers and stop MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 (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. --- src/channel.rs | 55 +++++++++------- src/lib.rs | 2 +- src/runtime.rs | 55 +++++++++++++++- src/scheduler.rs | 25 +++++++- tests/cross_thread_wake.rs | 125 +++++++++++++++++++++++++++++++++++++ 5 files changed, 238 insertions(+), 24 deletions(-) create mode 100644 tests/cross_thread_wake.rs diff --git a/src/channel.rs b/src/channel.rs index a2af72a..db90425 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -90,8 +90,9 @@ use crate::pid::Pid; use crate::raw_mutex::RawMutex; +use crate::runtime::RuntimeInner; use std::collections::VecDeque; -use std::sync::Arc; +use std::sync::{Arc, Weak}; /// Create a new channel and return its `(Sender, Receiver)` halves. /// @@ -114,12 +115,15 @@ pub fn channel() -> (Sender, Receiver) { struct Inner { queue: VecDeque, - /// The parked receiver's `(pid, park-epoch)`, if one is currently + /// The parked receiver's `(pid, park-epoch, runtime)`, if one is currently /// waiting. The epoch identifies exactly which wait this is, so a waker /// left over from a wait that already ended (a losing `select` arm, a /// `recv_timeout` whose timer fired after it was already satisfied) is - /// inert and does nothing when it fires. - parked_receiver: Option<(Pid, u32)>, + /// inert and does nothing when it fires. The `Weak` is the + /// receiver's runtime, captured while it parked (so provably alive then); + /// it lets a sender on a foreign OS thread wake the receiver without the + /// `RUNTIME` thread-local, which is unset off a scheduler thread. + parked_receiver: Option<(Pid, u32, Weak)>, senders: usize, receiver_alive: bool, } @@ -206,8 +210,8 @@ impl Drop for Sender { None } }; - if let Some((pid, epoch)) = unpark { - crate::scheduler::unpark_at(pid, epoch); + if let Some((pid, epoch, rt)) = unpark { + crate::scheduler::unpark_at_via(pid, epoch, &rt); } } } @@ -254,13 +258,13 @@ impl Sender { g.queue.push_back(value); g.parked_receiver.take() }; - if let Some((pid, epoch)) = unpark { + if let Some((pid, epoch, rt)) = unpark { crate::te!(crate::trace::Event::Send { sender: crate::actor::current_pid() .unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)), receiver: Some(pid) }); - crate::scheduler::unpark_at(pid, epoch); + crate::scheduler::unpark_at_via(pid, epoch, &rt); } else { crate::te!(crate::trace::Event::Send { sender: crate::actor::current_pid() @@ -293,13 +297,17 @@ impl Receiver { None => panic!("smarm: recv() called outside an actor"), }; debug_assert!( - g.parked_receiver.is_none_or(|(p, _)| p == me), + g.parked_receiver.as_ref().is_none_or(|(p, _, _)| *p == me), "channel has more than one receiver" ); // begin_wait is lock-free, so it's legal under the Channel lock; // registering in the same critical section makes the epoch // atomic with the senders' view of the registration. - g.parked_receiver = Some((me, crate::scheduler::begin_wait())); + g.parked_receiver = Some(( + me, + crate::scheduler::begin_wait(), + crate::scheduler::runtime_weak(), + )); crate::te!(crate::trace::Event::RecvPark(me)); } // Release the lock before parking: the unparker will need it. @@ -347,11 +355,11 @@ impl Receiver { return Err(RecvTimeoutError::Disconnected); } debug_assert!( - g.parked_receiver.is_none_or(|(p, _)| p == me), + g.parked_receiver.as_ref().is_none_or(|(p, _, _)| *p == me), "channel has more than one receiver" ); epoch = crate::scheduler::begin_wait(); - g.parked_receiver = Some((me, epoch)); + g.parked_receiver = Some((me, epoch, crate::scheduler::runtime_weak())); crate::te!(crate::trace::Event::RecvPark(me)); } @@ -423,10 +431,14 @@ impl Receiver { None => panic!("smarm: recv_match() called outside an actor"), }; debug_assert!( - g.parked_receiver.is_none_or(|(p, _)| p == me), + g.parked_receiver.as_ref().is_none_or(|(p, _, _)| *p == me), "channel has more than one receiver" ); - g.parked_receiver = Some((me, crate::scheduler::begin_wait())); + g.parked_receiver = Some(( + me, + crate::scheduler::begin_wait(), + crate::scheduler::runtime_weak(), + )); crate::te!(crate::trace::Event::RecvPark(me)); } // Release the lock before parking: the unparker will need it. @@ -497,11 +509,12 @@ impl crate::timer::TimerTarget for RawMutex> { // keeps the registration bookkeeping exact.) let unpark = { let mut g = self.lock(); - if g.parked_receiver == Some((pid, epoch)) { - g.parked_receiver = None; - true - } else { - false + match g.parked_receiver { + Some((p, e, _)) if p == pid && e == epoch => { + g.parked_receiver = None; + true + } + _ => false, } }; // Unpark outside the channel lock: it may take the run-queue lock; @@ -562,10 +575,10 @@ impl Selectable for Receiver { return Ok(false); } debug_assert!( - g.parked_receiver.is_none_or(|(p, _)| p == pid), + g.parked_receiver.as_ref().is_none_or(|(p, _, _)| *p == pid), "channel has more than one receiver" ); - g.parked_receiver = Some((pid, epoch)); + g.parked_receiver = Some((pid, epoch, crate::scheduler::runtime_weak())); Ok(true) } diff --git a/src/lib.rs b/src/lib.rs index 7dfb3f4..1b4047f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -85,7 +85,7 @@ pub use registry::{ install, lookup_as, register, resolve_name, send, send_dyn, send_to, unregister, whereis, NameResolution, RegisterError, SendError, }; -pub use runtime::{init, Config, Runtime}; +pub use runtime::{init, Config, Runtime, RuntimeHandle}; pub use scheduler::{ block_on_io, cancel_timer, 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, diff --git a/src/runtime.rs b/src/runtime.rs index be56e5e..869a1d4 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -127,7 +127,7 @@ use crate::supervisor::Signal; use crate::timer::Timers; use std::sync::atomic::{AtomicBool, AtomicPtr, AtomicU32, AtomicU64, AtomicUsize, Ordering}; -use std::sync::{Arc, Mutex}; +use std::sync::{Arc, Mutex, Weak}; use std::thread; // --------------------------------------------------------------------------- @@ -1457,6 +1457,59 @@ impl Runtime { inner: self.inner.clone(), } } + + /// A `Send + Sync` handle to this runtime, usable from any thread — + /// including threads that are not smarm schedulers (an OS-signal handler + /// thread, an external event source). Grab it *before* [`run`](Self::run) + /// and hand it to e.g. a signal thread; that thread can then + /// [`request_stop`](RuntimeHandle::request_stop) the runtime's top + /// supervisor to drive an ordered shutdown from outside the runtime. + /// + /// The in-runtime primitives ([`scheduler::request_stop`](crate::request_stop) + /// and friends) reach the runtime through a thread-local that is unset on + /// any non-scheduler thread, so they are silent no-ops off-runtime; this + /// handle carries its own reference and closes that gap. + pub fn handle(&self) -> RuntimeHandle { + RuntimeHandle { + inner: Arc::downgrade(&self.inner), + } + } +} + +// --------------------------------------------------------------------------- +// RuntimeHandle — off-runtime wake/stop +// --------------------------------------------------------------------------- + +/// A `Send + Sync` handle to a [`Runtime`], obtained from +/// [`Runtime::handle`]. Lets a thread that is *not* a smarm scheduler thread +/// drive a cooperative stop into the runtime — the off-runtime counterpart to +/// [`scheduler::request_stop`](crate::request_stop). +/// +/// Holds a [`Weak`] to the runtime, for the same reason the IO backend does +/// (RFC 018): a lingering handle can never keep the runtime's slot table alive +/// and can never block [`Runtime::run`] from finishing. Once the `Runtime` is +/// dropped every method is a harmless no-op — the same end state as calling +/// `request_stop` on an actor that has already exited. +#[derive(Clone)] +pub struct RuntimeHandle { + inner: Weak, +} + +impl RuntimeHandle { + /// Ask `pid` to stop cooperatively, from any thread. The off-runtime + /// equivalent of [`scheduler::request_stop`](crate::request_stop): it sets + /// the target's stop flag and wakes it, so a parked actor unwinds at its + /// next checkpoint exactly as it would for an in-runtime stop. A no-op if + /// the runtime has been dropped, or if the actor has already exited. + pub fn request_stop(&self, pid: Pid) { + let pid = pid.erase(); + // Upgrade the Weak per call, like the IO backend does (io.rs): a live + // runtime yields the inner and we drive the same stop the in-runtime + // path would; a dropped runtime makes this a no-op. + if let Some(inner) = self.inner.upgrade() { + crate::scheduler::request_stop_inner(&inner, pid); + } + } } // --------------------------------------------------------------------------- diff --git a/src/scheduler.rs b/src/scheduler.rs index 13122e9..d41239a 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -71,7 +71,7 @@ use crate::pid::{Name, Pid}; use crate::runtime::{self, RuntimeInner, YieldIntent, RUNTIME}; use crate::supervisor::Signal; use std::sync::atomic::Ordering; -use std::sync::Arc; +use std::sync::{Arc, Weak}; // --------------------------------------------------------------------------- // with_runtime / try_with_runtime @@ -544,6 +544,29 @@ pub(crate) fn unpark_at(pid: Pid, epoch: u32) { let _ = try_with_runtime(|inner| inner.unpark_at(pid, epoch)); } +// The current actor's runtime as a `Weak`, for a waker that must reach the +// runtime from a foreign thread later. A channel captures this when its +// receiver parks, so a cross-thread `send` can wake without the `RUNTIME` +// thread-local (unset off a scheduler thread). Panics outside `Runtime::run()`, +// the same contract as `begin_wait`. +pub(crate) fn runtime_weak() -> Weak { + with_runtime(Arc::downgrade) +} + +// Epoch-matched wake of `pid` from a waker that may or may not be on a +// scheduler thread. On a scheduler thread we take the thread-local path +// (preemption-gated, slot-eligible); off one that path is a silent no-op, so +// we reach the runtime through `rt` — the `Weak` the waker captured while it +// was in-runtime. Mirrors the IO backend's cross-context wake (io.rs, RFC 018). +pub(crate) fn unpark_at_via(pid: Pid, epoch: u32, rt: &Weak) { + if try_with_runtime(|inner| inner.unpark_at(pid, epoch)).is_some() { + return; + } + if let Some(inner) = rt.upgrade() { + inner.unpark_at(pid, epoch); + } +} + // Open a new wait for the current actor and return its wait identity // ("epoch"). Call once per wait, before registering with any waker. Lock-free, // so it's legal to call while already holding another internal lock. diff --git a/tests/cross_thread_wake.rs b/tests/cross_thread_wake.rs new file mode 100644 index 0000000..f84221a --- /dev/null +++ b/tests/cross_thread_wake.rs @@ -0,0 +1,125 @@ +//! 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` 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() {} + +/// 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::(); + // 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::(); + + 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::(); + 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::(); + 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); +}