Compare commits

...
Author SHA1 Message Date
Carl Lerche ab1ce8706a num spawns counter 2023-08-17 10:56:45 -07:00
Carl Lerche 481e78e912 explicitly dump metrics on shutdown 2023-08-17 10:27:00 -07:00
Carl Lerche 63856b7bc1 add some more rt metrics 2023-08-17 10:01:44 -07:00
3 changed files with 54 additions and 4 deletions
@@ -2,6 +2,8 @@
mod imp {
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering::Relaxed;
use std::time::Duration;
use crate::runtime::WorkerMetrics;
static NUM_MAINTENANCE: AtomicUsize = AtomicUsize::new(0);
static NUM_NOTIFY_LOCAL: AtomicUsize = AtomicUsize::new(0);
@@ -21,9 +23,12 @@ mod imp {
static NUM_RELAY_SEARCH: AtomicUsize = AtomicUsize::new(0);
static NUM_SPIN_STALL: AtomicUsize = AtomicUsize::new(0);
static NUM_NO_LOCAL_WORK: AtomicUsize = AtomicUsize::new(0);
static NUM_DRIVER_WAKE: AtomicUsize = AtomicUsize::new(0);
static NUM_DRIVER_WAKE_NO_CORE: AtomicUsize = AtomicUsize::new(0);
static NUM_SPAWNS: AtomicUsize = AtomicUsize::new(0);
impl Drop for super::Counters {
fn drop(&mut self) {
impl super::Counters {
pub(crate) fn dump(&self, worker_metrics: &[WorkerMetrics]) {
let notifies_local = NUM_NOTIFY_LOCAL.load(Relaxed);
let notifies_remote = NUM_NOTIFY_REMOTE.load(Relaxed);
let unparks_local = NUM_UNPARKS_LOCAL.load(Relaxed);
@@ -42,6 +47,9 @@ mod imp {
let num_relay_search = NUM_RELAY_SEARCH.load(Relaxed);
let num_spin_stall = NUM_SPIN_STALL.load(Relaxed);
let num_no_local_work = NUM_NO_LOCAL_WORK.load(Relaxed);
let num_driver_wake = NUM_DRIVER_WAKE.load(Relaxed);
let num_driver_wake_no_core = NUM_DRIVER_WAKE_NO_CORE.load(Relaxed);
let num_spawns = NUM_SPAWNS.load(Relaxed);
println!("---");
println!("notifies (remote): {}", notifies_remote);
@@ -62,6 +70,19 @@ mod imp {
println!(" relay search: {}", num_relay_search);
println!(" spin stall: {}", num_spin_stall);
println!(" no local work: {}", num_no_local_work);
println!(" driver wakes: {}", num_driver_wake);
println!(" (no core): {}", num_driver_wake_no_core);
println!(" num spawns: {}", num_spawns);
println!("");
println!("worker metrics:");
for (i, worker) in worker_metrics.iter().enumerate() {
let mean_poll_time = Duration::from_nanos(worker.mean_poll_time.load(Relaxed));
println!("");
println!("{}:", i);
println!(" mean poll time: {:?}", mean_poll_time);
}
}
}
@@ -136,10 +157,24 @@ mod imp {
pub(crate) fn inc_num_no_local_work() {
NUM_NO_LOCAL_WORK.fetch_add(1, Relaxed);
}
pub(crate) fn inc_num_driver_wakes() {
NUM_DRIVER_WAKE.fetch_add(1, Relaxed);
}
pub(crate) fn inc_num_driver_wakes_no_core() {
NUM_DRIVER_WAKE_NO_CORE.fetch_add(1, Relaxed);
}
pub(crate) fn inc_num_spawns() {
NUM_SPAWNS.fetch_add(1, Relaxed);
}
}
#[cfg(not(tokio_internal_mt_counters))]
mod imp {
use crate::runtime::WorkerMetrics;
pub(crate) fn inc_num_inc_notify_local() {}
pub(crate) fn inc_num_notify_remote() {}
pub(crate) fn inc_num_unparks_local() {}
@@ -158,6 +193,13 @@ mod imp {
pub(crate) fn inc_num_relay_search() {}
pub(crate) fn inc_num_spin_stall() {}
pub(crate) fn inc_num_no_local_work() {}
pub(crate) fn inc_num_driver_wakes() {}
pub(crate) fn inc_num_driver_wakes_no_core() {}
pub(crate) fn inc_num_spawns() {}
impl super::Counters {
pub(crate) fn dump(&self, _: &[WorkerMetrics]) {}
}
}
#[derive(Debug)]
@@ -50,6 +50,8 @@ impl Handle {
{
let (handle, notified) = me.shared.owned.bind(future, me.clone(), id);
super::counters::inc_num_spawns();
if let Some(notified) = notified {
me.shared.schedule_task(notified, false);
}
@@ -175,7 +175,7 @@ pub(crate) struct Shared {
/// runtime metrics that can be useful when doing performance
/// investigations. This does nothing (empty struct, no drop impl) unless
/// the `tokio_internal_mt_counters` cfg flag is set.
_counters: Counters,
counters: Counters,
}
/// Data synchronized by the scheduler mutex
@@ -321,7 +321,7 @@ pub(super) fn create(
config,
scheduler_metrics: SchedulerMetrics::new(),
worker_metrics: worker_metrics.into_boxed_slice(),
_counters: Counters,
counters: Counters,
},
driver: driver_handle,
blocking_spawner,
@@ -1215,6 +1215,7 @@ impl Worker {
if let Some(mut driver) = cx.shared().take_driver() {
// Wait for driver events
driver.park(&cx.handle.driver);
super::counters::inc_num_driver_wakes();
synced = cx.shared().synced.lock();
@@ -1233,6 +1234,8 @@ impl Worker {
// This may result in a task being run
self.schedule_deferred_with_core(cx, core, move || synced)
} else {
super::counters::inc_num_driver_wakes_no_core();
// Schedule any deferred tasks
self.schedule_deferred_without_core(cx, &mut synced);
@@ -1503,6 +1506,9 @@ impl Shared {
while let Some(task) = self.next_remote_task_synced(synced) {
drop(task);
}
// Drop counters if enabled
handle.shared.counters.dump(&handle.shared.worker_metrics);
}
}