mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-28 00:00:11 +02:00
Add initial io driver stats
This commit is contained in:
@@ -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<Driver> {
|
||||
pub(crate) fn new(stats: IoDriverStats) -> io::Result<Driver> {
|
||||
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(())
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<ProcessDriver, ParkThread>;
|
||||
pub(crate) type IoHandle = Option<crate::io::driver::Handle>;
|
||||
|
||||
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);
|
||||
|
||||
|
||||
@@ -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! {
|
||||
|
||||
@@ -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<IoDriverStatsInner>,
|
||||
}
|
||||
|
||||
#[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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user