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:
+45
-32
@@ -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
@@ -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);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user