diag(runtime): wake-path counters in SchedulerStats + RuntimeStats::wake_diag; RQDIAG/DIAG lines in rq_runtime and general (target 5, finding 17)
This commit is contained in:
+19
-4
@@ -433,9 +433,19 @@ fn bench_pp_tokio_multi() -> (u64, u128) {
|
|||||||
// 5. spawn_pair_control — section 4 minus the messages
|
// 5. spawn_pair_control — section 4 minus the messages
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
/// Target-5 instrumentation: emit the runtime's wake-path counters for one
|
||||||
|
/// run when SMARM_WAKE_DIAG is set. Off by default so sweep.py output is
|
||||||
|
/// unchanged.
|
||||||
|
fn wake_diag(section: &str, threads: usize, rt: &smarm::runtime::Runtime, us: u128) {
|
||||||
|
if std::env::var_os("SMARM_WAKE_DIAG").is_some() {
|
||||||
|
println!("DIAG,{section},{threads},{us},{}", rt.stats().wake_diag());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn bench_ctl_smarm(threads: usize) -> (u64, u128) {
|
fn bench_ctl_smarm(threads: usize) -> (u64, u128) {
|
||||||
let start = Instant::now();
|
let start = Instant::now();
|
||||||
smarm::runtime::init(bench_cfg(threads)).run(|| {
|
let rt = smarm::runtime::init(bench_cfg(threads));
|
||||||
|
rt.run(|| {
|
||||||
for _ in 0..PP_ROUNDS {
|
for _ in 0..PP_ROUNDS {
|
||||||
let hb = smarm::spawn(|| {});
|
let hb = smarm::spawn(|| {});
|
||||||
let ha = smarm::spawn(|| {});
|
let ha = smarm::spawn(|| {});
|
||||||
@@ -443,7 +453,9 @@ fn bench_ctl_smarm(threads: usize) -> (u64, u128) {
|
|||||||
hb.join().unwrap();
|
hb.join().unwrap();
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
(PP_ROUNDS, start.elapsed().as_micros())
|
let us = start.elapsed().as_micros();
|
||||||
|
wake_diag("spawn_pair_control", threads, &rt, us);
|
||||||
|
(PP_ROUNDS, us)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn bench_ctl_tokio_current() -> (u64, u128) {
|
fn bench_ctl_tokio_current() -> (u64, u128) {
|
||||||
@@ -488,7 +500,8 @@ const PP_STEADY: u64 = 10_000;
|
|||||||
|
|
||||||
fn bench_steady_smarm(threads: usize) -> (u64, u128) {
|
fn bench_steady_smarm(threads: usize) -> (u64, u128) {
|
||||||
let start = Instant::now();
|
let start = Instant::now();
|
||||||
smarm::runtime::init(bench_cfg(threads)).run(|| {
|
let rt = smarm::runtime::init(bench_cfg(threads));
|
||||||
|
rt.run(|| {
|
||||||
let (tx_ab, rx_ab) = smarm::channel::<u64>();
|
let (tx_ab, rx_ab) = smarm::channel::<u64>();
|
||||||
let (tx_ba, rx_ba) = smarm::channel::<u64>();
|
let (tx_ba, rx_ba) = smarm::channel::<u64>();
|
||||||
let echo = smarm::spawn(move || {
|
let echo = smarm::spawn(move || {
|
||||||
@@ -504,7 +517,9 @@ fn bench_steady_smarm(threads: usize) -> (u64, u128) {
|
|||||||
}
|
}
|
||||||
echo.join().unwrap();
|
echo.join().unwrap();
|
||||||
});
|
});
|
||||||
(PP_STEADY, start.elapsed().as_micros())
|
let us = start.elapsed().as_micros();
|
||||||
|
wake_diag("ping_pong_steady", threads, &rt, us);
|
||||||
|
(PP_STEADY, us)
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn steady_tokio_body() {
|
async fn steady_tokio_body() {
|
||||||
|
|||||||
@@ -90,6 +90,7 @@ struct Sample {
|
|||||||
us: u128,
|
us: u128,
|
||||||
hits: u64,
|
hits: u64,
|
||||||
displacements: u64,
|
displacements: u64,
|
||||||
|
diag: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
fn yield_storm(threads: usize, slot: bool, actors: usize, yields: usize) -> Sample {
|
fn yield_storm(threads: usize, slot: bool, actors: usize, yields: usize) -> Sample {
|
||||||
@@ -116,6 +117,7 @@ fn yield_storm(threads: usize, slot: bool, actors: usize, yields: usize) -> Samp
|
|||||||
us,
|
us,
|
||||||
hits: stats.slot_hits(),
|
hits: stats.slot_hits(),
|
||||||
displacements: stats.slot_displacements(),
|
displacements: stats.slot_displacements(),
|
||||||
|
diag: stats.wake_diag(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -159,6 +161,7 @@ fn ping_pong_pairs(threads: usize, slot: bool, pairs: usize, roundtrips: usize)
|
|||||||
us,
|
us,
|
||||||
hits: stats.slot_hits(),
|
hits: stats.slot_hits(),
|
||||||
displacements: stats.slot_displacements(),
|
displacements: stats.slot_displacements(),
|
||||||
|
diag: stats.wake_diag(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -185,6 +188,7 @@ fn spawn_storm(threads: usize, slot: bool, spawns: usize) -> Sample {
|
|||||||
us,
|
us,
|
||||||
hits: stats.slot_hits(),
|
hits: stats.slot_hits(),
|
||||||
displacements: stats.slot_displacements(),
|
displacements: stats.slot_displacements(),
|
||||||
|
diag: stats.wake_diag(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -253,6 +257,7 @@ fn main() {
|
|||||||
mid.us,
|
mid.us,
|
||||||
per_s
|
per_s
|
||||||
);
|
);
|
||||||
|
println!("RQDIAG,{},{},{},{},{}", variant(), slot_str, name, t, mid.diag);
|
||||||
if slot {
|
if slot {
|
||||||
println!(
|
println!(
|
||||||
"RQSLOT,{},{},{},{},{}",
|
"RQSLOT,{},{},{},{},{}",
|
||||||
|
|||||||
+93
-6
@@ -356,6 +356,25 @@ pub struct SchedulerStats {
|
|||||||
pub slot_hits: AtomicU64,
|
pub slot_hits: AtomicU64,
|
||||||
/// RFC 005: slot occupants displaced to the shared queue by a newer wake.
|
/// RFC 005: slot occupants displaced to the shared queue by a newer wake.
|
||||||
pub slot_displacements: AtomicU64,
|
pub slot_displacements: AtomicU64,
|
||||||
|
// --- wake-path diagnostics (target 5). Relaxed counters, cheap. ---
|
||||||
|
/// unpark hit Parked from actor context → slot_push.
|
||||||
|
pub unpark_slot: AtomicU64,
|
||||||
|
/// unpark hit Parked from scheduler/foreign context → shared enqueue.
|
||||||
|
pub unpark_queue: AtomicU64,
|
||||||
|
/// unpark hit Running → RunningNotified (peer had not parked yet).
|
||||||
|
pub unpark_notified: AtomicU64,
|
||||||
|
/// park_return found the flag consumed → shared re-enqueue.
|
||||||
|
pub park_flag_consumed: AtomicU64,
|
||||||
|
/// Yield-intent re-enqueues.
|
||||||
|
pub yield_requeues: AtomicU64,
|
||||||
|
/// `enqueue` tail wake actually delivered a futex permit.
|
||||||
|
pub enqueue_wakes: AtomicU64,
|
||||||
|
/// Chain-rule wake actually delivered a permit.
|
||||||
|
pub chain_wakes: AtomicU64,
|
||||||
|
/// Times this scheduler entered Pop::Idle (about to futex-park).
|
||||||
|
pub idle_parks: AtomicU64,
|
||||||
|
/// Idle parks whose recheck aborted (WorkFound) — no futex.
|
||||||
|
pub idle_recheck_hits: AtomicU64,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl SchedulerStats {
|
impl SchedulerStats {
|
||||||
@@ -365,8 +384,42 @@ impl SchedulerStats {
|
|||||||
run_queue_len: AtomicU64::new(0),
|
run_queue_len: AtomicU64::new(0),
|
||||||
slot_hits: AtomicU64::new(0),
|
slot_hits: AtomicU64::new(0),
|
||||||
slot_displacements: AtomicU64::new(0),
|
slot_displacements: AtomicU64::new(0),
|
||||||
|
unpark_slot: AtomicU64::new(0),
|
||||||
|
unpark_queue: AtomicU64::new(0),
|
||||||
|
unpark_notified: AtomicU64::new(0),
|
||||||
|
park_flag_consumed: AtomicU64::new(0),
|
||||||
|
yield_requeues: AtomicU64::new(0),
|
||||||
|
enqueue_wakes: AtomicU64::new(0),
|
||||||
|
chain_wakes: AtomicU64::new(0),
|
||||||
|
idle_parks: AtomicU64::new(0),
|
||||||
|
idle_recheck_hits: AtomicU64::new(0),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn reset(&self) {
|
||||||
|
for c in [
|
||||||
|
&self.slot_hits,
|
||||||
|
&self.slot_displacements,
|
||||||
|
&self.unpark_slot,
|
||||||
|
&self.unpark_queue,
|
||||||
|
&self.unpark_notified,
|
||||||
|
&self.park_flag_consumed,
|
||||||
|
&self.yield_requeues,
|
||||||
|
&self.enqueue_wakes,
|
||||||
|
&self.chain_wakes,
|
||||||
|
&self.idle_parks,
|
||||||
|
&self.idle_recheck_hits,
|
||||||
|
] {
|
||||||
|
c.store(0, Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Bump a per-thread diagnostic counter on the calling scheduler thread.
|
||||||
|
macro_rules! diag {
|
||||||
|
($inner:expr, $field:ident) => {
|
||||||
|
SCHED_SLOT.with(|s| $inner.stats[s.get()].$field.fetch_add(1, Ordering::Relaxed))
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
@@ -422,6 +475,30 @@ impl RuntimeStats {
|
|||||||
.map(|s| s.slot_displacements.load(Ordering::Relaxed))
|
.map(|s| s.slot_displacements.load(Ordering::Relaxed))
|
||||||
.sum()
|
.sum()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Wake-path diagnostics summed across threads, as `name=value` pairs
|
||||||
|
/// (target 5 instrumentation). Reset at the start of each `run()`.
|
||||||
|
pub fn wake_diag(&self) -> String {
|
||||||
|
let sum = |f: fn(&SchedulerStats) -> &AtomicU64| -> u64 {
|
||||||
|
self.inner.stats.iter().map(|s| f(s).load(Ordering::Relaxed)).sum()
|
||||||
|
};
|
||||||
|
format!(
|
||||||
|
"slot_hits={} displaced={} unpark_slot={} unpark_queue={} unpark_notified={} \
|
||||||
|
flag_consumed={} yield_requeues={} enqueue_wakes={} chain_wakes={} \
|
||||||
|
idle_parks={} idle_recheck_hits={}",
|
||||||
|
sum(|s| &s.slot_hits),
|
||||||
|
sum(|s| &s.slot_displacements),
|
||||||
|
sum(|s| &s.unpark_slot),
|
||||||
|
sum(|s| &s.unpark_queue),
|
||||||
|
sum(|s| &s.unpark_notified),
|
||||||
|
sum(|s| &s.park_flag_consumed),
|
||||||
|
sum(|s| &s.yield_requeues),
|
||||||
|
sum(|s| &s.enqueue_wakes),
|
||||||
|
sum(|s| &s.chain_wakes),
|
||||||
|
sum(|s| &s.idle_parks),
|
||||||
|
sum(|s| &s.idle_recheck_hits),
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
@@ -1080,7 +1157,9 @@ impl RuntimeInner {
|
|||||||
// pure-compute hot path pays (almost) nothing. Bias is over-wake:
|
// pure-compute hot path pays (almost) nothing. Bias is over-wake:
|
||||||
// a spurious wake costs one futex round-trip and a failed pop; a
|
// a spurious wake costs one futex round-trip and a failed pop; a
|
||||||
// missed wake would cost a stranded actor.
|
// missed wake would cost a stranded actor.
|
||||||
self.coord.wake_one_if_idle();
|
if self.coord.wake_one_if_idle() {
|
||||||
|
diag!(self, enqueue_wakes);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Make `pid` runnable if it is parked; coalesce or defer otherwise.
|
/// Make `pid` runnable if it is parked; coalesce or defer otherwise.
|
||||||
@@ -1137,12 +1216,15 @@ impl RuntimeInner {
|
|||||||
// (install_actor enqueues directly), so they bypass the
|
// (install_actor enqueues directly), so they bypass the
|
||||||
// slot by construction.
|
// slot by construction.
|
||||||
if self.wake_slot && crate::actor::current_pid().is_some() {
|
if self.wake_slot && crate::actor::current_pid().is_some() {
|
||||||
|
diag!(self, unpark_slot);
|
||||||
self.slot_push(pid);
|
self.slot_push(pid);
|
||||||
} else {
|
} else {
|
||||||
|
diag!(self, unpark_queue);
|
||||||
self.enqueue(pid);
|
self.enqueue(pid);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Unpark::Notified => {
|
Unpark::Notified => {
|
||||||
|
diag!(self, unpark_notified);
|
||||||
crate::te!(crate::trace::Event::UnparkDeferred(pid));
|
crate::te!(crate::trace::Event::UnparkDeferred(pid));
|
||||||
}
|
}
|
||||||
Unpark::Noop => {}
|
Unpark::Noop => {}
|
||||||
@@ -1330,8 +1412,7 @@ impl Runtime {
|
|||||||
// RFC 005: slot counters reset at the START of a run (not the end),
|
// RFC 005: slot counters reset at the START of a run (not the end),
|
||||||
// so `stats()` read after `run()` returns reports that run's totals.
|
// so `stats()` read after `run()` returns reports that run's totals.
|
||||||
for stat in &self.inner.stats {
|
for stat in &self.inner.stats {
|
||||||
stat.slot_hits.store(0, Ordering::Relaxed);
|
stat.reset();
|
||||||
stat.slot_displacements.store(0, Ordering::Relaxed);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Spawn the initial actor through the public spawn path (which
|
// Spawn the initial actor through the public spawn path (which
|
||||||
@@ -2171,13 +2252,17 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
|||||||
// The mandatory post-publish re-check: a producer that
|
// The mandatory post-publish re-check: a producer that
|
||||||
// enqueued (or a verdict input that flipped) before it
|
// enqueued (or a verdict input that flipped) before it
|
||||||
// could see our idle bit has left us the evidence.
|
// could see our idle bit has left us the evidence.
|
||||||
let _ = inner.coord.park(slot_idx, tk_deadline, || {
|
stats.idle_parks.fetch_add(1, Ordering::Relaxed);
|
||||||
|
let pr = inner.coord.park(slot_idx, tk_deadline, || {
|
||||||
!inner.run_queue.is_empty()
|
!inner.run_queue.is_empty()
|
||||||
|| (inner.live_actors.load(Ordering::Acquire) == 0
|
|| (inner.live_actors.load(Ordering::Acquire) == 0
|
||||||
&& inner.io_outstanding.load(Ordering::Acquire) == 0
|
&& inner.io_outstanding.load(Ordering::Acquire) == 0
|
||||||
&& inner.io_fd_waiters.load(Ordering::Acquire) == 0)
|
&& inner.io_fd_waiters.load(Ordering::Acquire) == 0)
|
||||||
|| inner.coord.deadline_due()
|
|| inner.coord.deadline_due()
|
||||||
});
|
});
|
||||||
|
if matches!(pr, crate::park::ParkResult::WorkFound) {
|
||||||
|
stats.idle_recheck_hits.fetch_add(1, Ordering::Relaxed);
|
||||||
|
}
|
||||||
if tk_deadline.is_some() {
|
if tk_deadline.is_some() {
|
||||||
// Hand the role back BEFORE firing: pop_due can run
|
// Hand the role back BEFORE firing: pop_due can run
|
||||||
// `Send` thunks that insert new timers, and the
|
// `Send` thunks that insert new timers, and the
|
||||||
@@ -2212,8 +2297,8 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
|||||||
// miss here is safe — the enqueue that created the surplus already
|
// miss here is safe — the enqueue that created the surplus already
|
||||||
// issued its own wake (RFC 018 no-lost-wake); this only sharpens
|
// issued its own wake (RFC 018 no-lost-wake); this only sharpens
|
||||||
// parallelism latency.
|
// parallelism latency.
|
||||||
if !inner.run_queue.is_empty() {
|
if !inner.run_queue.is_empty() && inner.coord.wake_one_if_idle() {
|
||||||
inner.coord.wake_one_if_idle();
|
stats.chain_wakes.fetch_add(1, Ordering::Relaxed);
|
||||||
}
|
}
|
||||||
|
|
||||||
let slot = match inner.slot_at(pid) {
|
let slot = match inner.slot_at(pid) {
|
||||||
@@ -2302,6 +2387,7 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
|||||||
// arriving mid-run coalesces into the re-queue.
|
// arriving mid-run coalesces into the re-queue.
|
||||||
crate::te!(crate::trace::Event::Yield(pid));
|
crate::te!(crate::trace::Event::Yield(pid));
|
||||||
slot.word.yield_return(gen);
|
slot.word.yield_return(gen);
|
||||||
|
diag!(inner, yield_requeues);
|
||||||
inner.enqueue(pid);
|
inner.enqueue(pid);
|
||||||
}
|
}
|
||||||
YieldIntent::Park => {
|
YieldIntent::Park => {
|
||||||
@@ -2342,6 +2428,7 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
|
|||||||
#[cfg(feature = "smarm-causal")]
|
#[cfg(feature = "smarm-causal")]
|
||||||
crate::causal::on_deschedule(slot, false);
|
crate::causal::on_deschedule(slot, false);
|
||||||
crate::te!(crate::trace::Event::UnparkFlagConsumed(pid));
|
crate::te!(crate::trace::Event::UnparkFlagConsumed(pid));
|
||||||
|
diag!(inner, park_flag_consumed);
|
||||||
inner.enqueue(pid);
|
inner.enqueue(pid);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user