1use core::time::Duration;
7use std::time::Instant;
8
9use crate::model::{MetricState, UnavailableReason};
10use crate::rates::keyed::DeltaTracker;
11use crate::units::Rate;
12
13#[derive(Clone, Copy, Debug, Default, Eq, Hash, PartialEq)]
22#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
23#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
24pub enum CounterWidth {
25 #[default]
30 Unknown,
31 Bits32,
34 Bits64,
40}
41
42impl CounterWidth {
43 #[must_use]
45 pub const fn bits(self) -> Option<u32> {
46 match self {
47 Self::Unknown => None,
48 Self::Bits32 => Some(32),
49 Self::Bits64 => Some(64),
50 }
51 }
52
53 #[must_use]
55 pub const fn max_value(self) -> Option<u64> {
56 match self {
57 Self::Unknown => None,
58 Self::Bits32 => Some(u32::MAX as u64),
60 Self::Bits64 => Some(u64::MAX),
61 }
62 }
63
64 const fn modulus(self) -> Option<u128> {
68 match self {
69 Self::Unknown => None,
70 Self::Bits32 => Some(1u128 << 32),
71 Self::Bits64 => Some(1u128 << 64),
72 }
73 }
74}
75
76fn forward_delta(previous: u64, current: u64, width: CounterWidth) -> Option<u64> {
88 if current >= previous {
89 return Some(current - previous);
90 }
91 let modulus = width.modulus()?;
92 if u128::from(previous) >= modulus || u128::from(current) >= modulus {
95 return None;
96 }
97 let backwards = u128::from(previous) - u128::from(current);
98 let wrapped = modulus - backwards;
99 if wrapped < backwards {
100 u64::try_from(wrapped).ok()
101 } else {
102 None
103 }
104}
105
106#[derive(Clone, Copy, Debug, Eq, PartialEq)]
113pub enum CounterDelta {
114 FirstSample,
119 Advanced {
121 delta: u64,
123 elapsed: Duration,
129 wrapped: bool,
134 },
135 Reset,
140}
141
142impl CounterDelta {
143 #[must_use]
145 pub fn rate(self) -> MetricState<Rate> {
146 match self {
147 Self::FirstSample => MetricState::WarmingUp,
148 Self::Advanced { delta, elapsed, .. } => match Rate::from_delta(delta, elapsed) {
149 Some(rate) => MetricState::Available(rate),
150 None => MetricState::WarmingUp,
155 },
156 Self::Reset => MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset),
157 }
158 }
159
160 #[must_use]
165 pub const fn advanced_by(self) -> Option<u64> {
166 match self {
167 Self::Advanced { delta, .. } => Some(delta),
168 Self::FirstSample | Self::Reset => None,
169 }
170 }
171
172 #[must_use]
174 pub const fn wrapped(self) -> bool {
175 matches!(self, Self::Advanced { wrapped: true, .. })
176 }
177}
178
179#[derive(Clone, Copy, Debug)]
209pub struct CounterTracker {
210 width: CounterWidth,
211 last: Option<Reading>,
212}
213
214#[derive(Clone, Copy, Debug)]
216struct Reading {
217 value: u64,
218 at: Instant,
219}
220
221impl CounterTracker {
222 #[must_use]
224 pub const fn new(width: CounterWidth) -> Self {
225 Self { width, last: None }
226 }
227
228 #[must_use]
230 pub const fn width(&self) -> CounterWidth {
231 self.width
232 }
233
234 #[must_use]
236 pub const fn is_warming_up(&self) -> bool {
237 self.last.is_none()
238 }
239
240 #[must_use]
242 pub const fn last_value(&self) -> Option<u64> {
243 match self.last {
244 Some(reading) => Some(reading.value),
245 None => None,
246 }
247 }
248
249 #[must_use]
251 pub const fn last_observed_at(&self) -> Option<Instant> {
252 match self.last {
253 Some(reading) => Some(reading.at),
254 None => None,
255 }
256 }
257
258 pub fn forget_baseline(&mut self) {
264 self.last = None;
265 }
266
267 pub fn observe(&mut self, value: u64, at: Instant) -> CounterDelta {
271 let Some(previous) = self.last.replace(Reading { value, at }) else {
272 return CounterDelta::FirstSample;
273 };
274 let elapsed = at.saturating_duration_since(previous.at);
279 match forward_delta(previous.value, value, self.width) {
280 Some(delta) => CounterDelta::Advanced {
281 delta,
282 elapsed,
283 wrapped: value < previous.value,
286 },
287 None => CounterDelta::Reset,
290 }
291 }
292
293 pub fn rate(&mut self, value: u64, at: Instant) -> MetricState<Rate> {
295 self.observe(value, at).rate()
296 }
297}
298
299impl DeltaTracker for CounterTracker {
300 type Config = CounterWidth;
301 type Reading = u64;
302 type Value = Rate;
303
304 fn with_config(config: Self::Config) -> Self {
305 Self::new(config)
306 }
307
308 fn observe_reading(&mut self, reading: Self::Reading, at: Instant) -> MetricState<Self::Value> {
309 self.rate(reading, at)
310 }
311
312 fn last_observed_at(&self) -> Option<Instant> {
315 self.last.map(|reading| reading.at)
316 }
317
318 fn forget_baseline(&mut self) {
319 self.last = None;
320 }
321}
322
323#[cfg(test)]
324mod tests {
325 use super::*;
326
327 fn origin() -> Instant {
329 Instant::now()
330 }
331
332 fn assert_rate(state: &MetricState<Rate>, expected: f64) {
336 let actual = state
337 .fresh()
338 .expect("expected a measured rate")
339 .per_second();
340 assert!(
341 (actual - expected).abs() < 1e-6,
342 "expected {expected}/s, got {actual}/s"
343 );
344 }
345
346 #[test]
347 fn a_first_sample_is_warming_up_and_not_zero() {
348 let mut tracker = CounterTracker::new(CounterWidth::Bits64);
349 let state = tracker.rate(4_096, origin());
350 assert!(state.is_warming_up());
351 assert_eq!(state.fresh(), None);
352 assert_ne!(state, MetricState::Available(Rate::ZERO));
353 }
354
355 #[test]
356 fn a_second_sample_divides_by_the_real_interval() {
357 let t0 = origin();
358 let mut tracker = CounterTracker::new(CounterWidth::Bits64);
359 assert_eq!(tracker.observe(1_000, t0), CounterDelta::FirstSample);
360 let state = tracker.rate(3_000, t0 + Duration::from_secs(2));
361 assert_rate(&state, 1_000.0);
362 }
363
364 #[test]
365 fn the_same_delta_over_different_intervals_gives_different_rates() {
366 let t0 = origin();
367 let mut fast = CounterTracker::new(CounterWidth::Bits64);
368 let mut slow = CounterTracker::new(CounterWidth::Bits64);
369 fast.rate(0, t0);
370 slow.rate(0, t0);
371
372 let half = fast.rate(1_000, t0 + Duration::from_millis(500));
373 let double = slow.rate(1_000, t0 + Duration::from_secs(2));
374
375 assert_rate(&half, 2_000.0);
376 assert_rate(&double, 500.0);
377 }
378
379 #[test]
380 fn a_counter_that_does_not_move_is_a_real_zero_rate() {
381 let t0 = origin();
384 let mut tracker = CounterTracker::new(CounterWidth::Bits64);
385 tracker.rate(7_777, t0);
386 let state = tracker.rate(7_777, t0 + Duration::from_secs(1));
387 assert_eq!(state, MetricState::Available(Rate::ZERO));
388 }
389
390 #[test]
391 fn zero_elapsed_is_warming_up_rather_than_a_division_by_zero() {
392 let t0 = origin();
393 let mut tracker = CounterTracker::new(CounterWidth::Bits64);
394 tracker.rate(100, t0);
395 let state = tracker.rate(900, t0);
396 assert!(state.is_warming_up());
397 let mut totals = CounterTracker::new(CounterWidth::Bits64);
399 totals.observe(100, t0);
400 assert_eq!(totals.observe(900, t0).advanced_by(), Some(800));
401 }
402
403 #[test]
404 fn a_reversed_instant_cannot_produce_a_negative_or_huge_rate() {
405 let t0 = origin();
408 let mut tracker = CounterTracker::new(CounterWidth::Bits64);
409 tracker.rate(0, t0 + Duration::from_secs(10));
410 let state = tracker.rate(1_000_000, t0);
411 assert!(state.is_warming_up());
412 }
413
414 #[test]
415 fn a_backwards_counter_of_unknown_width_is_a_typed_reset() {
416 let t0 = origin();
417 let mut tracker = CounterTracker::new(CounterWidth::Unknown);
418 tracker.rate(9_000_000, t0);
419 let state = tracker.rate(12, t0 + Duration::from_secs(1));
420 assert_eq!(
421 state,
422 MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset)
423 );
424 assert_eq!(state.fresh(), None);
425 }
426
427 #[test]
428 fn the_sample_after_a_reset_is_valid_again() {
429 let t0 = origin();
430 let mut tracker = CounterTracker::new(CounterWidth::Unknown);
431 tracker.rate(9_000_000, t0);
432 let reset = tracker.rate(12, t0 + Duration::from_secs(1));
433 assert!(!reset.is_available());
434
435 let recovered = tracker.rate(1_012, t0 + Duration::from_secs(2));
437 assert_rate(&recovered, 1_000.0);
438 }
439
440 #[test]
441 fn a_reset_never_reports_a_rate_derived_from_the_new_value() {
442 let t0 = origin();
445 let mut tracker = CounterTracker::new(CounterWidth::Unknown);
446 tracker.rate(9_000_000, t0);
447 let delta = tracker.observe(12, t0 + Duration::from_secs(1));
448 assert_eq!(delta, CounterDelta::Reset);
449 assert_eq!(delta.advanced_by(), None);
450 }
451
452 #[test]
453 fn a_known_width_counter_wraps_instead_of_resetting() {
454 let t0 = origin();
456 let ceiling = u64::from(u32::MAX) + 1;
457 let previous = ceiling - 300;
458 let mut tracker = CounterTracker::new(CounterWidth::Bits32);
459 tracker.observe(previous, t0);
460
461 let delta = tracker.observe(700, t0 + Duration::from_secs(1));
462 assert_eq!(
463 delta,
464 CounterDelta::Advanced {
465 delta: 1_000,
466 elapsed: Duration::from_secs(1),
467 wrapped: true,
468 }
469 );
470 assert_rate(&delta.rate(), 1_000.0);
471 assert!(delta.wrapped());
472 }
473
474 #[test]
475 fn the_same_movement_is_a_reset_when_the_width_is_unknown() {
476 let t0 = origin();
477 let previous = u64::from(u32::MAX) + 1 - 300;
478 let mut tracker = CounterTracker::new(CounterWidth::Unknown);
479 tracker.observe(previous, t0);
480 assert_eq!(
481 tracker.observe(700, t0 + Duration::from_secs(1)),
482 CounterDelta::Reset
483 );
484 }
485
486 #[test]
487 fn a_wrap_at_the_exact_boundary_is_reconstructed_exactly() {
488 let t0 = origin();
489 let mut tracker = CounterTracker::new(CounterWidth::Bits32);
490 tracker.observe(u64::from(u32::MAX), t0);
491 assert_eq!(
492 tracker
493 .observe(0, t0 + Duration::from_secs(1))
494 .advanced_by(),
495 Some(1),
496 "u32::MAX -> 0 is a single step forward"
497 );
498 }
499
500 #[test]
501 fn a_small_backwards_move_is_a_reset_even_at_a_known_width() {
502 let t0 = origin();
505 let mut tracker = CounterTracker::new(CounterWidth::Bits32);
506 tracker.observe(3_000_000_000, t0);
507 assert_eq!(
508 tracker.observe(2_999_000_000, t0 + Duration::from_secs(1)),
509 CounterDelta::Reset
510 );
511 }
512
513 #[test]
514 fn a_reading_outside_the_declared_width_is_a_reset_not_a_wrap() {
515 let t0 = origin();
517 let mut tracker = CounterTracker::new(CounterWidth::Bits32);
518 tracker.observe(u64::from(u32::MAX) + 5_000, t0);
519 assert_eq!(
520 tracker.observe(10, t0 + Duration::from_secs(1)),
521 CounterDelta::Reset
522 );
523 }
524
525 #[test]
526 fn a_wrapped_delta_can_never_exceed_half_the_counter_range() {
527 let half = 1u64 << 31;
529 let t0 = origin();
530 for previous in [u64::from(u32::MAX), 3_000_000_000, half + 1] {
531 for current in [0, 1, 1_000, half - 1] {
532 let mut tracker = CounterTracker::new(CounterWidth::Bits32);
533 tracker.observe(previous, t0);
534 if let Some(delta) = tracker
535 .observe(current, t0 + Duration::from_secs(1))
536 .advanced_by()
537 && current < previous
538 {
539 assert!(delta < half, "{previous} -> {current} produced {delta}");
540 }
541 }
542 }
543 }
544
545 #[test]
546 fn a_sixty_four_bit_counter_still_rejects_an_absurd_backwards_jump() {
547 let t0 = origin();
548 let mut tracker = CounterTracker::new(CounterWidth::Bits64);
549 tracker.observe(1_000_000, t0);
550 assert_eq!(
551 tracker.observe(9, t0 + Duration::from_secs(1)),
552 CounterDelta::Reset
553 );
554 }
555
556 #[test]
557 fn forgetting_the_baseline_makes_the_next_reading_warm_up() {
558 let t0 = origin();
559 let mut tracker = CounterTracker::new(CounterWidth::Bits64);
560 tracker.rate(1_000, t0);
561 assert!(!tracker.is_warming_up());
562
563 tracker.forget_baseline();
564 assert!(tracker.is_warming_up());
565 assert_eq!(tracker.last_value(), None);
566 assert_eq!(tracker.last_observed_at(), None);
567 assert!(
568 tracker
569 .rate(500_000, t0 + Duration::from_secs(1))
570 .is_warming_up(),
571 "a dropped baseline must not be reconstructed from the old value"
572 );
573 }
574
575 #[test]
576 fn the_baseline_tracks_the_most_recent_reading() {
577 let t0 = origin();
578 let at = t0 + Duration::from_secs(3);
579 let mut tracker = CounterTracker::new(CounterWidth::Bits32);
580 tracker.observe(42, t0);
581 tracker.observe(84, at);
582 assert_eq!(tracker.last_value(), Some(84));
583 assert_eq!(tracker.last_observed_at(), Some(at));
584 assert_eq!(tracker.width(), CounterWidth::Bits32);
585 }
586
587 #[test]
588 fn widths_report_their_own_limits() {
589 assert_eq!(CounterWidth::Unknown.bits(), None);
590 assert_eq!(CounterWidth::Unknown.max_value(), None);
591 assert_eq!(CounterWidth::Bits32.bits(), Some(32));
592 assert_eq!(CounterWidth::Bits32.max_value(), Some(u64::from(u32::MAX)));
593 assert_eq!(CounterWidth::Bits64.bits(), Some(64));
594 assert_eq!(CounterWidth::Bits64.max_value(), Some(u64::MAX));
595 assert_eq!(CounterWidth::default(), CounterWidth::Unknown);
596 }
597
598 #[test]
599 fn a_forward_move_is_never_treated_as_a_wrap() {
600 let t0 = origin();
601 let mut tracker = CounterTracker::new(CounterWidth::Bits32);
602 tracker.observe(10, t0);
603 let delta = tracker.observe(4_000_000_000, t0 + Duration::from_secs(1));
604 assert_eq!(delta.advanced_by(), Some(3_999_999_990));
605 assert!(!delta.wrapped());
606 }
607}