use core::fmt;
use core::hash::Hash;
use core::time::Duration;
use std::collections::HashMap;
use std::time::Instant;
use crate::model::{MetricState, ProcessIdentity, UnavailableReason};
use crate::rates::counter::CounterTracker;
use crate::rates::cpu::ProcessCpuTracker;
pub const DEFAULT_MAX_TRACKED: usize = 16_384;
pub trait DeltaTracker {
type Config: Copy + fmt::Debug;
type Reading;
type Value;
fn with_config(config: Self::Config) -> Self;
fn observe_reading(&mut self, reading: Self::Reading, at: Instant) -> MetricState<Self::Value>;
fn last_observed_at(&self) -> Option<Instant>;
fn forget_baseline(&mut self);
}
#[derive(Debug)]
pub struct KeyedTrackers<K, T: DeltaTracker> {
config: T::Config,
max_tracked: usize,
max_gap: Option<Duration>,
evictions: u64,
entries: HashMap<K, T>,
}
impl<K, T> KeyedTrackers<K, T>
where
K: Clone + Eq + Hash,
T: DeltaTracker,
{
#[must_use]
pub fn new(config: T::Config) -> Self {
Self {
config,
max_tracked: DEFAULT_MAX_TRACKED,
max_gap: None,
evictions: 0,
entries: HashMap::new(),
}
}
#[must_use]
pub fn with_max_tracked(mut self, max_tracked: usize) -> Self {
self.max_tracked = max_tracked;
self
}
#[must_use]
pub fn with_max_gap(mut self, max_gap: Duration) -> Self {
self.max_gap = Some(max_gap);
self
}
pub fn observe(&mut self, key: K, reading: T::Reading, at: Instant) -> MetricState<T::Value> {
if let Some(tracker) = self.entries.get_mut(&key) {
let gapped = match (self.max_gap, DeltaTracker::last_observed_at(tracker)) {
(Some(max_gap), Some(previous)) => at.saturating_duration_since(previous) > max_gap,
_ => false,
};
if !gapped {
return tracker.observe_reading(reading, at);
}
tracker.forget_baseline();
let _ = tracker.observe_reading(reading, at);
return MetricState::TemporarilyUnavailable(UnavailableReason::DeviceDisappeared);
}
if self.entries.len() >= self.max_tracked && !self.evict_oldest() {
return MetricState::TemporarilyUnavailable(UnavailableReason::SkippedUnderLoad);
}
let config = self.config;
self.entries
.entry(key)
.or_insert_with(|| T::with_config(config))
.observe_reading(reading, at)
}
pub fn forget(&mut self, key: &K) -> bool {
self.entries.remove(key).is_some()
}
pub fn retain(&mut self, mut keep: impl FnMut(&K) -> bool) -> usize {
let before = self.entries.len();
self.entries.retain(|key, _| keep(key));
before.saturating_sub(self.entries.len())
}
pub fn prune_idle(&mut self, now: Instant, max_idle: Duration) -> usize {
let before = self.entries.len();
self.entries.retain(|_, tracker| {
DeltaTracker::last_observed_at(tracker)
.is_some_and(|at| now.saturating_duration_since(at) <= max_idle)
});
let dropped = before.saturating_sub(self.entries.len());
self.evictions = self
.evictions
.saturating_add(u64::try_from(dropped).unwrap_or(u64::MAX));
dropped
}
pub fn clear(&mut self) {
self.entries.clear();
}
#[must_use]
pub fn len(&self) -> usize {
self.entries.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
#[must_use]
pub fn contains_key(&self, key: &K) -> bool {
self.entries.contains_key(key)
}
#[must_use]
pub fn tracker(&self, key: &K) -> Option<&T> {
self.entries.get(key)
}
#[must_use]
pub const fn max_tracked(&self) -> usize {
self.max_tracked
}
#[must_use]
pub const fn evictions(&self) -> u64 {
self.evictions
}
fn evict_oldest(&mut self) -> bool {
let victim = self
.entries
.iter()
.min_by_key(|(_, tracker)| DeltaTracker::last_observed_at(*tracker))
.map(|(key, _)| key.clone());
let Some(key) = victim else {
return false;
};
self.entries.remove(&key);
self.evictions = self.evictions.saturating_add(1);
true
}
}
impl<K, T> Default for KeyedTrackers<K, T>
where
K: Clone + Eq + Hash,
T: DeltaTracker,
T::Config: Default,
{
fn default() -> Self {
Self::new(T::Config::default())
}
}
pub type KeyedRateTrackers<K> = KeyedTrackers<K, CounterTracker>;
pub type KeyedProcessCpuTrackers = KeyedTrackers<ProcessIdentity, ProcessCpuTracker>;
#[cfg(test)]
mod tests {
use super::*;
use crate::rates::counter::CounterWidth;
use crate::rates::cpu::{CpuTimeTotals, SystemCpuTracker};
use crate::units::{Percent, Rate};
fn origin() -> Instant {
Instant::now()
}
fn secs(seconds: u64) -> Duration {
Duration::from_secs(seconds)
}
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"
);
}
fn bytes() -> KeyedRateTrackers<&'static str> {
KeyedRateTrackers::new(CounterWidth::Bits64)
}
#[test]
fn an_unseen_key_warms_up_instead_of_reporting_zero() {
let mut set = bytes();
let state = set.observe("eth0", 4_096, origin());
assert!(state.is_warming_up());
assert_eq!(set.len(), 1);
assert!(set.contains_key(&"eth0"));
}
#[test]
fn a_second_reading_for_a_key_yields_a_rate_over_the_real_interval() {
let t0 = origin();
let mut set = bytes();
set.observe("eth0", 1_000, t0);
let state = set.observe("eth0", 2_000, t0 + Duration::from_millis(500));
assert_rate(&state, 2_000.0);
}
#[test]
fn keys_keep_independent_baselines() {
let t0 = origin();
let mut set = bytes();
set.observe("eth0", 0, t0);
set.observe("wlan0", 1_000_000, t0);
let eth0 = set.observe("eth0", 100, t0 + secs(1));
let wlan0 = set.observe("wlan0", 1_000_300, t0 + secs(1));
assert_rate(ð0, 100.0);
assert_rate(&wlan0, 300.0);
}
#[test]
fn a_forgotten_key_rebaselines_when_it_reappears() {
let t0 = origin();
let mut set = bytes();
set.observe("sdb", 900_000, t0);
assert!(set.forget(&"sdb"));
assert!(!set.forget(&"sdb"), "forgetting twice is not an error");
assert!(set.is_empty());
let state = set.observe("sdb", 40_000, t0 + secs(1));
assert!(state.is_warming_up());
}
#[test]
fn retain_drops_absent_keys_so_a_reappearance_cannot_produce_a_bogus_delta() {
let t0 = origin();
let mut set = bytes();
set.observe("eth0", 1_000, t0);
set.observe("tun0", 5_000, t0);
assert_eq!(set.retain(|key| *key == "eth0"), 1);
assert_eq!(set.len(), 1);
assert_eq!(set.evictions(), 0, "deliberate removal is not an eviction");
let state = set.observe("tun0", 9_000_000, t0 + secs(300));
assert!(state.is_warming_up());
let recovered = set.observe("tun0", 9_000_100, t0 + secs(301));
assert_rate(&recovered, 100.0);
}
#[test]
fn a_gap_longer_than_the_guard_is_reported_as_a_disappearance() {
let t0 = origin();
let mut set = bytes().with_max_gap(secs(3));
set.observe("eth0", 1_000, t0);
let gapped = set.observe("eth0", 9_000_000, t0 + secs(10));
assert_eq!(
gapped,
MetricState::TemporarilyUnavailable(UnavailableReason::DeviceDisappeared)
);
let recovered = set.observe("eth0", 9_000_500, t0 + secs(11));
assert_rate(&recovered, 500.0);
}
#[test]
fn a_gap_inside_the_guard_is_an_ordinary_sample() {
let t0 = origin();
let mut set = bytes().with_max_gap(secs(3));
set.observe("eth0", 1_000, t0);
let state = set.observe("eth0", 3_000, t0 + secs(2));
assert_rate(&state, 1_000.0);
}
#[test]
fn without_a_guard_no_gap_is_ever_treated_as_a_disappearance() {
let t0 = origin();
let mut set = bytes();
set.observe("eth0", 1_000, t0);
let state = set.observe("eth0", 3_000, t0 + secs(600));
assert!(state.is_available());
}
#[test]
fn the_set_never_grows_past_its_cap_as_keys_churn() {
let t0 = origin();
let mut set: KeyedRateTrackers<u64> =
KeyedRateTrackers::new(CounterWidth::Bits64).with_max_tracked(64);
for pid in 0..10_000u64 {
set.observe(pid, pid, t0 + Duration::from_millis(pid));
}
assert_eq!(set.max_tracked(), 64);
assert!(set.len() <= 64, "len was {}", set.len());
assert!(set.evictions() > 0, "eviction must be reported");
}
#[test]
fn the_least_recently_observed_key_is_evicted_first() {
let t0 = origin();
let mut set = bytes().with_max_tracked(2);
set.observe("oldest", 1, t0);
set.observe("newer", 1, t0 + secs(5));
set.observe("newest", 1, t0 + secs(10));
assert!(!set.contains_key(&"oldest"));
assert!(set.contains_key(&"newer"));
assert!(set.contains_key(&"newest"));
assert_eq!(set.evictions(), 1);
}
#[test]
fn a_zero_cap_reports_skipped_rather_than_zero() {
let mut set = bytes().with_max_tracked(0);
let state = set.observe("eth0", 1_000, origin());
assert_eq!(
state,
MetricState::TemporarilyUnavailable(UnavailableReason::SkippedUnderLoad)
);
assert!(set.is_empty());
}
#[test]
fn prune_idle_drops_keys_that_stopped_reporting() {
let t0 = origin();
let mut set = bytes();
set.observe("alive", 0, t0);
set.observe("exited", 0, t0);
set.observe("alive", 100, t0 + secs(30));
assert_eq!(set.prune_idle(t0 + secs(30), secs(5)), 1);
assert!(set.contains_key(&"alive"));
assert!(!set.contains_key(&"exited"));
assert_eq!(set.evictions(), 1);
}
#[test]
fn pruning_an_empty_set_is_a_no_op() {
let mut set = bytes();
assert_eq!(set.prune_idle(origin(), secs(1)), 0);
assert!(set.is_empty());
}
#[test]
fn pruning_keeps_a_key_observed_exactly_at_the_idle_limit() {
let t0 = origin();
let mut set = bytes();
set.observe("eth0", 0, t0);
assert_eq!(set.prune_idle(t0 + secs(1), secs(1)), 0);
assert!(set.contains_key(&"eth0"));
assert_eq!(
set.prune_idle(t0 + secs(1) + Duration::from_nanos(1), secs(1)),
1
);
}
#[test]
fn clearing_drops_every_baseline() {
let t0 = origin();
let mut set = bytes();
set.observe("a", 1, t0);
set.observe("b", 1, t0);
set.clear();
assert!(set.is_empty());
assert!(set.observe("a", 1_000_000, t0 + secs(1)).is_warming_up());
}
#[test]
fn the_tracker_behind_a_key_is_inspectable() {
let t0 = origin();
let mut set = bytes();
set.observe("eth0", 4_096, t0);
let tracker = set.tracker(&"eth0").expect("tracked");
assert_eq!(tracker.last_value(), Some(4_096));
assert_eq!(tracker.width(), CounterWidth::Bits64);
assert!(set.tracker(&"missing").is_none());
}
#[test]
fn a_default_set_uses_the_default_counter_width_and_cap() {
let set: KeyedRateTrackers<&'static str> = KeyedRateTrackers::default();
assert_eq!(set.max_tracked(), DEFAULT_MAX_TRACKED);
assert!(set.is_empty());
}
#[test]
fn a_known_width_wrap_still_works_through_the_keyed_set() {
let t0 = origin();
let mut set: KeyedRateTrackers<&'static str> = KeyedRateTrackers::new(CounterWidth::Bits32);
let previous = u64::from(u32::MAX) - 99;
set.observe("eth0", previous, t0);
let state = set.observe("eth0", 400, t0 + secs(1));
assert_rate(&state, 500.0);
}
#[test]
fn process_cpu_trackers_are_keyed_on_identity_so_a_reused_pid_rebaselines() {
let t0 = origin();
let original = ProcessIdentity::new(4_242, 900_100);
let recycled = ProcessIdentity::new(4_242, 977_400);
let mut set = KeyedProcessCpuTrackers::default();
set.observe(original, secs(30), t0);
let measured = set.observe(original, secs(31), t0 + secs(1));
assert!(
(measured
.fresh()
.copied()
.map(Percent::value)
.expect("measured")
- 100.0)
.abs()
< f32::EPSILON
);
let reused = set.observe(recycled, Duration::from_millis(10), t0 + secs(2));
assert!(
reused.is_warming_up(),
"a recycled PID must warm up, not report a negative or reset delta"
);
assert_eq!(set.len(), 2, "the two identities are distinct keys");
}
#[test]
fn exited_processes_are_prunable_so_the_set_stays_bounded() {
let t0 = origin();
let mut set = KeyedProcessCpuTrackers::default();
for pid in 0..500u32 {
set.observe(ProcessIdentity::new(pid, 1), Duration::ZERO, t0);
}
let survivor = ProcessIdentity::new(1, 1);
set.observe(survivor, Duration::from_millis(1), t0 + secs(10));
assert_eq!(set.prune_idle(t0 + secs(10), secs(2)), 499);
assert_eq!(set.len(), 1);
assert!(set.contains_key(&survivor));
}
#[test]
fn per_core_cpu_trackers_share_the_set_without_a_counter_width() {
let t0 = origin();
let mut set: KeyedTrackers<u16, SystemCpuTracker> = KeyedTrackers::default();
for core in 0..4u16 {
assert!(
set.observe(core, CpuTimeTotals::new(secs(0), secs(0)), t0)
.is_warming_up()
);
}
let busy = set.observe(0, CpuTimeTotals::new(secs(1), secs(3)), t0 + secs(4));
assert!(
(busy.fresh().copied().map(Percent::value).expect("measured") - 25.0).abs()
< f32::EPSILON
);
assert_eq!(set.len(), 4);
}
}