use std::{
collections::VecDeque,
time::{Duration, Instant},
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct HostTelemetry {
rtt_history: VecDeque<(Instant, Duration)>,
pub max_samples: usize,
pub ttl: Option<u8>,
pub distance_hops: Option<u8>,
}
impl HostTelemetry {
pub fn new(max_samples: usize) -> Self {
Self {
rtt_history: VecDeque::with_capacity(max_samples),
max_samples,
ttl: None,
distance_hops: None,
}
}
pub fn history(&self) -> &VecDeque<(Instant, Duration)> {
&self.rtt_history
}
pub fn add_rtt(&mut self, rtt: Duration) {
self.add_rtt_at(Instant::now(), rtt);
}
pub fn add_rtt_at(&mut self, time: Instant, rtt: Duration) {
if self.max_samples == 0 {
return;
}
self.rtt_history.push_back((time, rtt));
while self.rtt_history.len() > self.max_samples {
self.rtt_history.pop_front();
}
}
pub fn lartt(&self) -> Option<Duration> {
self.rtt_history.back().map(|&(_, rtt)| rtt)
}
pub fn min_rtt(&self) -> Option<Duration> {
self.rtt_history.iter().map(|(_, rtt)| *rtt).min()
}
pub fn max_rtt(&self) -> Option<Duration> {
self.rtt_history.iter().map(|(_, rtt)| *rtt).max()
}
pub fn average_rtt(&self) -> Option<Duration> {
if self.rtt_history.is_empty() {
return None;
}
let sum: Duration = self.rtt_history.iter().map(|(_, rtt)| *rtt).sum();
Some(sum / self.rtt_history.len() as u32)
}
pub fn jitter(&self) -> Option<Duration> {
if self.rtt_history.len() < 2 {
return None;
}
let mut total_diff = Duration::ZERO;
let mut prev = self.rtt_history[0].1;
for &(_, curr) in self.rtt_history.iter().skip(1) {
total_diff += curr.abs_diff(prev);
prev = curr;
}
Some(total_diff / (self.rtt_history.len() - 1) as u32)
}
pub fn merge(&mut self, mut other: HostTelemetry) {
if other.max_samples > self.max_samples {
self.max_samples = other.max_samples;
}
if self.max_samples == 0 {
return;
}
let mut combined: Vec<_> = self
.rtt_history
.drain(..)
.chain(other.rtt_history.drain(..))
.collect();
combined.sort_by_key(|&(time, _)| time);
let start_idx = combined.len().saturating_sub(self.max_samples);
self.rtt_history
.extend(combined.into_iter().skip(start_idx));
if self.ttl.is_none() {
self.ttl = other.ttl;
}
if self.distance_hops.is_none() {
self.distance_hops = other.distance_hops;
}
}
}
impl std::fmt::Display for HostTelemetry {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self.average_rtt() {
Some(avg) => write!(
f,
"avg={:?}, jitter={:?}",
avg,
self.jitter().unwrap_or(Duration::ZERO)
),
None => write!(f, "no telemetry"),
}
}
}
impl Default for HostTelemetry {
fn default() -> Self {
Self::new(10)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn telemetry_math_safety() {
let t = HostTelemetry::new(10);
assert_eq!(t.average_rtt(), None);
assert_eq!(t.jitter(), None);
assert_eq!(t.lartt(), None);
}
#[test]
fn telemetry_averaging_logic() {
let mut t = HostTelemetry::new(5);
t.add_rtt(Duration::from_millis(10));
t.add_rtt(Duration::from_millis(20));
assert_eq!(t.average_rtt(), Some(Duration::from_millis(15)));
}
#[test]
fn jitter_calculation_consistency() {
let mut t = HostTelemetry::new(5);
t.add_rtt(Duration::from_millis(100)); t.add_rtt(Duration::from_millis(110)); t.add_rtt(Duration::from_millis(105)); assert_eq!(
t.jitter(),
Some(Duration::from_millis(7) + Duration::from_micros(500))
);
}
#[test]
fn merge_capacity_upgrade() {
let mut t1 = HostTelemetry::new(3);
let t2 = HostTelemetry::new(10);
t1.merge(t2);
assert_eq!(t1.max_samples, 10);
}
#[test]
fn merge_zero_capacity_safety() {
let mut t1 = HostTelemetry::new(0);
let t2 = HostTelemetry::new(0);
t1.merge(t2);
assert_eq!(t1.rtt_history.len(), 0);
}
}