use std::time::Duration;
use crate::ctx::{Ctx, SourceIdx};
use crate::error::Result;
use crate::stats::{CaptureStats, DropBreakdown};
#[derive(Debug, Clone, Copy, PartialEq)]
#[non_exhaustive]
pub struct CaptureTelemetry {
pub source: SourceIdx,
pub packets: u64,
pub drops: u64,
pub freezes: u64,
pub drop_rate: f64,
pub detail: DropBreakdown,
}
impl CaptureTelemetry {
#[inline]
pub fn lifetime_drop_rate(&self) -> f64 {
let total = self.packets + self.drops;
if total == 0 {
0.0
} else {
self.drops as f64 / total as f64
}
}
#[inline]
pub fn is_degraded(&self, threshold: f64) -> bool {
self.drop_rate >= threshold
}
#[cfg(feature = "metrics")]
pub fn record_metrics(&self) {
let source = self.source.0.to_string();
metrics::gauge!(crate::metrics::GAUGE_PACKETS, "source" => source.clone())
.set(self.packets as f64);
metrics::gauge!(crate::metrics::GAUGE_DROPS, "source" => source.clone())
.set(self.drops as f64);
metrics::gauge!(crate::metrics::GAUGE_FREEZES, "source" => source.clone())
.set(self.freezes as f64);
metrics::gauge!(crate::metrics::GAUGE_DROP_RATE, "source" => source).set(self.drop_rate);
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize))]
#[non_exhaustive]
pub struct CaptureHealth {
pub source: u8,
pub packets: u64,
pub drops: u64,
pub freezes: u64,
pub drop_rate: f64,
pub lifetime_drop_rate: f64,
pub detail: DropBreakdown,
}
impl crate::report::Report for CaptureHealth {
const NAME: &'static str = "capture_health";
}
impl From<CaptureTelemetry> for CaptureHealth {
fn from(t: CaptureTelemetry) -> Self {
Self {
source: t.source.0,
packets: t.packets,
drops: t.drops,
freezes: t.freezes,
drop_rate: t.drop_rate,
lifetime_drop_rate: t.lifetime_drop_rate(),
detail: t.detail,
}
}
}
pub(crate) struct TelemetrySampler {
last: Vec<(u64, u64)>,
}
impl TelemetrySampler {
pub(crate) fn new(num_sources: usize) -> Self {
Self {
last: vec![(0, 0); num_sources],
}
}
pub(crate) fn sample(
&mut self,
source: usize,
cum: CaptureStats,
detail: DropBreakdown,
) -> CaptureTelemetry {
let packets = cum.packets as u64;
let drops = cum.drops as u64;
let (last_packets, last_drops) = self.last[source];
let window_packets = packets.saturating_sub(last_packets);
let window_drops = drops.saturating_sub(last_drops);
self.last[source] = (packets, drops);
let window_total = window_packets + window_drops;
let drop_rate = if window_total == 0 {
0.0
} else {
window_drops as f64 / window_total as f64
};
CaptureTelemetry {
source: SourceIdx(source as u8),
packets,
drops,
freezes: cum.freeze_count as u64,
drop_rate,
detail,
}
}
}
pub(crate) type BoxedCaptureStatsHandler =
Box<dyn FnMut(&CaptureTelemetry, &mut Ctx<'_>) -> Result<()> + Send>;
pub(crate) struct CaptureStatsRegistration {
pub(crate) period: Duration,
pub(crate) handler: BoxedCaptureStatsHandler,
}
impl std::fmt::Debug for CaptureStatsRegistration {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("CaptureStatsRegistration")
.field("period", &self.period)
.finish_non_exhaustive()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn stats(packets: u32, drops: u32, freezes: u32) -> CaptureStats {
CaptureStats {
packets,
drops,
freeze_count: freezes,
}
}
fn afp(freezes: u32) -> DropBreakdown {
DropBreakdown::AfPacket {
freezes: freezes as u64,
}
}
#[test]
fn windowed_drop_rate_uses_the_delta_not_the_lifetime_total() {
let mut s = TelemetrySampler::new(1);
let t0 = s.sample(0, stats(1000, 0, 0), afp(0));
assert_eq!(t0.packets, 1000);
assert_eq!(t0.drops, 0);
assert_eq!(t0.drop_rate, 0.0);
assert_eq!(t0.lifetime_drop_rate(), 0.0);
let t1 = s.sample(0, stats(1100, 900, 3), afp(3));
assert_eq!(t1.packets, 1100);
assert_eq!(t1.drops, 900);
assert_eq!(t1.freezes, 3);
assert!(
(t1.drop_rate - 0.9).abs() < 1e-9,
"window rate = {}",
t1.drop_rate
);
assert!(
(t1.lifetime_drop_rate() - 0.45).abs() < 1e-9,
"lifetime rate = {}",
t1.lifetime_drop_rate()
);
assert!(t1.is_degraded(0.5));
assert!(!t1.is_degraded(0.95));
}
#[test]
fn idle_window_reports_zero_drop_rate_not_nan() {
let mut s = TelemetrySampler::new(1);
let _ = s.sample(0, stats(500, 10, 0), afp(0));
let t = s.sample(0, stats(500, 10, 0), afp(0));
assert_eq!(t.drop_rate, 0.0);
assert!(!t.drop_rate.is_nan());
}
#[test]
fn counter_going_backwards_saturates_to_zero_window() {
let mut s = TelemetrySampler::new(1);
let _ = s.sample(0, stats(1000, 50, 0), afp(0));
let t = s.sample(0, stats(10, 1, 0), afp(0));
assert_eq!(t.drop_rate, 0.0);
assert_eq!(t.packets, 10);
}
#[test]
fn capture_health_flattens_telemetry_including_lifetime_rate() {
let mut s = TelemetrySampler::new(1);
let _ = s.sample(0, stats(1000, 0, 0), afp(0));
let t = s.sample(0, stats(1100, 900, 3), afp(3));
let h = CaptureHealth::from(t);
assert_eq!(h.source, 0);
assert_eq!(h.packets, 1100);
assert_eq!(h.drops, 900);
assert_eq!(h.freezes, 3);
assert!((h.drop_rate - 0.9).abs() < 1e-9);
assert!((h.lifetime_drop_rate - 0.45).abs() < 1e-9);
}
#[cfg(feature = "metrics")]
#[test]
fn record_metrics_is_a_noop_without_a_recorder() {
let mut s = TelemetrySampler::new(1);
let t = s.sample(0, stats(1000, 10, 0), afp(0));
t.record_metrics();
}
#[cfg(feature = "serde")]
#[test]
fn capture_health_serializes_to_a_json_line() {
let h = CaptureHealth {
source: 1,
packets: 42,
drops: 7,
freezes: 0,
drop_rate: 0.25,
lifetime_drop_rate: 0.14,
detail: DropBreakdown::AfPacket { freezes: 0 },
};
let line = serde_json::to_string(&h).expect("serialize");
assert!(line.contains("\"source\":1"));
assert!(line.contains("\"packets\":42"));
assert!(line.contains("\"drop_rate\":0.25"));
assert!(line.contains("\"detail\""));
assert!(line.contains("\"AfPacket\""));
}
#[test]
fn xdp_detail_is_carried_through_the_sample_uncollapsed() {
let mut s = TelemetrySampler::new(1);
let detail = DropBreakdown::Xdp {
rx_dropped: 1,
rx_invalid_descs: 2,
rx_ring_full: 3,
rx_fill_ring_empty_descs: 4,
tx_invalid_descs: 5,
tx_ring_empty_descs: 6,
};
let t = s.sample(0, stats(0, 8, 0), detail);
assert_eq!(t.detail, detail);
match t.detail {
DropBreakdown::Xdp {
rx_invalid_descs, ..
} => assert_eq!(rx_invalid_descs, 2),
_ => panic!("expected an Xdp breakdown"),
}
}
#[test]
fn sources_are_tracked_independently() {
let mut s = TelemetrySampler::new(2);
let _ = s.sample(0, stats(100, 0, 0), afp(0));
let _ = s.sample(1, stats(0, 0, 0), afp(0));
let a = s.sample(0, stats(200, 0, 0), afp(0));
let b = s.sample(1, stats(100, 100, 0), afp(0));
assert_eq!(a.source, SourceIdx(0));
assert_eq!(a.drop_rate, 0.0);
assert_eq!(b.source, SourceIdx(1));
assert!((b.drop_rate - 0.5).abs() < 1e-9);
}
}