1use core::fmt;
13use core::hash::Hash;
14use core::time::Duration;
15use std::collections::HashMap;
16use std::time::Instant;
17
18use crate::model::{MetricState, ProcessIdentity, UnavailableReason};
19use crate::rates::counter::CounterTracker;
20use crate::rates::cpu::ProcessCpuTracker;
21
22pub const DEFAULT_MAX_TRACKED: usize = 16_384;
29
30pub trait DeltaTracker {
38 type Config: Copy + fmt::Debug;
43 type Reading;
45 type Value;
47
48 fn with_config(config: Self::Config) -> Self;
50
51 fn observe_reading(&mut self, reading: Self::Reading, at: Instant) -> MetricState<Self::Value>;
53
54 fn last_observed_at(&self) -> Option<Instant>;
56
57 fn forget_baseline(&mut self);
60}
61
62#[derive(Debug)]
95pub struct KeyedTrackers<K, T: DeltaTracker> {
96 config: T::Config,
97 max_tracked: usize,
98 max_gap: Option<Duration>,
99 evictions: u64,
100 entries: HashMap<K, T>,
101}
102
103impl<K, T> KeyedTrackers<K, T>
104where
105 K: Clone + Eq + Hash,
106 T: DeltaTracker,
107{
108 #[must_use]
110 pub fn new(config: T::Config) -> Self {
111 Self {
112 config,
113 max_tracked: DEFAULT_MAX_TRACKED,
114 max_gap: None,
115 evictions: 0,
116 entries: HashMap::new(),
117 }
118 }
119
120 #[must_use]
125 pub fn with_max_tracked(mut self, max_tracked: usize) -> Self {
126 self.max_tracked = max_tracked;
127 self
128 }
129
130 #[must_use]
141 pub fn with_max_gap(mut self, max_gap: Duration) -> Self {
142 self.max_gap = Some(max_gap);
143 self
144 }
145
146 pub fn observe(&mut self, key: K, reading: T::Reading, at: Instant) -> MetricState<T::Value> {
151 if let Some(tracker) = self.entries.get_mut(&key) {
152 let gapped = match (self.max_gap, DeltaTracker::last_observed_at(tracker)) {
153 (Some(max_gap), Some(previous)) => at.saturating_duration_since(previous) > max_gap,
154 _ => false,
155 };
156 if !gapped {
157 return tracker.observe_reading(reading, at);
158 }
159 tracker.forget_baseline();
163 let _ = tracker.observe_reading(reading, at);
164 return MetricState::TemporarilyUnavailable(UnavailableReason::DeviceDisappeared);
165 }
166
167 if self.entries.len() >= self.max_tracked && !self.evict_oldest() {
168 return MetricState::TemporarilyUnavailable(UnavailableReason::SkippedUnderLoad);
171 }
172 let config = self.config;
173 self.entries
174 .entry(key)
175 .or_insert_with(|| T::with_config(config))
176 .observe_reading(reading, at)
177 }
178
179 pub fn forget(&mut self, key: &K) -> bool {
187 self.entries.remove(key).is_some()
188 }
189
190 pub fn retain(&mut self, mut keep: impl FnMut(&K) -> bool) -> usize {
196 let before = self.entries.len();
197 self.entries.retain(|key, _| keep(key));
198 before.saturating_sub(self.entries.len())
199 }
200
201 pub fn prune_idle(&mut self, now: Instant, max_idle: Duration) -> usize {
208 let before = self.entries.len();
209 self.entries.retain(|_, tracker| {
210 DeltaTracker::last_observed_at(tracker)
211 .is_some_and(|at| now.saturating_duration_since(at) <= max_idle)
212 });
213 let dropped = before.saturating_sub(self.entries.len());
214 self.evictions = self
215 .evictions
216 .saturating_add(u64::try_from(dropped).unwrap_or(u64::MAX));
217 dropped
218 }
219
220 pub fn clear(&mut self) {
222 self.entries.clear();
223 }
224
225 #[must_use]
227 pub fn len(&self) -> usize {
228 self.entries.len()
229 }
230
231 #[must_use]
233 pub fn is_empty(&self) -> bool {
234 self.entries.is_empty()
235 }
236
237 #[must_use]
239 pub fn contains_key(&self, key: &K) -> bool {
240 self.entries.contains_key(key)
241 }
242
243 #[must_use]
245 pub fn tracker(&self, key: &K) -> Option<&T> {
246 self.entries.get(key)
247 }
248
249 #[must_use]
251 pub const fn max_tracked(&self) -> usize {
252 self.max_tracked
253 }
254
255 #[must_use]
261 pub const fn evictions(&self) -> u64 {
262 self.evictions
263 }
264
265 fn evict_oldest(&mut self) -> bool {
272 let victim = self
273 .entries
274 .iter()
275 .min_by_key(|(_, tracker)| DeltaTracker::last_observed_at(*tracker))
276 .map(|(key, _)| key.clone());
277 let Some(key) = victim else {
278 return false;
279 };
280 self.entries.remove(&key);
281 self.evictions = self.evictions.saturating_add(1);
282 true
283 }
284}
285
286impl<K, T> Default for KeyedTrackers<K, T>
287where
288 K: Clone + Eq + Hash,
289 T: DeltaTracker,
290 T::Config: Default,
291{
292 fn default() -> Self {
293 Self::new(T::Config::default())
294 }
295}
296
297pub type KeyedRateTrackers<K> = KeyedTrackers<K, CounterTracker>;
303
304pub type KeyedProcessCpuTrackers = KeyedTrackers<ProcessIdentity, ProcessCpuTracker>;
311
312#[cfg(test)]
313mod tests {
314 use super::*;
315 use crate::rates::counter::CounterWidth;
316 use crate::rates::cpu::{CpuTimeTotals, SystemCpuTracker};
317 use crate::units::{Percent, Rate};
318
319 fn origin() -> Instant {
320 Instant::now()
321 }
322
323 fn secs(seconds: u64) -> Duration {
324 Duration::from_secs(seconds)
325 }
326
327 fn assert_rate(state: &MetricState<Rate>, expected: f64) {
330 let actual = state
331 .fresh()
332 .expect("expected a measured rate")
333 .per_second();
334 assert!(
335 (actual - expected).abs() < 1e-6,
336 "expected {expected}/s, got {actual}/s"
337 );
338 }
339
340 fn bytes() -> KeyedRateTrackers<&'static str> {
341 KeyedRateTrackers::new(CounterWidth::Bits64)
342 }
343
344 #[test]
345 fn an_unseen_key_warms_up_instead_of_reporting_zero() {
346 let mut set = bytes();
347 let state = set.observe("eth0", 4_096, origin());
348 assert!(state.is_warming_up());
349 assert_eq!(set.len(), 1);
350 assert!(set.contains_key(&"eth0"));
351 }
352
353 #[test]
354 fn a_second_reading_for_a_key_yields_a_rate_over_the_real_interval() {
355 let t0 = origin();
356 let mut set = bytes();
357 set.observe("eth0", 1_000, t0);
358 let state = set.observe("eth0", 2_000, t0 + Duration::from_millis(500));
359 assert_rate(&state, 2_000.0);
360 }
361
362 #[test]
363 fn keys_keep_independent_baselines() {
364 let t0 = origin();
365 let mut set = bytes();
366 set.observe("eth0", 0, t0);
367 set.observe("wlan0", 1_000_000, t0);
368
369 let eth0 = set.observe("eth0", 100, t0 + secs(1));
370 let wlan0 = set.observe("wlan0", 1_000_300, t0 + secs(1));
371 assert_rate(ð0, 100.0);
372 assert_rate(&wlan0, 300.0);
373 }
374
375 #[test]
376 fn a_forgotten_key_rebaselines_when_it_reappears() {
377 let t0 = origin();
378 let mut set = bytes();
379 set.observe("sdb", 900_000, t0);
380 assert!(set.forget(&"sdb"));
381 assert!(!set.forget(&"sdb"), "forgetting twice is not an error");
382 assert!(set.is_empty());
383
384 let state = set.observe("sdb", 40_000, t0 + secs(1));
387 assert!(state.is_warming_up());
388 }
389
390 #[test]
391 fn retain_drops_absent_keys_so_a_reappearance_cannot_produce_a_bogus_delta() {
392 let t0 = origin();
393 let mut set = bytes();
394 set.observe("eth0", 1_000, t0);
395 set.observe("tun0", 5_000, t0);
396
397 assert_eq!(set.retain(|key| *key == "eth0"), 1);
399 assert_eq!(set.len(), 1);
400 assert_eq!(set.evictions(), 0, "deliberate removal is not an eviction");
401
402 let state = set.observe("tun0", 9_000_000, t0 + secs(300));
404 assert!(state.is_warming_up());
405 let recovered = set.observe("tun0", 9_000_100, t0 + secs(301));
406 assert_rate(&recovered, 100.0);
407 }
408
409 #[test]
410 fn a_gap_longer_than_the_guard_is_reported_as_a_disappearance() {
411 let t0 = origin();
412 let mut set = bytes().with_max_gap(secs(3));
413 set.observe("eth0", 1_000, t0);
414
415 let gapped = set.observe("eth0", 9_000_000, t0 + secs(10));
418 assert_eq!(
419 gapped,
420 MetricState::TemporarilyUnavailable(UnavailableReason::DeviceDisappeared)
421 );
422
423 let recovered = set.observe("eth0", 9_000_500, t0 + secs(11));
425 assert_rate(&recovered, 500.0);
426 }
427
428 #[test]
429 fn a_gap_inside_the_guard_is_an_ordinary_sample() {
430 let t0 = origin();
431 let mut set = bytes().with_max_gap(secs(3));
432 set.observe("eth0", 1_000, t0);
433 let state = set.observe("eth0", 3_000, t0 + secs(2));
434 assert_rate(&state, 1_000.0);
435 }
436
437 #[test]
438 fn without_a_guard_no_gap_is_ever_treated_as_a_disappearance() {
439 let t0 = origin();
441 let mut set = bytes();
442 set.observe("eth0", 1_000, t0);
443 let state = set.observe("eth0", 3_000, t0 + secs(600));
444 assert!(state.is_available());
445 }
446
447 #[test]
448 fn the_set_never_grows_past_its_cap_as_keys_churn() {
449 let t0 = origin();
452 let mut set: KeyedRateTrackers<u64> =
453 KeyedRateTrackers::new(CounterWidth::Bits64).with_max_tracked(64);
454 for pid in 0..10_000u64 {
455 set.observe(pid, pid, t0 + Duration::from_millis(pid));
456 }
457 assert_eq!(set.max_tracked(), 64);
458 assert!(set.len() <= 64, "len was {}", set.len());
459 assert!(set.evictions() > 0, "eviction must be reported");
460 }
461
462 #[test]
463 fn the_least_recently_observed_key_is_evicted_first() {
464 let t0 = origin();
465 let mut set = bytes().with_max_tracked(2);
466 set.observe("oldest", 1, t0);
467 set.observe("newer", 1, t0 + secs(5));
468
469 set.observe("newest", 1, t0 + secs(10));
470 assert!(!set.contains_key(&"oldest"));
471 assert!(set.contains_key(&"newer"));
472 assert!(set.contains_key(&"newest"));
473 assert_eq!(set.evictions(), 1);
474 }
475
476 #[test]
477 fn a_zero_cap_reports_skipped_rather_than_zero() {
478 let mut set = bytes().with_max_tracked(0);
479 let state = set.observe("eth0", 1_000, origin());
480 assert_eq!(
481 state,
482 MetricState::TemporarilyUnavailable(UnavailableReason::SkippedUnderLoad)
483 );
484 assert!(set.is_empty());
485 }
486
487 #[test]
488 fn prune_idle_drops_keys_that_stopped_reporting() {
489 let t0 = origin();
490 let mut set = bytes();
491 set.observe("alive", 0, t0);
492 set.observe("exited", 0, t0);
493
494 set.observe("alive", 100, t0 + secs(30));
496 assert_eq!(set.prune_idle(t0 + secs(30), secs(5)), 1);
497 assert!(set.contains_key(&"alive"));
498 assert!(!set.contains_key(&"exited"));
499 assert_eq!(set.evictions(), 1);
500 }
501
502 #[test]
503 fn pruning_an_empty_set_is_a_no_op() {
504 let mut set = bytes();
505 assert_eq!(set.prune_idle(origin(), secs(1)), 0);
506 assert!(set.is_empty());
507 }
508
509 #[test]
510 fn pruning_keeps_a_key_observed_exactly_at_the_idle_limit() {
511 let t0 = origin();
514 let mut set = bytes();
515 set.observe("eth0", 0, t0);
516 assert_eq!(set.prune_idle(t0 + secs(1), secs(1)), 0);
517 assert!(set.contains_key(&"eth0"));
518 assert_eq!(
519 set.prune_idle(t0 + secs(1) + Duration::from_nanos(1), secs(1)),
520 1
521 );
522 }
523
524 #[test]
525 fn clearing_drops_every_baseline() {
526 let t0 = origin();
527 let mut set = bytes();
528 set.observe("a", 1, t0);
529 set.observe("b", 1, t0);
530 set.clear();
531 assert!(set.is_empty());
532 assert!(set.observe("a", 1_000_000, t0 + secs(1)).is_warming_up());
533 }
534
535 #[test]
536 fn the_tracker_behind_a_key_is_inspectable() {
537 let t0 = origin();
538 let mut set = bytes();
539 set.observe("eth0", 4_096, t0);
540 let tracker = set.tracker(&"eth0").expect("tracked");
541 assert_eq!(tracker.last_value(), Some(4_096));
542 assert_eq!(tracker.width(), CounterWidth::Bits64);
543 assert!(set.tracker(&"missing").is_none());
544 }
545
546 #[test]
547 fn a_default_set_uses_the_default_counter_width_and_cap() {
548 let set: KeyedRateTrackers<&'static str> = KeyedRateTrackers::default();
549 assert_eq!(set.max_tracked(), DEFAULT_MAX_TRACKED);
550 assert!(set.is_empty());
551 }
552
553 #[test]
554 fn a_known_width_wrap_still_works_through_the_keyed_set() {
555 let t0 = origin();
556 let mut set: KeyedRateTrackers<&'static str> = KeyedRateTrackers::new(CounterWidth::Bits32);
557 let previous = u64::from(u32::MAX) - 99;
558 set.observe("eth0", previous, t0);
559 let state = set.observe("eth0", 400, t0 + secs(1));
560 assert_rate(&state, 500.0);
561 }
562
563 #[test]
564 fn process_cpu_trackers_are_keyed_on_identity_so_a_reused_pid_rebaselines() {
565 let t0 = origin();
568 let original = ProcessIdentity::new(4_242, 900_100);
569 let recycled = ProcessIdentity::new(4_242, 977_400);
570
571 let mut set = KeyedProcessCpuTrackers::default();
572 set.observe(original, secs(30), t0);
573 let measured = set.observe(original, secs(31), t0 + secs(1));
574 assert!(
575 (measured
576 .fresh()
577 .copied()
578 .map(Percent::value)
579 .expect("measured")
580 - 100.0)
581 .abs()
582 < f32::EPSILON
583 );
584
585 let reused = set.observe(recycled, Duration::from_millis(10), t0 + secs(2));
586 assert!(
587 reused.is_warming_up(),
588 "a recycled PID must warm up, not report a negative or reset delta"
589 );
590 assert_eq!(set.len(), 2, "the two identities are distinct keys");
591 }
592
593 #[test]
594 fn exited_processes_are_prunable_so_the_set_stays_bounded() {
595 let t0 = origin();
596 let mut set = KeyedProcessCpuTrackers::default();
597 for pid in 0..500u32 {
598 set.observe(ProcessIdentity::new(pid, 1), Duration::ZERO, t0);
599 }
600 let survivor = ProcessIdentity::new(1, 1);
601 set.observe(survivor, Duration::from_millis(1), t0 + secs(10));
602
603 assert_eq!(set.prune_idle(t0 + secs(10), secs(2)), 499);
604 assert_eq!(set.len(), 1);
605 assert!(set.contains_key(&survivor));
606 }
607
608 #[test]
609 fn per_core_cpu_trackers_share_the_set_without_a_counter_width() {
610 let t0 = origin();
613 let mut set: KeyedTrackers<u16, SystemCpuTracker> = KeyedTrackers::default();
614 for core in 0..4u16 {
615 assert!(
616 set.observe(core, CpuTimeTotals::new(secs(0), secs(0)), t0)
617 .is_warming_up()
618 );
619 }
620 let busy = set.observe(0, CpuTimeTotals::new(secs(1), secs(3)), t0 + secs(4));
621 assert!(
622 (busy.fresh().copied().map(Percent::value).expect("measured") - 25.0).abs()
623 < f32::EPSILON
624 );
625 assert_eq!(set.len(), 4);
626 }
627}