use core::time::Duration;
use std::time::Instant;
use crate::model::{CpuNormalization, MetricState, UnavailableReason};
use crate::rates::keyed::DeltaTracker;
use crate::units::Percent;
fn percent_of_duration(part: Duration, whole: Duration) -> Option<Percent> {
let part = u64::try_from(part.as_nanos()).ok()?;
let whole = u64::try_from(whole.as_nanos()).ok()?;
Percent::ratio(part, whole)
}
#[derive(Clone, Copy, Debug, Default, Eq, Ord, PartialEq, PartialOrd)]
pub struct CpuTimeTotals {
pub busy: Duration,
pub idle: Duration,
}
impl CpuTimeTotals {
#[must_use]
pub const fn new(busy: Duration, idle: Duration) -> Self {
Self { busy, idle }
}
#[must_use]
pub const fn total(self) -> Duration {
self.busy.saturating_add(self.idle)
}
}
#[derive(Clone, Copy, Debug, Default)]
pub struct SystemCpuTracker {
last: Option<(CpuTimeTotals, Instant)>,
}
impl SystemCpuTracker {
#[must_use]
pub const fn new() -> Self {
Self { last: None }
}
#[must_use]
pub const fn is_warming_up(&self) -> bool {
self.last.is_none()
}
#[must_use]
pub const fn last_observed_at(&self) -> Option<Instant> {
match self.last {
Some((_, at)) => Some(at),
None => None,
}
}
pub fn forget_baseline(&mut self) {
self.last = None;
}
pub fn observe(&mut self, totals: CpuTimeTotals, at: Instant) -> MetricState<Percent> {
let Some((previous, _)) = self.last.replace((totals, at)) else {
return MetricState::WarmingUp;
};
let (Some(busy), Some(idle)) = (
totals.busy.checked_sub(previous.busy),
totals.idle.checked_sub(previous.idle),
) else {
return MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset);
};
match percent_of_duration(busy, busy.saturating_add(idle)) {
Some(percent) => MetricState::Available(percent.clamped_to_100()),
None => MetricState::WarmingUp,
}
}
}
impl DeltaTracker for SystemCpuTracker {
type Config = ();
type Reading = CpuTimeTotals;
type Value = Percent;
fn with_config(_config: Self::Config) -> Self {
Self::new()
}
fn observe_reading(&mut self, reading: Self::Reading, at: Instant) -> MetricState<Self::Value> {
self.observe(reading, at)
}
fn last_observed_at(&self) -> Option<Instant> {
self.last.map(|(_, at)| at)
}
fn forget_baseline(&mut self) {
self.last = None;
}
}
#[derive(Clone, Copy, Debug, Default)]
pub struct ProcessCpuTracker {
last: Option<(Duration, Instant)>,
}
impl ProcessCpuTracker {
#[must_use]
pub const fn new() -> Self {
Self { last: None }
}
#[must_use]
pub const fn is_warming_up(&self) -> bool {
self.last.is_none()
}
#[must_use]
pub const fn last_observed_at(&self) -> Option<Instant> {
match self.last {
Some((_, at)) => Some(at),
None => None,
}
}
pub fn forget_baseline(&mut self) {
self.last = None;
}
pub fn observe(&mut self, cpu_time: Duration, at: Instant) -> MetricState<Percent> {
let Some((previous_cpu, previous_at)) = self.last.replace((cpu_time, at)) else {
return MetricState::WarmingUp;
};
let Some(cpu_delta) = cpu_time.checked_sub(previous_cpu) else {
return MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset);
};
match percent_of_duration(cpu_delta, at.saturating_duration_since(previous_at)) {
Some(percent) => MetricState::Available(percent),
None => MetricState::WarmingUp,
}
}
pub fn observe_normalized(
&mut self,
cpu_time: Duration,
at: Instant,
normalization: CpuNormalization,
logical_cpus: u16,
) -> MetricState<Percent> {
let core_normalized = self.observe(cpu_time, at);
let Some(percent) = core_normalized.fresh().copied() else {
return core_normalized;
};
match normalization.apply(percent, logical_cpus) {
Some(scaled) => MetricState::Available(scaled),
None => MetricState::TemporarilyUnavailable(UnavailableReason::ReadFailed),
}
}
}
impl DeltaTracker for ProcessCpuTracker {
type Config = ();
type Reading = Duration;
type Value = Percent;
fn with_config(_config: Self::Config) -> Self {
Self::new()
}
fn observe_reading(&mut self, reading: Self::Reading, at: Instant) -> MetricState<Self::Value> {
self.observe(reading, at)
}
fn last_observed_at(&self) -> Option<Instant> {
self.last.map(|(_, at)| at)
}
fn forget_baseline(&mut self) {
self.last = None;
}
}
#[cfg(test)]
mod tests {
use super::*;
fn origin() -> Instant {
Instant::now()
}
fn percent(state: &MetricState<Percent>) -> f32 {
state
.fresh()
.expect("expected a measured percentage")
.value()
}
fn secs(seconds: u64) -> Duration {
Duration::from_secs(seconds)
}
#[test]
fn totals_sum_busy_and_idle_without_overflowing() {
let totals = CpuTimeTotals::new(secs(3), secs(5));
assert_eq!(totals.total(), secs(8));
let extreme = CpuTimeTotals::new(Duration::MAX, secs(1));
assert_eq!(extreme.total(), Duration::MAX);
}
#[test]
fn a_first_system_sample_is_warming_up_and_not_zero_percent() {
let mut tracker = SystemCpuTracker::new();
assert!(tracker.is_warming_up());
let state = tracker.observe(CpuTimeTotals::new(secs(100), secs(900)), origin());
assert!(state.is_warming_up());
assert_eq!(state.fresh(), None);
assert!(!tracker.is_warming_up());
}
#[test]
fn system_cpu_is_the_busy_share_of_total_cpu_time() {
let t0 = origin();
let mut tracker = SystemCpuTracker::new();
tracker.observe(CpuTimeTotals::new(secs(0), secs(0)), t0);
let state = tracker.observe(CpuTimeTotals::new(secs(1), secs(7)), t0 + secs(1));
assert!((percent(&state) - 12.5).abs() < f32::EPSILON);
}
#[test]
fn system_cpu_is_an_aggregate_capped_at_one_hundred_percent() {
let t0 = origin();
let mut tracker = SystemCpuTracker::new();
tracker.observe(CpuTimeTotals::new(secs(0), secs(0)), t0);
let state = tracker.observe(CpuTimeTotals::new(secs(8), secs(0)), t0 + secs(1));
assert!((percent(&state) - 100.0).abs() < f32::EPSILON);
}
#[test]
fn system_cpu_does_not_depend_on_the_sample_interval_length() {
let t0 = origin();
let before = CpuTimeTotals::new(secs(10), secs(70));
let after = CpuTimeTotals::new(secs(11), secs(77));
let mut quick = SystemCpuTracker::new();
quick.observe(before, t0);
let quick_state = quick.observe(after, t0 + Duration::from_millis(250));
let mut slow = SystemCpuTracker::new();
slow.observe(before, t0);
let slow_state = slow.observe(after, t0 + secs(4));
assert!((percent(&quick_state) - percent(&slow_state)).abs() < f32::EPSILON);
}
#[test]
fn a_stalled_system_counter_is_warming_up_not_zero_percent() {
let t0 = origin();
let totals = CpuTimeTotals::new(secs(10), secs(70));
let mut tracker = SystemCpuTracker::new();
tracker.observe(totals, t0);
let state = tracker.observe(totals, t0 + secs(1));
assert!(state.is_warming_up());
assert_ne!(state, MetricState::Available(Percent::ZERO));
}
#[test]
fn system_cpu_time_going_backwards_is_a_reset_and_recovers_next_sample() {
let t0 = origin();
let mut tracker = SystemCpuTracker::new();
tracker.observe(CpuTimeTotals::new(secs(100), secs(700)), t0);
let reset = tracker.observe(CpuTimeTotals::new(secs(2), secs(6)), t0 + secs(1));
assert_eq!(
reset,
MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset)
);
let recovered = tracker.observe(CpuTimeTotals::new(secs(3), secs(13)), t0 + secs(2));
assert!((percent(&recovered) - 12.5).abs() < f32::EPSILON);
}
#[test]
fn an_idle_only_reset_is_detected_as_well_as_a_busy_only_one() {
let t0 = origin();
let mut tracker = SystemCpuTracker::new();
tracker.observe(CpuTimeTotals::new(secs(100), secs(700)), t0);
let reset = tracker.observe(CpuTimeTotals::new(secs(101), secs(1)), t0 + secs(1));
assert_eq!(
reset,
MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset)
);
}
#[test]
fn a_first_process_sample_is_warming_up_and_not_zero_percent() {
let mut tracker = ProcessCpuTracker::new();
let state = tracker.observe(secs(42), origin());
assert!(state.is_warming_up());
assert_ne!(state, MetricState::Available(Percent::ZERO));
}
#[test]
fn a_fully_busy_single_core_reads_one_hundred_percent_on_an_eight_cpu_machine() {
let t0 = origin();
let mut tracker = ProcessCpuTracker::new();
tracker.observe(secs(0), t0);
let state = tracker.observe(secs(1), t0 + secs(1));
assert!((percent(&state) - 100.0).abs() < f32::EPSILON);
let machine =
CpuNormalization::Machine.apply(Percent::new(percent(&state)).expect("valid"), 8);
assert!((machine.expect("valid").value() - 12.5).abs() < f32::EPSILON);
}
#[test]
fn a_process_on_four_cores_reads_four_hundred_percent_under_core_normalization() {
let t0 = origin();
let mut tracker = ProcessCpuTracker::new();
tracker.observe(secs(0), t0);
let state = tracker.observe(secs(4), t0 + secs(1));
assert!((percent(&state) - 400.0).abs() < f32::EPSILON);
}
#[test]
fn both_normalizations_are_reachable_from_one_reading() {
let t0 = origin();
let mut core = ProcessCpuTracker::new();
core.observe(secs(0), t0);
let core_state = core.observe_normalized(secs(4), t0 + secs(1), CpuNormalization::Core, 8);
assert!((percent(&core_state) - 400.0).abs() < f32::EPSILON);
let mut machine = ProcessCpuTracker::new();
machine.observe(secs(0), t0);
let machine_state =
machine.observe_normalized(secs(4), t0 + secs(1), CpuNormalization::Machine, 8);
assert!((percent(&machine_state) - 50.0).abs() < f32::EPSILON);
}
#[test]
fn machine_normalization_without_a_cpu_count_is_unavailable_not_zero() {
let t0 = origin();
let mut tracker = ProcessCpuTracker::new();
tracker.observe(secs(0), t0);
let state = tracker.observe_normalized(secs(1), t0 + secs(1), CpuNormalization::Machine, 0);
assert_eq!(
state,
MetricState::TemporarilyUnavailable(UnavailableReason::ReadFailed)
);
assert_eq!(state.fresh(), None);
}
#[test]
fn normalization_preserves_an_unavailable_state_rather_than_scaling_it() {
let mut tracker = ProcessCpuTracker::new();
let first = tracker.observe_normalized(secs(5), origin(), CpuNormalization::Machine, 8);
assert!(first.is_warming_up());
}
#[test]
fn process_cpu_uses_the_actual_elapsed_interval() {
let t0 = origin();
let cpu = Duration::from_millis(500);
let mut quick = ProcessCpuTracker::new();
quick.observe(Duration::ZERO, t0);
let quick_state = quick.observe(cpu, t0 + Duration::from_millis(500));
let mut slow = ProcessCpuTracker::new();
slow.observe(Duration::ZERO, t0);
let slow_state = slow.observe(cpu, t0 + secs(2));
assert!((percent(&quick_state) - 100.0).abs() < f32::EPSILON);
assert!((percent(&slow_state) - 25.0).abs() < f32::EPSILON);
}
#[test]
fn a_process_that_used_no_cpu_reads_a_real_zero_percent() {
let t0 = origin();
let mut tracker = ProcessCpuTracker::new();
tracker.observe(secs(9), t0);
let state = tracker.observe(secs(9), t0 + secs(1));
assert_eq!(state, MetricState::Available(Percent::ZERO));
}
#[test]
fn zero_elapsed_process_sample_is_warming_up_not_a_division_by_zero() {
let t0 = origin();
let mut tracker = ProcessCpuTracker::new();
tracker.observe(secs(1), t0);
let state = tracker.observe(secs(2), t0);
assert!(state.is_warming_up());
}
#[test]
fn a_reversed_instant_cannot_make_a_process_percentage_negative_or_huge() {
let t0 = origin();
let mut tracker = ProcessCpuTracker::new();
tracker.observe(secs(0), t0 + secs(10));
let state = tracker.observe(secs(5), t0);
assert!(state.is_warming_up());
}
#[test]
fn process_cpu_time_going_backwards_is_a_reset_and_recovers_next_sample() {
let t0 = origin();
let mut tracker = ProcessCpuTracker::new();
tracker.observe(secs(30), t0);
let reset = tracker.observe(secs(1), t0 + secs(1));
assert_eq!(
reset,
MetricState::TemporarilyUnavailable(UnavailableReason::CounterReset)
);
let recovered = tracker.observe(secs(2), t0 + secs(2));
assert!((percent(&recovered) - 100.0).abs() < f32::EPSILON);
}
#[test]
fn forgetting_a_process_baseline_prevents_a_delta_across_the_gap() {
let t0 = origin();
let mut tracker = ProcessCpuTracker::new();
tracker.observe(secs(0), t0);
tracker.forget_baseline();
assert!(tracker.is_warming_up());
assert!(tracker.observe(secs(600), t0 + secs(1)).is_warming_up());
}
#[test]
fn trackers_report_when_they_last_saw_a_reading() {
let t0 = origin();
let at = t0 + secs(5);
let mut system = SystemCpuTracker::new();
assert_eq!(system.last_observed_at(), None);
system.observe(CpuTimeTotals::default(), at);
assert_eq!(system.last_observed_at(), Some(at));
let mut process = ProcessCpuTracker::new();
assert_eq!(process.last_observed_at(), None);
process.observe(Duration::ZERO, at);
assert_eq!(process.last_observed_at(), Some(at));
}
}