From 8c076cb00d27c31b68071308b12412ca7cecceab Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Thu, 18 May 2023 13:28:47 -0700 Subject: [PATCH] rt: add internal counters to threaded runtime. (#5700) These counters are enabled using the `tokio_internal_mt_counters` and are intended to help with debugging performance issues. Whenever I work on the threaded runtime, I often find myself adding these counters, then removing them before submitting a PR. I think keeping them in will save time in the long run and shouldn't impact dev much. --- .../scheduler/multi_thread/counters.rs | 46 +++++++++++++++++++ .../src/runtime/scheduler/multi_thread/mod.rs | 3 ++ .../runtime/scheduler/multi_thread/worker.rs | 34 ++++++++++---- 3 files changed, 75 insertions(+), 8 deletions(-) create mode 100644 tokio/src/runtime/scheduler/multi_thread/counters.rs diff --git a/tokio/src/runtime/scheduler/multi_thread/counters.rs b/tokio/src/runtime/scheduler/multi_thread/counters.rs new file mode 100644 index 000000000..bfc7d6da2 --- /dev/null +++ b/tokio/src/runtime/scheduler/multi_thread/counters.rs @@ -0,0 +1,46 @@ +#[cfg(tokio_internal_mt_counters)] +mod imp { + use std::sync::atomic::AtomicUsize; + use std::sync::atomic::Ordering::Relaxed; + + static NUM_MAINTENANCE: AtomicUsize = AtomicUsize::new(0); + static NUM_NOTIFY_LOCAL: AtomicUsize = AtomicUsize::new(0); + static NUM_UNPARKS_LOCAL: AtomicUsize = AtomicUsize::new(0); + + impl Drop for super::Counters { + fn drop(&mut self) { + let notifies_local = NUM_NOTIFY_LOCAL.load(Relaxed); + let unparks_local = NUM_UNPARKS_LOCAL.load(Relaxed); + let maintenance = NUM_MAINTENANCE.load(Relaxed); + + println!("---"); + println!("notifies (local): {}", notifies_local); + println!(" unparks (local): {}", unparks_local); + println!(" maintenance: {}", maintenance); + } + } + + pub(crate) fn inc_num_inc_notify_local() { + NUM_NOTIFY_LOCAL.fetch_add(1, Relaxed); + } + + pub(crate) fn inc_num_unparks_local() { + NUM_UNPARKS_LOCAL.fetch_add(1, Relaxed); + } + + pub(crate) fn inc_num_maintenance() { + NUM_MAINTENANCE.fetch_add(1, Relaxed); + } +} + +#[cfg(not(tokio_internal_mt_counters))] +mod imp { + pub(crate) fn inc_num_inc_notify_local() {} + pub(crate) fn inc_num_unparks_local() {} + pub(crate) fn inc_num_maintenance() {} +} + +#[derive(Debug)] +pub(crate) struct Counters; + +pub(super) use imp::*; diff --git a/tokio/src/runtime/scheduler/multi_thread/mod.rs b/tokio/src/runtime/scheduler/multi_thread/mod.rs index 47cd1f3d7..0a1c63de9 100644 --- a/tokio/src/runtime/scheduler/multi_thread/mod.rs +++ b/tokio/src/runtime/scheduler/multi_thread/mod.rs @@ -1,5 +1,8 @@ //! Multi-threaded runtime +mod counters; +use counters::Counters; + mod handle; pub(crate) use handle::Handle; diff --git a/tokio/src/runtime/scheduler/multi_thread/worker.rs b/tokio/src/runtime/scheduler/multi_thread/worker.rs index c59f19e59..eb0ed705e 100644 --- a/tokio/src/runtime/scheduler/multi_thread/worker.rs +++ b/tokio/src/runtime/scheduler/multi_thread/worker.rs @@ -59,7 +59,7 @@ use crate::loom::sync::{Arc, Mutex}; use crate::runtime; use crate::runtime::context; -use crate::runtime::scheduler::multi_thread::{queue, Handle, Idle, Parker, Unparker}; +use crate::runtime::scheduler::multi_thread::{queue, Counters, Handle, Idle, Parker, Unparker}; use crate::runtime::task::{Inject, OwnedTasks}; use crate::runtime::{ blocking, coop, driver, scheduler, task, Config, MetricsBatch, SchedulerMetrics, WorkerMetrics, @@ -148,6 +148,12 @@ pub(super) struct Shared { pub(super) scheduler_metrics: SchedulerMetrics, pub(super) worker_metrics: Box<[WorkerMetrics]>, + + /// Only held to trigger some code on drop. This is used to get internal + /// 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, } /// Used to communicate with a worker from other threads. @@ -230,6 +236,7 @@ pub(super) fn create( config, scheduler_metrics: SchedulerMetrics::new(), worker_metrics: worker_metrics.into_boxed_slice(), + _counters: Counters, }, driver: driver_handle, blocking_spawner, @@ -512,6 +519,8 @@ impl Context { fn maintenance(&self, mut core: Box) -> Box { if core.tick % self.worker.handle.shared.config.event_interval == 0 { + super::counters::inc_num_maintenance(); + // Call `park` with a 0 timeout. This enables the I/O driver, timer, ... // to run without actually putting the thread to sleep. core = self.park_timeout(core, Some(Duration::from_millis(0))); @@ -585,7 +594,7 @@ impl Context { // If there are tasks available to steal, but this worker is not // looking for tasks to steal, notify another worker. if !core.is_searching && core.run_queue.is_stealable() { - self.worker.handle.notify_parked(); + self.worker.handle.notify_parked_local(); } core @@ -786,7 +795,7 @@ impl Handle { // Otherwise, use the inject queue. self.shared.inject.push(task); self.shared.scheduler_metrics.inc_remote_schedule_count(); - self.notify_parked(); + self.notify_parked_remote(); }) } @@ -820,7 +829,7 @@ impl Handle { // scheduling is from a resource driver. As notifications often come in // batches, the notification is delayed until the park is complete. if should_notify && core.park.is_some() { - self.notify_parked(); + self.notify_parked_local(); } } @@ -830,7 +839,16 @@ impl Handle { } } - fn notify_parked(&self) { + fn notify_parked_local(&self) { + super::counters::inc_num_inc_notify_local(); + + if let Some(index) = self.shared.idle.worker_to_notify() { + super::counters::inc_num_unparks_local(); + self.shared.remotes[index].unpark.unpark(&self.driver); + } + } + + fn notify_parked_remote(&self) { if let Some(index) = self.shared.idle.worker_to_notify() { self.shared.remotes[index].unpark.unpark(&self.driver); } @@ -845,13 +863,13 @@ impl Handle { fn notify_if_work_pending(&self) { for remote in &self.shared.remotes[..] { if !remote.steal.is_empty() { - self.notify_parked(); + self.notify_parked_local(); return; } } if !self.shared.inject.is_empty() { - self.notify_parked(); + self.notify_parked_local(); } } @@ -859,7 +877,7 @@ impl Handle { if self.shared.idle.transition_worker_from_searching() { // We are the final searching worker. Because work was found, we // need to notify another worker. - self.notify_parked(); + self.notify_parked_local(); } }