Skip to main content

tokio_metrics/runtime/
poll_time_histogram.rs

1use std::time::Duration;
2
3/// A histogram of task poll durations, pairing each bucket's count with its
4/// time range from the runtime configuration.
5///
6/// This type is returned as part of [`RuntimeMetrics`][super::RuntimeMetrics]
7/// when the runtime has poll time histograms enabled via
8/// [`enable_metrics_poll_time_histogram`][tokio::runtime::Builder::enable_metrics_poll_time_histogram].
9///
10/// Each bucket contains the [`Duration`] range configured for that bucket and
11/// the count of task polls that fell into that range during the sampling
12/// interval.
13#[derive(Debug, Clone, Default)]
14#[non_exhaustive]
15pub struct PollTimeHistogram {
16    buckets: Vec<HistogramBucket>,
17}
18
19impl PollTimeHistogram {
20    // Only used to populate the histogram, which requires `tokio_unstable`.
21    #[cfg_attr(not(tokio_unstable), allow(dead_code))]
22    pub(crate) fn new(buckets: Vec<HistogramBucket>) -> Self {
23        Self { buckets }
24    }
25
26    /// Returns the histogram buckets.
27    pub fn buckets(&self) -> &[HistogramBucket] {
28        &self.buckets
29    }
30
31    // Only used to populate the histogram, which requires `tokio_unstable`.
32    #[cfg_attr(not(tokio_unstable), allow(dead_code))]
33    pub(crate) fn buckets_mut(&mut self) -> &mut [HistogramBucket] {
34        &mut self.buckets
35    }
36
37    /// Returns just the bucket counts as a `Vec<u64>`.
38    pub fn as_counts(&self) -> Vec<u64> {
39        self.buckets.iter().map(|b| b.count).collect()
40    }
41}
42
43/// A single bucket in a [`PollTimeHistogram`].
44#[derive(Debug, Clone, Copy, Default)]
45#[non_exhaustive]
46pub struct HistogramBucket {
47    range_start: Duration,
48    range_end: Duration,
49    count: u64,
50}
51
52impl HistogramBucket {
53    // Only used to populate the histogram, which requires `tokio_unstable`.
54    #[cfg_attr(not(tokio_unstable), allow(dead_code))]
55    pub(crate) fn new(range_start: Duration, range_end: Duration, count: u64) -> Self {
56        Self { range_start, range_end, count }
57    }
58
59    /// The start of the time range for this bucket (inclusive).
60    pub fn range_start(&self) -> Duration {
61        self.range_start
62    }
63
64    /// The end of the time range for this bucket (exclusive).
65    pub fn range_end(&self) -> Duration {
66        self.range_end
67    }
68
69    /// Returns the poll count for this bucket during the interval.
70    pub fn count(&self) -> u64 {
71        self.count
72    }
73
74    /// Adds to the count of this bucket.
75    // Only used to populate the histogram, which requires `tokio_unstable`.
76    #[cfg_attr(not(tokio_unstable), allow(dead_code))]
77    pub(crate) fn add_count(&mut self, delta: u64) {
78        self.count = self.count.saturating_add(delta);
79    }
80}
81
82#[cfg(feature = "metrique-integration")]
83impl metrique::writer::Value for PollTimeHistogram {
84    // Emitted as a distribution of bucket midpoints in microseconds, so the
85    // closed shape is a float rather than the `Opaque` default.
86    const SHAPE: metrique::writer::core::FieldShape<'static> =
87        metrique::writer::core::FieldShape::Known(metrique::writer::core::KnownShape::F64);
88    const UNIT: metrique::writer::Unit = metrique::writer::Unit::Second(
89        metrique::writer::unit::NegativeScale::Micro,
90    );
91
92    fn write(&self, writer: impl metrique::writer::ValueWriter) {
93        use metrique::writer::{MetricFlags, Observation};
94
95        // Use the bucket midpoint as the representative value.
96        // Tokio's last bucket has range_end of Duration::from_nanos(u64::MAX),
97        // so use range_start for it since the midpoint wouldn't be representative.
98        const LAST_BUCKET_END: Duration = Duration::from_nanos(u64::MAX);
99        writer.metric(
100            self.buckets.iter().filter(|b| b.count > 0).map(|b| {
101                let value_us = if b.range_end == LAST_BUCKET_END {
102                    b.range_start.as_micros() as f64
103                } else {
104                    #[allow(clippy::incompatible_msrv)] // metrique-integration requires 1.89+
105                    f64::midpoint(
106                        b.range_start.as_micros() as f64,
107                        b.range_end.as_micros() as f64,
108                    )
109                };
110                Observation::Repeated {
111                    total: value_us * b.count as f64,
112                    occurrences: b.count,
113                }
114            }),
115            Self::UNIT,
116            [],
117            MetricFlags::empty(),
118        );
119    }
120}
121
122#[cfg(feature = "metrique-integration")]
123impl metrique::CloseValue for PollTimeHistogram {
124    type Closed = Self;
125
126    fn close(self) -> Self {
127        self
128    }
129}
130
131// `poll_time_histogram_last_bucket_uses_range_start` constructs a
132// `RuntimeMetrics` and reads its `poll_time_histogram` field, both of which
133// require `tokio_unstable`.
134#[cfg(all(test, tokio_unstable, feature = "metrique-integration"))]
135mod tests {
136    use super::*;
137    use crate::runtime::RuntimeMetrics;
138    use metrique::CloseValue;
139    use metrique::test_util::test_metric;
140
141    #[test]
142    fn poll_time_histogram_close_value() {
143        let hist = PollTimeHistogram::new(vec![
144            HistogramBucket::new(Duration::from_micros(0), Duration::from_micros(100), 5),
145            HistogramBucket::new(Duration::from_micros(100), Duration::from_micros(200), 0),
146            HistogramBucket::new(Duration::from_micros(200), Duration::from_micros(500), 3),
147        ]);
148
149        let closed = hist.close();
150        let buckets = closed.buckets();
151        assert_eq!(buckets.len(), 3);
152        assert_eq!(buckets[0].count(), 5);
153        assert_eq!(buckets[0].range_start(), Duration::from_micros(0));
154        assert_eq!(buckets[0].range_end(), Duration::from_micros(100));
155        assert_eq!(buckets[1].count(), 0);
156        assert_eq!(buckets[2].count(), 3);
157        assert_eq!(buckets[2].range_start(), Duration::from_micros(200));
158        assert_eq!(buckets[2].range_end(), Duration::from_micros(500));
159    }
160
161    #[test]
162    fn poll_time_histogram_declares_shape_and_unit() {
163        use metrique::writer::Value;
164        use metrique::writer::core::{FieldShape, KnownShape};
165
166        assert_eq!(
167            <PollTimeHistogram as Value>::SHAPE,
168            FieldShape::Known(KnownShape::F64)
169        );
170
171        let metrics = RuntimeMetrics {
172            poll_time_histogram: PollTimeHistogram::new(vec![HistogramBucket::new(
173                Duration::from_micros(0),
174                Duration::from_micros(100),
175                1,
176            )]),
177            ..Default::default()
178        };
179
180        // The unit reaching the writer must be the one the impl declares.
181        let entry = test_metric(metrics);
182        assert_eq!(
183            entry.metrics["poll_time_histogram"].unit,
184            <PollTimeHistogram as Value>::UNIT
185        );
186    }
187
188    #[test]
189    fn poll_time_histogram_last_bucket_uses_range_start() {
190        let last_bucket_start = Duration::from_millis(500);
191        let metrics = RuntimeMetrics {
192            poll_time_histogram: PollTimeHistogram::new(vec![
193                HistogramBucket::new(Duration::from_micros(0), Duration::from_micros(100), 0),
194                HistogramBucket::new(last_bucket_start, Duration::from_nanos(u64::MAX), 2),
195            ]),
196            ..Default::default()
197        };
198
199        let entry = test_metric(metrics);
200        let hist = &entry.metrics["poll_time_histogram"];
201        assert_eq!(hist.distribution.len(), 1);
202
203        match hist.distribution[0] {
204            metrique::writer::Observation::Repeated { total, occurrences } => {
205                assert_eq!(occurrences, 2);
206                let expected = last_bucket_start.as_micros() as f64 * 2.0;
207                assert!((total - expected).abs() < 0.01);
208            }
209            other => panic!("expected Repeated, got {other:?}"),
210        }
211    }
212}