use std::sync::atomic::{AtomicU64, Ordering::Relaxed};
pub(crate) struct PipelineMetrics {
pub raw_send_wait_ns: AtomicU64,
pub decode_admit_wait_ns: AtomicU64,
pub decode_admit_blocked: AtomicU64,
pub decoded_send_wait_ns: AtomicU64,
pub decoded_recv_wait_ns: AtomicU64,
pub reorder_high_water: AtomicU64,
pub reorder_window_high_water: AtomicU64,
pub scratch_st_capacity_peak_bytes: AtomicU64,
pub scratch_gr_capacity_peak_bytes: AtomicU64,
pub decode_tasks: AtomicU64,
pub blobs_skipped_by_filter: AtomicU64,
}
impl PipelineMetrics {
const fn new() -> Self {
Self {
raw_send_wait_ns: AtomicU64::new(0),
decode_admit_wait_ns: AtomicU64::new(0),
decode_admit_blocked: AtomicU64::new(0),
decoded_send_wait_ns: AtomicU64::new(0),
decoded_recv_wait_ns: AtomicU64::new(0),
reorder_high_water: AtomicU64::new(0),
reorder_window_high_water: AtomicU64::new(0),
scratch_st_capacity_peak_bytes: AtomicU64::new(0),
scratch_gr_capacity_peak_bytes: AtomicU64::new(0),
decode_tasks: AtomicU64::new(0),
blobs_skipped_by_filter: AtomicU64::new(0),
}
}
pub fn record_reorder_levels(&self, filled: usize, window: usize) {
cas_max(&self.reorder_high_water, filled as u64);
cas_max(&self.reorder_window_high_water, window as u64);
}
pub fn record_scratch_capacity(&self, st_bytes: usize, gr_bytes: usize) {
cas_max(&self.scratch_st_capacity_peak_bytes, st_bytes as u64);
cas_max(&self.scratch_gr_capacity_peak_bytes, gr_bytes as u64);
}
pub fn emit(&self) {
macro_rules! emit {
($name:literal, $field:ident) => {
crate::debug::emit_counter(
$name,
i64::try_from(self.$field.load(Relaxed)).unwrap_or(i64::MAX),
);
};
}
emit!("pipeline_raw_send_wait_ns", raw_send_wait_ns);
emit!("pipeline_decode_admit_wait_ns", decode_admit_wait_ns);
emit!("pipeline_decode_admit_blocked", decode_admit_blocked);
emit!("pipeline_decoded_send_wait_ns", decoded_send_wait_ns);
emit!("pipeline_decoded_recv_wait_ns", decoded_recv_wait_ns);
emit!("pipeline_reorder_high_water", reorder_high_water);
emit!(
"pipeline_reorder_window_high_water",
reorder_window_high_water
);
emit!(
"pipeline_scratch_st_capacity_peak_bytes",
scratch_st_capacity_peak_bytes
);
emit!(
"pipeline_scratch_gr_capacity_peak_bytes",
scratch_gr_capacity_peak_bytes
);
emit!("pipeline_decode_tasks", decode_tasks);
emit!("pipeline_blobs_skipped_by_filter", blobs_skipped_by_filter);
}
}
fn cas_max(field: &AtomicU64, candidate: u64) {
let mut current = field.load(Relaxed);
while candidate > current {
match field.compare_exchange_weak(current, candidate, Relaxed, Relaxed) {
Ok(_) => break,
Err(observed) => current = observed,
}
}
}
pub(crate) static PIPELINE_METRICS: PipelineMetrics = PipelineMetrics::new();
#[inline]
pub(crate) fn elapsed_ns_u64(start: std::time::Instant) -> u64 {
u64::try_from(start.elapsed().as_nanos()).unwrap_or(u64::MAX)
}