use std::sync::Mutex;
use std::time::{Duration, Instant};
use hdrhistogram::Histogram as HdrHistogram;
#[derive(Clone, Copy, Debug)]
pub struct HdrBounds {
pub low: u64,
pub high: u64,
pub sig_digits: u8,
}
impl HdrBounds {
pub const LATENCY_DEFAULT: Self = Self {
low: 1_000, high: 60_000_000_000, sig_digits: 3,
};
fn build(&self) -> HdrHistogram<u64> {
HdrHistogram::new_with_bounds(self.low, self.high, self.sig_digits)
.expect("HDR bounds must be valid")
}
}
#[derive(Clone, Copy, Debug)]
pub struct LiveWindowConfig {
pub window: Duration,
pub slot_count: usize,
pub bounds: HdrBounds,
}
impl Default for LiveWindowConfig {
fn default() -> Self {
Self {
window: Duration::from_secs(1),
slot_count: 10,
bounds: HdrBounds::LATENCY_DEFAULT,
}
}
}
struct Slot {
started_at: Instant,
hist: HdrHistogram<u64>,
}
pub struct LiveWindowHistogram {
config: LiveWindowConfig,
slot_duration_ns: u128,
slots: Vec<Mutex<Slot>>,
scratch: Mutex<HdrHistogram<u64>>,
anchor: Instant,
}
impl LiveWindowHistogram {
pub fn new(config: LiveWindowConfig) -> Self {
assert!(config.slot_count > 0, "slot_count must be > 0");
assert!(config.window.as_nanos() > 0, "window must be > 0");
let slot_duration = config.window / config.slot_count as u32;
let anchor = Instant::now();
let stale_start = anchor.checked_sub(config.window).unwrap_or(anchor);
let slots = (0..config.slot_count)
.map(|_| {
Mutex::new(Slot {
started_at: stale_start,
hist: config.bounds.build(),
})
})
.collect();
let scratch = Mutex::new(config.bounds.build());
Self {
config,
slot_duration_ns: slot_duration.as_nanos(),
slots,
scratch,
anchor,
}
}
pub fn config(&self) -> &LiveWindowConfig {
&self.config
}
pub fn record(&self, value: u64) {
let now = Instant::now();
let idx = self.slot_index(now);
let boundary = self.slot_boundary_for(now, idx);
let mut slot = self.slots[idx].lock().unwrap_or_else(|e| e.into_inner());
if slot.started_at < boundary {
slot.hist.reset();
slot.started_at = boundary;
}
let v = value.min(self.config.bounds.high);
let _ = slot.hist.record(v.max(self.config.bounds.low));
}
pub fn peek(&self) -> HdrHistogram<u64> {
let now = Instant::now();
let mut scratch = self.scratch.lock().unwrap_or_else(|e| e.into_inner());
scratch.reset();
for slot_mu in &self.slots {
let slot = slot_mu.lock().unwrap_or_else(|e| e.into_inner());
if now.duration_since(slot.started_at) < self.config.window {
let _ = scratch.add(&slot.hist);
}
}
scratch.clone()
}
pub fn len(&self) -> u64 {
let now = Instant::now();
let mut total = 0u64;
for slot_mu in &self.slots {
let slot = slot_mu.lock().unwrap_or_else(|e| e.into_inner());
if now.duration_since(slot.started_at) < self.config.window {
total = total.saturating_add(slot.hist.len());
}
}
total
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
fn slot_index(&self, now: Instant) -> usize {
let ns = now.duration_since(self.anchor).as_nanos();
((ns / self.slot_duration_ns) as usize) % self.config.slot_count
}
fn slot_boundary_for(&self, now: Instant, idx: usize) -> Instant {
let ns = now.duration_since(self.anchor).as_nanos();
let current_slot_start_ns = (ns / self.slot_duration_ns) * self.slot_duration_ns;
let current_idx = ((ns / self.slot_duration_ns) as usize) % self.config.slot_count;
let offset = ((current_idx + self.config.slot_count) - idx) % self.config.slot_count;
let boundary_ns =
current_slot_start_ns.saturating_sub((offset as u128) * self.slot_duration_ns);
self.anchor + Duration::from_nanos(boundary_ns as u64)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn fast_config() -> LiveWindowConfig {
LiveWindowConfig {
window: Duration::from_millis(500),
slot_count: 5,
bounds: HdrBounds::LATENCY_DEFAULT,
}
}
#[test]
fn peek_returns_merged_fresh_slots() {
let ring = LiveWindowHistogram::new(fast_config());
ring.record(1_000_000);
ring.record(2_000_000);
ring.record(3_000_000);
let snap = ring.peek();
assert_eq!(snap.len(), 3);
assert!(snap.max() >= 3_000_000);
}
#[test]
fn peek_does_not_drain_slots() {
let ring = LiveWindowHistogram::new(fast_config());
ring.record(100_000);
let a = ring.peek();
let b = ring.peek();
assert_eq!(a.len(), b.len(), "peek must be idempotent");
}
#[test]
fn len_matches_peek_len() {
let ring = LiveWindowHistogram::new(fast_config());
for _ in 0..50 {
ring.record(100_000);
}
assert_eq!(ring.len(), ring.peek().len());
}
#[test]
fn stale_slots_do_not_contribute() {
let cfg = LiveWindowConfig {
window: Duration::from_millis(200),
slot_count: 4,
bounds: HdrBounds::LATENCY_DEFAULT,
};
let ring = LiveWindowHistogram::new(cfg);
ring.record(1_000_000);
assert_eq!(ring.len(), 1);
std::thread::sleep(Duration::from_millis(260));
let late = ring.peek();
assert_eq!(
late.len(),
0,
"every slot older than window should be skipped"
);
}
#[test]
fn lazy_reset_only_on_next_record_into_slot() {
let cfg = LiveWindowConfig {
window: Duration::from_millis(150),
slot_count: 3,
bounds: HdrBounds::LATENCY_DEFAULT,
};
let ring = LiveWindowHistogram::new(cfg);
ring.record(500_000);
assert_eq!(ring.len(), 1);
std::thread::sleep(Duration::from_millis(200));
ring.record(700_000);
let snap = ring.peek();
assert_eq!(snap.len(), 1, "stale slots must not leak old samples");
assert!(snap.max() >= 700_000);
}
#[test]
fn record_bounds_clamp_out_of_range_values() {
let ring = LiveWindowHistogram::new(fast_config());
ring.record(0);
ring.record(u64::MAX);
let snap = ring.peek();
assert_eq!(snap.len(), 2);
assert!(snap.max() <= 60_000_000_000 * 2);
assert!(
snap.min() > 0,
"below-low clamp must produce a positive recorded value"
);
}
}