use std::collections::VecDeque;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(crate) struct TimingSample {
pub(crate) wall_ns: u64,
pub(crate) active_ns: u64,
}
impl TimingSample {
pub(crate) fn new(wall: Duration, active: Duration) -> Self {
Self {
wall_ns: wall.as_nanos().min(u64::MAX as u128) as u64,
active_ns: active.as_nanos().min(u64::MAX as u128) as u64,
}
}
pub(crate) fn suspended_ns(self) -> u64 {
self.wall_ns.saturating_sub(self.active_ns)
}
}
#[derive(Debug, Default)]
pub(crate) struct RollingSamples {
values: VecDeque<TimingSample>,
total_calls: u64,
}
impl RollingSamples {
pub(crate) fn record(
&mut self,
sample: TimingSample,
window_size: usize,
dropped_samples: &AtomicU64,
) {
self.total_calls = self.total_calls.saturating_add(1);
if window_size == 0 {
dropped_samples.fetch_add(1, Ordering::Relaxed);
return;
}
if self.values.len() == window_size {
self.values.pop_front();
}
self.values.push_back(sample);
}
pub(crate) fn snapshot(&self) -> RollingSnapshot {
RollingSnapshot {
total_calls: self.total_calls,
values: self.values.iter().copied().collect(),
}
}
}
#[derive(Clone, Debug, Default)]
pub(crate) struct RollingSnapshot {
pub(crate) total_calls: u64,
pub(crate) values: Vec<TimingSample>,
}
#[cfg(test)]
mod tests {
use super::{RollingSamples, TimingSample};
use std::sync::atomic::AtomicU64;
use std::time::Duration;
#[test]
fn rolling_window_evicts_oldest_samples() {
let dropped = AtomicU64::new(0);
let mut samples = RollingSamples::default();
for value in 1..=150 {
samples.record(
TimingSample::new(Duration::from_nanos(value), Duration::from_nanos(value)),
100,
&dropped,
);
}
let snapshot = samples.snapshot();
assert_eq!(snapshot.total_calls, 150);
assert_eq!(snapshot.values.len(), 100);
assert_eq!(snapshot.values.first().unwrap().wall_ns, 51);
assert_eq!(snapshot.values.last().unwrap().wall_ns, 150);
}
}