Skip to main content

subms_timer_wheel/features/
metrics.rs

1//! Metered timer wheel: thin wrapper around the base `TimerWheel`
2//! that tracks per-instance counters. Counters are plain `u64`
3//! fields - no atomics, no locks - because the underlying wheel is
4//! itself single-threaded. Pair with the `concurrent` feature for a
5//! thread-safe metered surface (wrap a `MeteredTimerWheel` inside a
6//! mutex of your own).
7//!
8//! `cascade_events` is always 0 for the single-level base wheel.
9//! It's tracked here so downstream code that swaps the wheel for
10//! the hierarchical variant doesn't need a schema change in the
11//! metrics snapshot. The hierarchical wheel exposes its own
12//! `cascades()` counter directly; see that module.
13
14use crate::{TimerError, TimerWheel};
15
16#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
17pub struct TimerMetrics {
18    pub scheduled: u64,
19    pub fired: u64,
20    pub cancelled: u64,
21    pub rescheduled: u64,
22    pub drained: u64,
23    pub ticks: u64,
24    pub cascade_events: u64,
25}
26
27pub struct MeteredTimerWheel<V> {
28    wheel: TimerWheel<V>,
29    metrics: TimerMetrics,
30}
31
32impl<V> MeteredTimerWheel<V> {
33    pub fn new(num_slots: usize) -> Self {
34        Self {
35            wheel: TimerWheel::new(num_slots),
36            metrics: TimerMetrics::default(),
37        }
38    }
39
40    pub fn num_slots(&self) -> usize {
41        self.wheel.num_slots()
42    }
43
44    pub fn max_delay(&self) -> u64 {
45        self.wheel.max_delay()
46    }
47
48    pub fn pending(&self) -> usize {
49        self.wheel.pending()
50    }
51
52    pub fn is_empty(&self) -> bool {
53        self.wheel.is_empty()
54    }
55
56    pub fn slot_len(&self, slot: usize) -> usize {
57        self.wheel.slot_len(slot)
58    }
59
60    pub fn metrics(&self) -> TimerMetrics {
61        self.metrics
62    }
63
64    pub fn schedule(&mut self, delay_ticks: usize, value: V) -> u64 {
65        self.metrics.scheduled += 1;
66        self.wheel.schedule(delay_ticks, value)
67    }
68
69    pub fn try_schedule(&mut self, delay_ticks: usize, value: V) -> Result<u64, TimerError> {
70        let id = self.wheel.try_schedule(delay_ticks, value)?;
71        self.metrics.scheduled += 1;
72        Ok(id)
73    }
74
75    pub fn cancel(&mut self, id: u64) -> bool {
76        let ok = self.wheel.cancel(id);
77        if ok {
78            self.metrics.cancelled += 1;
79        }
80        ok
81    }
82
83    pub fn reschedule(&mut self, id: u64, delay_ticks: usize) -> bool {
84        let ok = self.wheel.reschedule(id, delay_ticks);
85        if ok {
86            self.metrics.rescheduled += 1;
87        }
88        ok
89    }
90
91    pub fn tick(&mut self) -> Vec<V> {
92        self.metrics.ticks += 1;
93        let fired = self.wheel.tick();
94        self.metrics.fired += fired.len() as u64;
95        fired
96    }
97
98    pub fn advance(&mut self, ticks: usize) -> Vec<V> {
99        let mut fired = Vec::new();
100        for _ in 0..ticks {
101            fired.append(&mut self.tick());
102        }
103        fired
104    }
105
106    /// Hand back every pending timer. Counted separately from `fired`: a
107    /// drained timer never came due, and folding the two together would make
108    /// a shutdown look like a burst of expiries.
109    pub fn drain(&mut self) -> Vec<V> {
110        let out = self.wheel.drain();
111        self.metrics.drained += out.len() as u64;
112        out
113    }
114
115    pub fn clear(&mut self) {
116        self.metrics.drained += self.wheel.pending() as u64;
117        self.wheel.clear();
118    }
119}
120
121#[cfg(test)]
122#[path = "metrics_tests.rs"]
123mod tests;