tokio_metrics/runtime/
poll_time_histogram.rs1use std::time::Duration;
2
3#[derive(Debug, Clone, Default)]
14#[non_exhaustive]
15pub struct PollTimeHistogram {
16 buckets: Vec<HistogramBucket>,
17}
18
19impl PollTimeHistogram {
20 #[cfg_attr(not(tokio_unstable), allow(dead_code))]
22 pub(crate) fn new(buckets: Vec<HistogramBucket>) -> Self {
23 Self { buckets }
24 }
25
26 pub fn buckets(&self) -> &[HistogramBucket] {
28 &self.buckets
29 }
30
31 #[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 pub fn as_counts(&self) -> Vec<u64> {
39 self.buckets.iter().map(|b| b.count).collect()
40 }
41}
42
43#[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 #[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 pub fn range_start(&self) -> Duration {
61 self.range_start
62 }
63
64 pub fn range_end(&self) -> Duration {
66 self.range_end
67 }
68
69 pub fn count(&self) -> u64 {
71 self.count
72 }
73
74 #[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 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 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)] 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#[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 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}