From 7001f04b650cf24b013c2a83f5d3aac8dc1c1e37 Mon Sep 17 00:00:00 2001 From: claude-asm-audit Date: Tue, 18 Aug 2026 07:04:31 +0000 Subject: [PATCH] diag(runtime): wake-path counters in SchedulerStats + RuntimeStats::wake_diag; RQDIAG/DIAG lines in rq_runtime and general (target 5, finding 17) --- benches/general.rs | 23 ++++++++-- benches/rq_runtime.rs | 5 +++ src/runtime.rs | 99 ++++++++++++++++++++++++++++++++++++++++--- 3 files changed, 117 insertions(+), 10 deletions(-) diff --git a/benches/general.rs b/benches/general.rs index 3f10d07..7b6187b 100644 --- a/benches/general.rs +++ b/benches/general.rs @@ -433,9 +433,19 @@ fn bench_pp_tokio_multi() -> (u64, u128) { // 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) { 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 { let hb = smarm::spawn(|| {}); let ha = smarm::spawn(|| {}); @@ -443,7 +453,9 @@ fn bench_ctl_smarm(threads: usize) -> (u64, u128) { 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) { @@ -488,7 +500,8 @@ const PP_STEADY: u64 = 10_000; fn bench_steady_smarm(threads: usize) -> (u64, u128) { 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::(); let (tx_ba, rx_ba) = smarm::channel::(); let echo = smarm::spawn(move || { @@ -504,7 +517,9 @@ fn bench_steady_smarm(threads: usize) -> (u64, u128) { } 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() { diff --git a/benches/rq_runtime.rs b/benches/rq_runtime.rs index 64f2cca..cc7fcbc 100644 --- a/benches/rq_runtime.rs +++ b/benches/rq_runtime.rs @@ -90,6 +90,7 @@ struct Sample { us: u128, hits: u64, displacements: u64, + diag: String, } 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, hits: stats.slot_hits(), 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, hits: stats.slot_hits(), displacements: stats.slot_displacements(), + diag: stats.wake_diag(), } } @@ -185,6 +188,7 @@ fn spawn_storm(threads: usize, slot: bool, spawns: usize) -> Sample { us, hits: stats.slot_hits(), displacements: stats.slot_displacements(), + diag: stats.wake_diag(), } } @@ -253,6 +257,7 @@ fn main() { mid.us, per_s ); + println!("RQDIAG,{},{},{},{},{}", variant(), slot_str, name, t, mid.diag); if slot { println!( "RQSLOT,{},{},{},{},{}", diff --git a/src/runtime.rs b/src/runtime.rs index fe57066..9b4d482 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -356,6 +356,25 @@ pub struct SchedulerStats { pub slot_hits: AtomicU64, /// RFC 005: slot occupants displaced to the shared queue by a newer wake. 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 { @@ -365,8 +384,42 @@ impl SchedulerStats { run_queue_len: AtomicU64::new(0), slot_hits: 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)) .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: // a spurious wake costs one futex round-trip and a failed pop; a // 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. @@ -1137,12 +1216,15 @@ impl RuntimeInner { // (install_actor enqueues directly), so they bypass the // slot by construction. if self.wake_slot && crate::actor::current_pid().is_some() { + diag!(self, unpark_slot); self.slot_push(pid); } else { + diag!(self, unpark_queue); self.enqueue(pid); } } Unpark::Notified => { + diag!(self, unpark_notified); crate::te!(crate::trace::Event::UnparkDeferred(pid)); } Unpark::Noop => {} @@ -1330,8 +1412,7 @@ impl Runtime { // 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. for stat in &self.inner.stats { - stat.slot_hits.store(0, Ordering::Relaxed); - stat.slot_displacements.store(0, Ordering::Relaxed); + stat.reset(); } // Spawn the initial actor through the public spawn path (which @@ -2171,13 +2252,17 @@ fn schedule_loop(inner: &Arc, slot_idx: usize) { // The mandatory post-publish re-check: a producer that // enqueued (or a verdict input that flipped) before it // 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.live_actors.load(Ordering::Acquire) == 0 && inner.io_outstanding.load(Ordering::Acquire) == 0 && inner.io_fd_waiters.load(Ordering::Acquire) == 0) || 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() { // Hand the role back BEFORE firing: pop_due can run // `Send` thunks that insert new timers, and the @@ -2212,8 +2297,8 @@ fn schedule_loop(inner: &Arc, slot_idx: usize) { // miss here is safe — the enqueue that created the surplus already // issued its own wake (RFC 018 no-lost-wake); this only sharpens // parallelism latency. - if !inner.run_queue.is_empty() { - inner.coord.wake_one_if_idle(); + if !inner.run_queue.is_empty() && inner.coord.wake_one_if_idle() { + stats.chain_wakes.fetch_add(1, Ordering::Relaxed); } let slot = match inner.slot_at(pid) { @@ -2302,6 +2387,7 @@ fn schedule_loop(inner: &Arc, slot_idx: usize) { // arriving mid-run coalesces into the re-queue. crate::te!(crate::trace::Event::Yield(pid)); slot.word.yield_return(gen); + diag!(inner, yield_requeues); inner.enqueue(pid); } YieldIntent::Park => { @@ -2342,6 +2428,7 @@ fn schedule_loop(inner: &Arc, slot_idx: usize) { #[cfg(feature = "smarm-causal")] crate::causal::on_deschedule(slot, false); crate::te!(crate::trace::Event::UnparkFlagConsumed(pid)); + diag!(inner, park_flag_consumed); inner.enqueue(pid); } }