use core::time::Duration;
use std::time::Instant;
use crate::model::{MetricState, UnavailableReason};
use crate::rates::keyed::DeltaTracker;
use crate::units::Rate;
#[derive(Clone, Copy, Debug, Default, Eq, Hash, PartialEq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
pub enum CounterWidth {
#[default]
Unknown,
Bits32,
Bits64,
}
impl CounterWidth {
#[must_use]
pub const fn bits(self) -> Option<u32> {
match self {
Self::Unknown => None,
Self::Bits32 => Some(32),
Self::Bits64 => Some(64),
}
}
#[must_use]
pub const fn max_value(self) -> Option<u64> {
match self {
Self::Unknown => None,
Self::Bits32 => Some(u32::MAX as u64),
Self::Bits64 => Some(u64::MAX),
}
}
const fn modulus(self) -> Option<u128> {
match self {
Self::Unknown => None,
Self::Bits32 => Some(1u128 << 32),
Self::Bits64 => Some(1u128 << 64),
}
}
}
fn forward_delta(previous: u64, current: u64, width: CounterWidth) -> Option<u64> {
if current >= previous {
return Some(current - previous);
}
let modulus = width.modulus()?;
if u128::from(previous) >= modulus || u128::from(current) >= modulus {
return None;
}
let backwards = u128::from(previous) - u128::from(current);
let wrapped = modulus - backwards;
if wrapped < backwards {
u64::try_from(wrapped).ok()
} else {
None
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CounterDelta {
FirstSample,
Advanced {
delta: u64,
elapsed: Duration,
wrapped: bool,
},
Reset,
}
impl CounterDelta {
#[must_use]
pub fn rate(self) -> MetricState<Rate> {
match self {
Self::FirstSample => MetricState::WarmingUp,
Self::Advanced { delta, elapsed, .. } => match Rate::from_delta(delta, elapsed) {
Some(rate) => MetricState::Available(rate),
None => MetricState::WarmingUp,
},
Self::Reset => MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset),
}
}
#[must_use]
pub const fn advanced_by(self) -> Option<u64> {
match self {
Self::Advanced { delta, .. } => Some(delta),
Self::FirstSample | Self::Reset => None,
}
}
#[must_use]
pub const fn wrapped(self) -> bool {
matches!(self, Self::Advanced { wrapped: true, .. })
}
}
#[derive(Clone, Copy, Debug)]
pub struct CounterTracker {
width: CounterWidth,
last: Option<Reading>,
}
#[derive(Clone, Copy, Debug)]
struct Reading {
value: u64,
at: Instant,
}
impl CounterTracker {
#[must_use]
pub const fn new(width: CounterWidth) -> Self {
Self { width, last: None }
}
#[must_use]
pub const fn width(&self) -> CounterWidth {
self.width
}
#[must_use]
pub const fn is_warming_up(&self) -> bool {
self.last.is_none()
}
#[must_use]
pub const fn last_value(&self) -> Option<u64> {
match self.last {
Some(reading) => Some(reading.value),
None => None,
}
}
#[must_use]
pub const fn last_observed_at(&self) -> Option<Instant> {
match self.last {
Some(reading) => Some(reading.at),
None => None,
}
}
pub fn forget_baseline(&mut self) {
self.last = None;
}
pub fn observe(&mut self, value: u64, at: Instant) -> CounterDelta {
let Some(previous) = self.last.replace(Reading { value, at }) else {
return CounterDelta::FirstSample;
};
let elapsed = at.saturating_duration_since(previous.at);
match forward_delta(previous.value, value, self.width) {
Some(delta) => CounterDelta::Advanced {
delta,
elapsed,
wrapped: value < previous.value,
},
None => CounterDelta::Reset,
}
}
pub fn rate(&mut self, value: u64, at: Instant) -> MetricState<Rate> {
self.observe(value, at).rate()
}
}
impl DeltaTracker for CounterTracker {
type Config = CounterWidth;
type Reading = u64;
type Value = Rate;
fn with_config(config: Self::Config) -> Self {
Self::new(config)
}
fn observe_reading(&mut self, reading: Self::Reading, at: Instant) -> MetricState<Self::Value> {
self.rate(reading, at)
}
fn last_observed_at(&self) -> Option<Instant> {
self.last.map(|reading| reading.at)
}
fn forget_baseline(&mut self) {
self.last = None;
}
}
#[cfg(test)]
mod tests {
use super::*;
fn origin() -> Instant {
Instant::now()
}
fn assert_rate(state: &MetricState<Rate>, expected: f64) {
let actual = state
.fresh()
.expect("expected a measured rate")
.per_second();
assert!(
(actual - expected).abs() < 1e-6,
"expected {expected}/s, got {actual}/s"
);
}
#[test]
fn a_first_sample_is_warming_up_and_not_zero() {
let mut tracker = CounterTracker::new(CounterWidth::Bits64);
let state = tracker.rate(4_096, origin());
assert!(state.is_warming_up());
assert_eq!(state.fresh(), None);
assert_ne!(state, MetricState::Available(Rate::ZERO));
}
#[test]
fn a_second_sample_divides_by_the_real_interval() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits64);
assert_eq!(tracker.observe(1_000, t0), CounterDelta::FirstSample);
let state = tracker.rate(3_000, t0 + Duration::from_secs(2));
assert_rate(&state, 1_000.0);
}
#[test]
fn the_same_delta_over_different_intervals_gives_different_rates() {
let t0 = origin();
let mut fast = CounterTracker::new(CounterWidth::Bits64);
let mut slow = CounterTracker::new(CounterWidth::Bits64);
fast.rate(0, t0);
slow.rate(0, t0);
let half = fast.rate(1_000, t0 + Duration::from_millis(500));
let double = slow.rate(1_000, t0 + Duration::from_secs(2));
assert_rate(&half, 2_000.0);
assert_rate(&double, 500.0);
}
#[test]
fn a_counter_that_does_not_move_is_a_real_zero_rate() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits64);
tracker.rate(7_777, t0);
let state = tracker.rate(7_777, t0 + Duration::from_secs(1));
assert_eq!(state, MetricState::Available(Rate::ZERO));
}
#[test]
fn zero_elapsed_is_warming_up_rather_than_a_division_by_zero() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits64);
tracker.rate(100, t0);
let state = tracker.rate(900, t0);
assert!(state.is_warming_up());
let mut totals = CounterTracker::new(CounterWidth::Bits64);
totals.observe(100, t0);
assert_eq!(totals.observe(900, t0).advanced_by(), Some(800));
}
#[test]
fn a_reversed_instant_cannot_produce_a_negative_or_huge_rate() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits64);
tracker.rate(0, t0 + Duration::from_secs(10));
let state = tracker.rate(1_000_000, t0);
assert!(state.is_warming_up());
}
#[test]
fn a_backwards_counter_of_unknown_width_is_a_typed_reset() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Unknown);
tracker.rate(9_000_000, t0);
let state = tracker.rate(12, t0 + Duration::from_secs(1));
assert_eq!(
state,
MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset)
);
assert_eq!(state.fresh(), None);
}
#[test]
fn the_sample_after_a_reset_is_valid_again() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Unknown);
tracker.rate(9_000_000, t0);
let reset = tracker.rate(12, t0 + Duration::from_secs(1));
assert!(!reset.is_available());
let recovered = tracker.rate(1_012, t0 + Duration::from_secs(2));
assert_rate(&recovered, 1_000.0);
}
#[test]
fn a_reset_never_reports_a_rate_derived_from_the_new_value() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Unknown);
tracker.rate(9_000_000, t0);
let delta = tracker.observe(12, t0 + Duration::from_secs(1));
assert_eq!(delta, CounterDelta::Reset);
assert_eq!(delta.advanced_by(), None);
}
#[test]
fn a_known_width_counter_wraps_instead_of_resetting() {
let t0 = origin();
let ceiling = u64::from(u32::MAX) + 1;
let previous = ceiling - 300;
let mut tracker = CounterTracker::new(CounterWidth::Bits32);
tracker.observe(previous, t0);
let delta = tracker.observe(700, t0 + Duration::from_secs(1));
assert_eq!(
delta,
CounterDelta::Advanced {
delta: 1_000,
elapsed: Duration::from_secs(1),
wrapped: true,
}
);
assert_rate(&delta.rate(), 1_000.0);
assert!(delta.wrapped());
}
#[test]
fn the_same_movement_is_a_reset_when_the_width_is_unknown() {
let t0 = origin();
let previous = u64::from(u32::MAX) + 1 - 300;
let mut tracker = CounterTracker::new(CounterWidth::Unknown);
tracker.observe(previous, t0);
assert_eq!(
tracker.observe(700, t0 + Duration::from_secs(1)),
CounterDelta::Reset
);
}
#[test]
fn a_wrap_at_the_exact_boundary_is_reconstructed_exactly() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits32);
tracker.observe(u64::from(u32::MAX), t0);
assert_eq!(
tracker
.observe(0, t0 + Duration::from_secs(1))
.advanced_by(),
Some(1),
"u32::MAX -> 0 is a single step forward"
);
}
#[test]
fn a_small_backwards_move_is_a_reset_even_at_a_known_width() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits32);
tracker.observe(3_000_000_000, t0);
assert_eq!(
tracker.observe(2_999_000_000, t0 + Duration::from_secs(1)),
CounterDelta::Reset
);
}
#[test]
fn a_reading_outside_the_declared_width_is_a_reset_not_a_wrap() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits32);
tracker.observe(u64::from(u32::MAX) + 5_000, t0);
assert_eq!(
tracker.observe(10, t0 + Duration::from_secs(1)),
CounterDelta::Reset
);
}
#[test]
fn a_wrapped_delta_can_never_exceed_half_the_counter_range() {
let half = 1u64 << 31;
let t0 = origin();
for previous in [u64::from(u32::MAX), 3_000_000_000, half + 1] {
for current in [0, 1, 1_000, half - 1] {
let mut tracker = CounterTracker::new(CounterWidth::Bits32);
tracker.observe(previous, t0);
if let Some(delta) = tracker
.observe(current, t0 + Duration::from_secs(1))
.advanced_by()
&& current < previous
{
assert!(delta < half, "{previous} -> {current} produced {delta}");
}
}
}
}
#[test]
fn a_sixty_four_bit_counter_still_rejects_an_absurd_backwards_jump() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits64);
tracker.observe(1_000_000, t0);
assert_eq!(
tracker.observe(9, t0 + Duration::from_secs(1)),
CounterDelta::Reset
);
}
#[test]
fn forgetting_the_baseline_makes_the_next_reading_warm_up() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits64);
tracker.rate(1_000, t0);
assert!(!tracker.is_warming_up());
tracker.forget_baseline();
assert!(tracker.is_warming_up());
assert_eq!(tracker.last_value(), None);
assert_eq!(tracker.last_observed_at(), None);
assert!(
tracker
.rate(500_000, t0 + Duration::from_secs(1))
.is_warming_up(),
"a dropped baseline must not be reconstructed from the old value"
);
}
#[test]
fn the_baseline_tracks_the_most_recent_reading() {
let t0 = origin();
let at = t0 + Duration::from_secs(3);
let mut tracker = CounterTracker::new(CounterWidth::Bits32);
tracker.observe(42, t0);
tracker.observe(84, at);
assert_eq!(tracker.last_value(), Some(84));
assert_eq!(tracker.last_observed_at(), Some(at));
assert_eq!(tracker.width(), CounterWidth::Bits32);
}
#[test]
fn widths_report_their_own_limits() {
assert_eq!(CounterWidth::Unknown.bits(), None);
assert_eq!(CounterWidth::Unknown.max_value(), None);
assert_eq!(CounterWidth::Bits32.bits(), Some(32));
assert_eq!(CounterWidth::Bits32.max_value(), Some(u64::from(u32::MAX)));
assert_eq!(CounterWidth::Bits64.bits(), Some(64));
assert_eq!(CounterWidth::Bits64.max_value(), Some(u64::MAX));
assert_eq!(CounterWidth::default(), CounterWidth::Unknown);
}
#[test]
fn a_forward_move_is_never_treated_as_a_wrap() {
let t0 = origin();
let mut tracker = CounterTracker::new(CounterWidth::Bits32);
tracker.observe(10, t0);
let delta = tracker.observe(4_000_000_000, t0 + Duration::from_secs(1));
assert_eq!(delta.advanced_by(), Some(3_999_999_990));
assert!(!delta.wrapped());
}
}