From 1d785fd66fce4a9125c6f0c403dfe1921d011539 Mon Sep 17 00:00:00 2001 From: Matilda Smeds Date: Thu, 27 Apr 2023 12:58:31 +0200 Subject: [PATCH] metrics: add metric for number of tasks (#5628) --- tokio/src/runtime/metrics/runtime.rs | 19 ++++++ tokio/src/runtime/scheduler/current_thread.rs | 4 ++ tokio/src/runtime/scheduler/mod.rs | 8 +++ .../runtime/scheduler/multi_thread/handle.rs | 4 ++ tokio/src/runtime/task/list.rs | 16 +++-- tokio/src/util/linked_list.rs | 67 +++++++++++++++++++ tokio/tests/rt_metrics.rs | 17 +++++ 7 files changed, 131 insertions(+), 4 deletions(-) diff --git a/tokio/src/runtime/metrics/runtime.rs b/tokio/src/runtime/metrics/runtime.rs index 8e52fead1..014995c7e 100644 --- a/tokio/src/runtime/metrics/runtime.rs +++ b/tokio/src/runtime/metrics/runtime.rs @@ -68,6 +68,25 @@ impl RuntimeMetrics { self.handle.inner.num_blocking_threads() } + /// Returns the number of active tasks in the runtime. + /// + /// # Examples + /// + /// ``` + /// use tokio::runtime::Handle; + /// + /// #[tokio::main] + /// async fn main() { + /// let metrics = Handle::current().metrics(); + /// + /// let n = metrics.active_tasks_count(); + /// println!("Runtime has {} active tasks", n); + /// } + /// ``` + pub fn active_tasks_count(&self) -> usize { + self.handle.inner.active_tasks_count() + } + /// Returns the number of idle threads, which have spawned by the runtime /// for `spawn_blocking` calls. /// diff --git a/tokio/src/runtime/scheduler/current_thread.rs b/tokio/src/runtime/scheduler/current_thread.rs index 375e47c41..64e1265c9 100644 --- a/tokio/src/runtime/scheduler/current_thread.rs +++ b/tokio/src/runtime/scheduler/current_thread.rs @@ -430,6 +430,10 @@ cfg_metrics! { pub(crate) fn blocking_queue_depth(&self) -> usize { self.blocking_spawner.queue_depth() } + + pub(crate) fn active_tasks_count(&self) -> usize { + self.shared.owned.active_tasks_count() + } } } diff --git a/tokio/src/runtime/scheduler/mod.rs b/tokio/src/runtime/scheduler/mod.rs index f45d8a80b..0ea207e1e 100644 --- a/tokio/src/runtime/scheduler/mod.rs +++ b/tokio/src/runtime/scheduler/mod.rs @@ -135,6 +135,14 @@ cfg_rt! { } } + pub(crate) fn active_tasks_count(&self) -> usize { + match self { + Handle::CurrentThread(handle) => handle.active_tasks_count(), + #[cfg(all(feature = "rt-multi-thread", not(tokio_wasi)))] + Handle::MultiThread(handle) => handle.active_tasks_count(), + } + } + pub(crate) fn scheduler_metrics(&self) -> &SchedulerMetrics { match self { Handle::CurrentThread(handle) => handle.scheduler_metrics(), diff --git a/tokio/src/runtime/scheduler/multi_thread/handle.rs b/tokio/src/runtime/scheduler/multi_thread/handle.rs index 69a4620c1..77baaeb06 100644 --- a/tokio/src/runtime/scheduler/multi_thread/handle.rs +++ b/tokio/src/runtime/scheduler/multi_thread/handle.rs @@ -69,6 +69,10 @@ cfg_metrics! { self.blocking_spawner.num_idle_threads() } + pub(crate) fn active_tasks_count(&self) -> usize { + self.shared.owned.active_tasks_count() + } + pub(crate) fn scheduler_metrics(&self) -> &SchedulerMetrics { &self.shared.scheduler_metrics } diff --git a/tokio/src/runtime/task/list.rs b/tokio/src/runtime/task/list.rs index 159c13e16..c520cb04e 100644 --- a/tokio/src/runtime/task/list.rs +++ b/tokio/src/runtime/task/list.rs @@ -10,7 +10,7 @@ use crate::future::Future; use crate::loom::cell::UnsafeCell; use crate::loom::sync::Mutex; use crate::runtime::task::{JoinHandle, LocalNotified, Notified, Schedule, Task}; -use crate::util::linked_list::{Link, LinkedList}; +use crate::util::linked_list::{CountedLinkedList, Link, LinkedList}; use std::marker::PhantomData; @@ -54,9 +54,13 @@ cfg_not_has_atomic_u64! { } pub(crate) struct OwnedTasks { - inner: Mutex>, + inner: Mutex>, id: u64, } +struct CountedOwnedTasksInner { + list: CountedLinkedList, as Link>::Target>, + closed: bool, +} pub(crate) struct LocalOwnedTasks { inner: UnsafeCell>, id: u64, @@ -70,8 +74,8 @@ struct OwnedTasksInner { impl OwnedTasks { pub(crate) fn new() -> Self { Self { - inner: Mutex::new(OwnedTasksInner { - list: LinkedList::new(), + inner: Mutex::new(CountedOwnedTasksInner { + list: CountedLinkedList::new(), closed: false, }), id: get_next_id(), @@ -153,6 +157,10 @@ impl OwnedTasks { } } + pub(crate) fn active_tasks_count(&self) -> usize { + self.inner.lock().list.count() + } + pub(crate) fn remove(&self, task: &Task) -> Option> { let task_id = task.header().get_owner_id(); if task_id == 0 { diff --git a/tokio/src/util/linked_list.rs b/tokio/src/util/linked_list.rs index 6bf268184..74a684d1a 100644 --- a/tokio/src/util/linked_list.rs +++ b/tokio/src/util/linked_list.rs @@ -228,6 +228,53 @@ impl fmt::Debug for LinkedList { } } +// ===== impl CountedLinkedList ==== + +// Delegates operations to the base LinkedList implementation, and adds a counter to the elements +// in the list. +pub(crate) struct CountedLinkedList { + list: LinkedList, + count: usize, +} + +impl CountedLinkedList { + pub(crate) fn new() -> CountedLinkedList { + CountedLinkedList { + list: LinkedList::new(), + count: 0, + } + } + + pub(crate) fn push_front(&mut self, val: L::Handle) { + self.list.push_front(val); + self.count += 1; + } + + pub(crate) fn pop_back(&mut self) -> Option { + let val = self.list.pop_back(); + if val.is_some() { + self.count -= 1; + } + val + } + + pub(crate) fn is_empty(&self) -> bool { + self.list.is_empty() + } + + pub(crate) unsafe fn remove(&mut self, node: NonNull) -> Option { + let val = self.list.remove(node); + if val.is_some() { + self.count -= 1; + } + val + } + + pub(crate) fn count(&self) -> usize { + self.count + } +} + #[cfg(any( feature = "fs", feature = "rt", @@ -719,6 +766,26 @@ pub(crate) mod tests { } } + #[test] + fn count() { + let mut list = CountedLinkedList::<&Entry, <&Entry as Link>::Target>::new(); + assert_eq!(0, list.count()); + + let a = entry(5); + let b = entry(7); + list.push_front(a.as_ref()); + list.push_front(b.as_ref()); + assert_eq!(2, list.count()); + + list.pop_back(); + assert_eq!(1, list.count()); + + unsafe { + list.remove(ptr(&b)); + } + assert_eq!(0, list.count()); + } + /// This is a fuzz test. You run it by entering `cargo fuzz run fuzz_linked_list` in CLI in `/tokio/` module. #[cfg(fuzzing)] pub fn fuzz_linked_list(ops: &[u8]) { diff --git a/tokio/tests/rt_metrics.rs b/tokio/tests/rt_metrics.rs index d238808c8..76aa0587f 100644 --- a/tokio/tests/rt_metrics.rs +++ b/tokio/tests/rt_metrics.rs @@ -81,6 +81,23 @@ fn blocking_queue_depth() { assert_eq!(0, rt.metrics().blocking_queue_depth()); } +#[test] +fn active_tasks_count() { + let rt = current_thread(); + let metrics = rt.metrics(); + assert_eq!(0, metrics.active_tasks_count()); + rt.spawn(async move { + assert_eq!(1, metrics.active_tasks_count()); + }); + + let rt = threaded(); + let metrics = rt.metrics(); + assert_eq!(0, metrics.active_tasks_count()); + rt.spawn(async move { + assert_eq!(1, metrics.active_tasks_count()); + }); +} + #[test] fn remote_schedule_count() { use std::thread;