diff --git a/tokio/src/io/driver/mod.rs b/tokio/src/io/driver/mod.rs index 19f67a24e..b7696ac49 100644 --- a/tokio/src/io/driver/mod.rs +++ b/tokio/src/io/driver/mod.rs @@ -15,6 +15,7 @@ mod scheduled_io; use scheduled_io::ScheduledIo; use crate::park::{Park, Unpark}; +use crate::runtime::stats::IoDriverStats; use crate::util::slab::{self, Slab}; use crate::{loom::sync::Mutex, util::bit}; @@ -74,6 +75,8 @@ pub(super) struct Inner { /// Used to wake up the reactor from a call to `turn`. waker: mio::Waker, + + stats: IoDriverStats, } #[derive(Debug, Eq, PartialEq, Clone, Copy)] @@ -112,7 +115,7 @@ fn _assert_kinds() { impl Driver { /// Creates a new event loop, returning any error that happened during the /// creation. - pub(crate) fn new() -> io::Result { + pub(crate) fn new(stats: IoDriverStats) -> io::Result { let poll = mio::Poll::new()?; let waker = mio::Waker::new(poll.registry(), TOKEN_WAKEUP)?; let registry = poll.registry().try_clone()?; @@ -130,6 +133,7 @@ impl Driver { registry, io_dispatch: allocator, waker, + stats, }), }) } @@ -153,7 +157,8 @@ impl Driver { self.tick = self.tick.wrapping_add(1); if self.tick == COMPACT_INTERVAL { - self.resources.as_mut().unwrap().compact() + self.resources.as_mut().unwrap().compact(); + self.inner.stats.incr_compact_count(); } let mut events = self.events.take().expect("i/o driver event store missing"); @@ -192,6 +197,14 @@ impl Driver { let res = io.set_readiness(Some(token.0), Tick::Set(self.tick), |curr| curr | ready); + if ready.is_readable() { + self.inner.stats.incr_read_ready_count(); + } + + if ready.is_writable() { + self.inner.stats.incr_write_ready_count(); + } + if res.is_err() { // token no longer valid! return; @@ -335,12 +348,18 @@ impl Inner { self.registry .register(source, mio::Token(token), interest.to_mio())?; + self.stats.incr_fd_count(); + Ok(shared) } /// Deregisters an I/O resource from the reactor. pub(super) fn deregister_source(&self, source: &mut impl mio::event::Source) -> io::Result<()> { - self.registry.deregister(source) + self.registry.deregister(source)?; + + self.stats.dec_fd_count(); + + Ok(()) } } diff --git a/tokio/src/runtime/driver.rs b/tokio/src/runtime/driver.rs index 7e459779b..c2032dc90 100644 --- a/tokio/src/runtime/driver.rs +++ b/tokio/src/runtime/driver.rs @@ -5,6 +5,8 @@ use crate::park::Park; use std::io; use std::time::Duration; +use super::stats::IoDriverStats; + // ===== io driver ===== cfg_io_driver! { @@ -12,14 +14,14 @@ cfg_io_driver! { type IoStack = crate::park::either::Either; pub(crate) type IoHandle = Option; - fn create_io_stack(enabled: bool) -> io::Result<(IoStack, IoHandle, SignalHandle)> { + fn create_io_stack(enabled: bool, stats: IoDriverStats) -> io::Result<(IoStack, IoHandle, SignalHandle)> { use crate::park::either::Either; #[cfg(loom)] assert!(!enabled); let ret = if enabled { - let io_driver = crate::io::driver::Driver::new()?; + let io_driver = crate::io::driver::Driver::new(stats)?; let io_handle = io_driver.handle(); let (signal_driver, signal_handle) = create_signal_driver(io_driver)?; @@ -38,7 +40,7 @@ cfg_not_io_driver! { pub(crate) type IoHandle = (); type IoStack = ParkThread; - fn create_io_stack(_enabled: bool) -> io::Result<(IoStack, IoHandle, SignalHandle)> { + fn create_io_stack(_enabled: bool, _stats: IoDriverStats) -> io::Result<(IoStack, IoHandle, SignalHandle)> { Ok((ParkThread::new(), Default::default(), Default::default())) } } @@ -166,8 +168,8 @@ pub(crate) struct Cfg { } impl Driver { - pub(crate) fn new(cfg: Cfg) -> io::Result<(Self, Resources)> { - let (io_stack, io_handle, signal_handle) = create_io_stack(cfg.enable_io)?; + pub(crate) fn new(cfg: Cfg, stats: IoDriverStats) -> io::Result<(Self, Resources)> { + let (io_stack, io_handle, signal_handle) = create_io_stack(cfg.enable_io, stats)?; let clock = create_clock(cfg.enable_pause_time, cfg.start_paused); diff --git a/tokio/src/runtime/stats/mod.rs b/tokio/src/runtime/stats/mod.rs index 355e40060..70ada117d 100644 --- a/tokio/src/runtime/stats/mod.rs +++ b/tokio/src/runtime/stats/mod.rs @@ -12,7 +12,7 @@ cfg_stats! { mod stats; pub use self::stats::{RuntimeStats, WorkerStats}; - pub(crate) use self::stats::WorkerStatsBatcher; + pub(crate) use self::stats::{WorkerStatsBatcher, IoDriverStats}; } cfg_not_stats! { diff --git a/tokio/src/runtime/stats/stats.rs b/tokio/src/runtime/stats/stats.rs index 375786300..f4b5aacda 100644 --- a/tokio/src/runtime/stats/stats.rs +++ b/tokio/src/runtime/stats/stats.rs @@ -2,6 +2,7 @@ use crate::loom::sync::atomic::{AtomicU64, Ordering::Relaxed}; use std::convert::TryFrom; +use std::sync::Arc; use std::time::{Duration, Instant}; /// This type contains methods to retrieve stats from a Tokio runtime. @@ -14,6 +15,7 @@ use std::time::{Duration, Instant}; #[derive(Debug)] pub struct RuntimeStats { workers: Box<[WorkerStats]>, + driver: IoDriverStats, } /// This type contains methods to retrieve stats from a worker thread on a Tokio runtime. @@ -46,6 +48,7 @@ impl RuntimeStats { Self { workers: workers.into_boxed_slice(), + driver: IoDriverStats::default(), } } @@ -132,3 +135,39 @@ impl WorkerStatsBatcher { self.poll_count += 1; } } + +#[derive(Debug, Default, Clone)] +pub(crate) struct IoDriverStats { + inner: Arc, +} + +#[derive(Debug, Default)] +#[repr(align(128))] +struct IoDriverStatsInner { + read_ready_count: AtomicU64, + write_ready_count: AtomicU64, + fd_count: AtomicU64, + compact_count: AtomicU64, +} + +impl IoDriverStats { + pub(crate) fn incr_read_ready_count(&self) { + self.inner.read_ready_count.fetch_add(1, Relaxed); + } + + pub(crate) fn incr_write_ready_count(&self) { + self.inner.write_ready_count.fetch_add(1, Relaxed); + } + + pub(crate) fn incr_fd_count(&self) { + self.inner.fd_count.fetch_add(1, Relaxed); + } + + pub(crate) fn dec_fd_count(&self) { + self.inner.fd_count.fetch_sub(1, Relaxed); + } + + pub(crate) fn incr_compact_count(&self) { + self.inner.compact_count.fetch_add(1, Relaxed); + } +}