perf(channel): capture the receiver's runtime Weak once per channel, not once per park

Upstream's off-runtime wake (1002777) captured a Weak<RuntimeInner> in the
parked_receiver tuple on every park, so every park/unpark round-trip paid an
Arc::downgrade plus the matching Weak drop — a locked RMW pair on one globally
shared counter, on the hot path of every channel workload.

The receiver never migrates between runtimes, so one capture is enough: the
Weak moves into Inner<T>, filled in on first park (a branch on an Option
thereafter), and unpark_at_via takes a closure so the fallback only re-locks
and clones off a scheduler thread, where the in-runtime fast path has already
declined. runtime_weak() returns Option and no longer panics off-runtime.

Sandbox 1T channel steady round-trip: 290 -> 270 ns (upstream's shape measured
427 -> 445 on its own base). Full suite green under reltest, including
upstream's tests/cross_thread_wake.rs.
This commit is contained in:
claude-asm-audit
2026-08-21 12:30:18 +00:00
parent a9341e2d82
commit 93bc83a5b3
2 changed files with 64 additions and 40 deletions
+45 -32
View File
@@ -102,6 +102,7 @@ pub fn channel<T>() -> (Sender<T>, Receiver<T>) {
let inner = Arc::new(RawMutex::new_channel(Inner { let inner = Arc::new(RawMutex::new_channel(Inner {
queue: VecDeque::new(), queue: VecDeque::new(),
parked_receiver: None, parked_receiver: None,
rt: None,
senders: 1, senders: 1,
receiver_alive: true, receiver_alive: true,
})); }));
@@ -115,19 +116,36 @@ pub fn channel<T>() -> (Sender<T>, Receiver<T>) {
struct Inner<T> { struct Inner<T> {
queue: VecDeque<T>, queue: VecDeque<T>,
/// The parked receiver's `(pid, park-epoch, runtime)`, if one is currently /// The parked receiver's `(pid, park-epoch)`, if one is currently
/// waiting. The epoch identifies exactly which wait this is, so a waker /// 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 /// left over from a wait that already ended (a losing `select` arm, a
/// `recv_timeout` whose timer fired after it was already satisfied) is /// `recv_timeout` whose timer fired after it was already satisfied) is
/// inert and does nothing when it fires. The `Weak<RuntimeInner>` is the /// inert and does nothing when it fires.
/// receiver's runtime, captured while it parked (so provably alive then); parked_receiver: Option<(Pid, u32)>,
/// it lets a sender on a foreign OS thread wake the receiver without the /// The receiver's runtime, captured the first time it parks (so provably
/// `RUNTIME` thread-local, which is unset off a scheduler thread. /// alive then) and kept for the life of the channel: it lets a sender on a
parked_receiver: Option<(Pid, u32, Weak<RuntimeInner>)>, /// foreign OS thread wake the receiver without the `RUNTIME` thread-local,
/// which is unset off a scheduler thread. Captured once rather than per
/// park because `Arc::downgrade` + drop is a locked RMW pair on a shared
/// counter, and parking is the channel hot path. A `Receiver` never
/// migrates between runtimes — it is pinned to its actor — so one capture
/// stays correct for every later park.
rt: Option<Weak<RuntimeInner>>,
senders: usize, senders: usize,
receiver_alive: bool, receiver_alive: bool,
} }
impl<T> Inner<T> {
/// Capture the receiver's runtime if we have not already. Called under the
/// channel lock at every park site; after the first park it is one branch
/// on an `Option`, no atomics.
fn note_runtime(&mut self) {
if self.rt.is_none() {
self.rt = crate::scheduler::runtime_weak();
}
}
}
/// The sending half of a channel, created by [`channel`]. Clonable: every /// The sending half of a channel, created by [`channel`]. Clonable: every
/// clone pushes onto the same queue, and the channel stays open as long as /// clone pushes onto the same queue, and the channel stays open as long as
/// any clone is alive. Dropping the last `Sender` closes the channel, which /// any clone is alive. Dropping the last `Sender` closes the channel, which
@@ -210,8 +228,8 @@ impl<T> Drop for Sender<T> {
None None
} }
}; };
if let Some((pid, epoch, rt)) = unpark { if let Some((pid, epoch)) = unpark {
crate::scheduler::unpark_at_via(pid, epoch, &rt); crate::scheduler::unpark_at_via(pid, epoch, || self.inner.lock().rt.clone());
} }
} }
} }
@@ -258,13 +276,13 @@ impl<T> Sender<T> {
g.queue.push_back(value); g.queue.push_back(value);
g.parked_receiver.take() g.parked_receiver.take()
}; };
if let Some((pid, epoch, rt)) = unpark { if let Some((pid, epoch)) = unpark {
crate::te!(crate::trace::Event::Send { crate::te!(crate::trace::Event::Send {
sender: crate::actor::current_pid() sender: crate::actor::current_pid()
.unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)), .unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)),
receiver: Some(pid) receiver: Some(pid)
}); });
crate::scheduler::unpark_at_via(pid, epoch, &rt); crate::scheduler::unpark_at_via(pid, epoch, || self.inner.lock().rt.clone());
} else { } else {
crate::te!(crate::trace::Event::Send { crate::te!(crate::trace::Event::Send {
sender: crate::actor::current_pid() sender: crate::actor::current_pid()
@@ -297,17 +315,14 @@ impl<T> Receiver<T> {
None => panic!("smarm: recv() called outside an actor"), None => panic!("smarm: recv() called outside an actor"),
}; };
debug_assert!( debug_assert!(
g.parked_receiver.as_ref().is_none_or(|(p, _, _)| *p == me), g.parked_receiver.is_none_or(|(p, _)| p == me),
"channel has more than one receiver" "channel has more than one receiver"
); );
// begin_wait is lock-free, so it's legal under the Channel lock; // begin_wait is lock-free, so it's legal under the Channel lock;
// registering in the same critical section makes the epoch // registering in the same critical section makes the epoch
// atomic with the senders' view of the registration. // atomic with the senders' view of the registration.
g.parked_receiver = Some(( g.note_runtime();
me, g.parked_receiver = Some((me, crate::scheduler::begin_wait()));
crate::scheduler::begin_wait(),
crate::scheduler::runtime_weak(),
));
crate::te!(crate::trace::Event::RecvPark(me)); crate::te!(crate::trace::Event::RecvPark(me));
} }
// Release the lock before parking: the unparker will need it. // Release the lock before parking: the unparker will need it.
@@ -355,11 +370,12 @@ impl<T> Receiver<T> {
return Err(RecvTimeoutError::Disconnected); return Err(RecvTimeoutError::Disconnected);
} }
debug_assert!( debug_assert!(
g.parked_receiver.as_ref().is_none_or(|(p, _, _)| *p == me), g.parked_receiver.is_none_or(|(p, _)| p == me),
"channel has more than one receiver" "channel has more than one receiver"
); );
epoch = crate::scheduler::begin_wait(); epoch = crate::scheduler::begin_wait();
g.parked_receiver = Some((me, epoch, crate::scheduler::runtime_weak())); g.note_runtime();
g.parked_receiver = Some((me, epoch));
crate::te!(crate::trace::Event::RecvPark(me)); crate::te!(crate::trace::Event::RecvPark(me));
} }
@@ -431,14 +447,11 @@ impl<T> Receiver<T> {
None => panic!("smarm: recv_match() called outside an actor"), None => panic!("smarm: recv_match() called outside an actor"),
}; };
debug_assert!( debug_assert!(
g.parked_receiver.as_ref().is_none_or(|(p, _, _)| *p == me), g.parked_receiver.is_none_or(|(p, _)| p == me),
"channel has more than one receiver" "channel has more than one receiver"
); );
g.parked_receiver = Some(( g.note_runtime();
me, g.parked_receiver = Some((me, crate::scheduler::begin_wait()));
crate::scheduler::begin_wait(),
crate::scheduler::runtime_weak(),
));
crate::te!(crate::trace::Event::RecvPark(me)); crate::te!(crate::trace::Event::RecvPark(me));
} }
// Release the lock before parking: the unparker will need it. // Release the lock before parking: the unparker will need it.
@@ -509,12 +522,11 @@ impl<T: Send + 'static> crate::timer::TimerTarget for RawMutex<Inner<T>> {
// keeps the registration bookkeeping exact.) // keeps the registration bookkeeping exact.)
let unpark = { let unpark = {
let mut g = self.lock(); let mut g = self.lock();
match g.parked_receiver { if g.parked_receiver == Some((pid, epoch)) {
Some((p, e, _)) if p == pid && e == epoch => { g.parked_receiver = None;
g.parked_receiver = None; true
true } else {
} false
_ => false,
} }
}; };
// Unpark outside the channel lock: it may take the run-queue lock; // Unpark outside the channel lock: it may take the run-queue lock;
@@ -575,10 +587,11 @@ impl<T> Selectable for Receiver<T> {
return Ok(false); return Ok(false);
} }
debug_assert!( debug_assert!(
g.parked_receiver.as_ref().is_none_or(|(p, _, _)| *p == pid), g.parked_receiver.is_none_or(|(p, _)| p == pid),
"channel has more than one receiver" "channel has more than one receiver"
); );
g.parked_receiver = Some((pid, epoch, crate::scheduler::runtime_weak())); g.note_runtime();
g.parked_receiver = Some((pid, epoch));
Ok(true) Ok(true)
} }
+19 -8
View File
@@ -595,12 +595,16 @@ pub(crate) fn unpark_at(pid: Pid, epoch: u32) {
} }
// The current actor's runtime as a `Weak`, for a waker that must reach the // 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 // runtime from a foreign thread later. A channel captures this ONCE, the first
// receiver parks, so a cross-thread `send` can wake without the `RUNTIME` // time its receiver parks, so a cross-thread `send` can wake without the
// thread-local (unset off a scheduler thread). Panics outside `Runtime::run()`, // `RUNTIME` thread-local (unset off a scheduler thread). `None` off a
// the same contract as `begin_wait`. // scheduler thread, where there is nothing to capture.
pub(crate) fn runtime_weak() -> Weak<RuntimeInner> { //
with_runtime(Arc::downgrade) // Deliberately not called per park: `Arc::downgrade` plus the matching drop is
// a locked RMW pair on one globally shared counter, and the park/unpark
// round-trip is the hot path of every channel workload.
pub(crate) fn runtime_weak() -> Option<Weak<RuntimeInner>> {
try_with_runtime(Arc::downgrade)
} }
// Epoch-matched wake of `pid` from a waker that may or may not be on a // Epoch-matched wake of `pid` from a waker that may or may not be on a
@@ -608,11 +612,18 @@ pub(crate) fn runtime_weak() -> Weak<RuntimeInner> {
// (preemption-gated, slot-eligible); off one that path is a silent no-op, so // (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 // 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). // 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<RuntimeInner>) { pub(crate) fn unpark_at_via(
pid: Pid,
epoch: u32,
rt: impl FnOnce() -> Option<Weak<RuntimeInner>>,
) {
if try_with_runtime(|inner| inner.unpark_at(pid, epoch)).is_some() { if try_with_runtime(|inner| inner.unpark_at(pid, epoch)).is_some() {
return; return;
} }
if let Some(inner) = rt.upgrade() { // Off a scheduler thread only: `rt` re-takes the waker's lock to read the
// captured `Weak`, which is why it is a closure and not a value — the
// in-runtime path above must not pay for it.
if let Some(inner) = rt().as_ref().and_then(Weak::upgrade) {
inner.unpark_at(pid, epoch); inner.unpark_at(pid, epoch);
} }
} }