use std::time::Duration;
pub const SAMPLE_WINDOW: usize = 256;
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct Distribution {
pub p50: u32,
pub p95: u32,
pub max: u32,
pub count: usize,
}
#[derive(Clone, Debug)]
pub struct RingStats<const N: usize> {
values: [u32; N],
len: usize,
next: usize,
}
impl<const N: usize> Default for RingStats<N> {
fn default() -> Self {
Self { values: [0; N], len: 0, next: 0 }
}
}
impl<const N: usize> RingStats<N> {
pub fn push(&mut self, value: u32) {
self.values[self.next] = value;
self.next = (self.next + 1) % N;
self.len = (self.len + 1).min(N);
}
pub fn stats(&self) -> Distribution {
if self.len == 0 {
return Distribution::default();
}
let mut sorted = self.values;
sorted[..self.len].sort_unstable();
let rank = |percentile: usize| sorted[(self.len * percentile / 100).min(self.len - 1)];
Distribution { p50: rank(50), p95: rank(95), max: sorted[self.len - 1], count: self.len }
}
}
pub fn duration_micros(elapsed: Duration) -> u32 {
elapsed.as_micros().min(u128::from(u32::MAX)) as u32
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct TickProfileSnapshot {
pub drain_observations_us: Distribution,
pub drain_ingress_us: Distribution,
pub step_control_plane_us: Distribution,
pub dispatch_effects_us: Distribution,
pub publish_step_us: Distribution,
pub provider_supervisors_us: Distribution,
}
#[derive(Clone, Debug, Default)]
pub struct TickPhaseProfiler {
drain_observations_us: RingStats<SAMPLE_WINDOW>,
drain_ingress_us: RingStats<SAMPLE_WINDOW>,
step_control_plane_us: RingStats<SAMPLE_WINDOW>,
dispatch_effects_us: RingStats<SAMPLE_WINDOW>,
publish_step_us: RingStats<SAMPLE_WINDOW>,
provider_supervisors_us: RingStats<SAMPLE_WINDOW>,
}
impl TickPhaseProfiler {
pub fn record_drain_observations(&mut self, elapsed: Duration) {
self.drain_observations_us.push(duration_micros(elapsed));
}
pub fn record_drain_ingress(&mut self, elapsed: Duration) {
self.drain_ingress_us.push(duration_micros(elapsed));
}
pub fn record_step_control_plane(&mut self, elapsed: Duration) {
self.step_control_plane_us.push(duration_micros(elapsed));
}
pub fn record_dispatch_effects(&mut self, elapsed: Duration) {
self.dispatch_effects_us.push(duration_micros(elapsed));
}
pub fn record_publish_step(&mut self, elapsed: Duration) {
self.publish_step_us.push(duration_micros(elapsed));
}
pub fn record_provider_supervisors(&mut self, elapsed: Duration) {
self.provider_supervisors_us.push(duration_micros(elapsed));
}
pub fn snapshot(&self) -> TickProfileSnapshot {
TickProfileSnapshot {
drain_observations_us: self.drain_observations_us.stats(),
drain_ingress_us: self.drain_ingress_us.stats(),
step_control_plane_us: self.step_control_plane_us.stats(),
dispatch_effects_us: self.dispatch_effects_us.stats(),
publish_step_us: self.publish_step_us.stats(),
provider_supervisors_us: self.provider_supervisors_us.stats(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ring_stats_reports_nearest_rank_percentiles() {
let mut ring: RingStats<8> = RingStats::default();
for value in [10, 20, 30, 40, 50, 60, 70, 80] {
ring.push(value);
}
let stats = ring.stats();
assert_eq!(stats.count, 8);
assert_eq!(stats.max, 80);
assert_eq!(stats.p50, 50);
}
#[test]
fn ring_stats_evicts_oldest_once_full() {
let mut ring: RingStats<4> = RingStats::default();
for value in 1..=6u32 {
ring.push(value);
}
let stats = ring.stats();
assert_eq!(stats.count, 4);
assert_eq!(stats.max, 6);
assert_eq!(stats.p50, 5);
}
#[test]
fn empty_ring_reports_zeroed_distribution() {
let ring: RingStats<SAMPLE_WINDOW> = RingStats::default();
assert_eq!(ring.stats(), Distribution::default());
}
#[test]
fn tick_phase_profiler_snapshot_reflects_every_recorded_phase() {
let mut profiler = TickPhaseProfiler::default();
profiler.record_drain_observations(Duration::from_micros(10));
profiler.record_drain_ingress(Duration::from_micros(20));
profiler.record_step_control_plane(Duration::from_micros(1_500));
profiler.record_dispatch_effects(Duration::from_micros(30));
profiler.record_publish_step(Duration::from_micros(40));
profiler.record_provider_supervisors(Duration::from_micros(50));
let snapshot = profiler.snapshot();
assert_eq!(snapshot.drain_observations_us.max, 10);
assert_eq!(snapshot.drain_ingress_us.max, 20);
assert_eq!(snapshot.step_control_plane_us.max, 1_500);
assert_eq!(snapshot.dispatch_effects_us.max, 30);
assert_eq!(snapshot.publish_step_us.max, 40);
assert_eq!(snapshot.provider_supervisors_us.max, 50);
for distribution in [
snapshot.drain_observations_us,
snapshot.drain_ingress_us,
snapshot.step_control_plane_us,
snapshot.dispatch_effects_us,
snapshot.publish_step_us,
snapshot.provider_supervisors_us,
] {
assert_eq!(distribution.count, 1);
}
}
}