Skip to main content

multiprobe/
analytics.rs

1//! Advanced Path Analytics
2//!
3//! This module provides sophisticated network path analysis including:
4//! - Path MTU Discovery (PMTUD)
5//! - Jitter and packet reordering metrics
6//! - Bufferbloat detection under load
7//! - Latency distribution analysis
8
9use std::collections::VecDeque;
10use std::net::{IpAddr, Ipv4Addr, SocketAddr};
11use std::time::{Duration, Instant};
12use std::mem::MaybeUninit;
13
14use socket2::{Domain, Protocol, Socket, Type};
15
16use crate::dns;
17use crate::error::Error;
18
19// ============================================================================
20// JITTER AND LATENCY ANALYSIS
21// ============================================================================
22
23/// Statistics from multiple probe samples
24#[derive(Debug, Clone)]
25pub struct LatencyStats {
26    /// Minimum RTT observed
27    pub min_rtt: Duration,
28    /// Maximum RTT observed
29    pub max_rtt: Duration,
30    /// Mean RTT
31    pub mean_rtt: Duration,
32    /// Median RTT
33    pub median_rtt: Duration,
34    /// Standard deviation
35    pub std_dev: Duration,
36    /// Jitter (mean absolute deviation)
37    pub jitter: Duration,
38    /// 95th percentile RTT
39    pub p95_rtt: Duration,
40    /// 99th percentile RTT
41    pub p99_rtt: Duration,
42    /// Packet loss rate (0.0 - 1.0)
43    pub loss_rate: f64,
44    /// Number of samples
45    pub sample_count: usize,
46    /// Number of successful probes
47    pub success_count: usize,
48}
49
50impl LatencyStats {
51    /// Calculate statistics from a list of RTT measurements
52    pub fn from_samples(samples: &[Option<Duration>]) -> Self {
53        let successful: Vec<Duration> = samples.iter()
54            .filter_map(|s| *s)
55            .collect();
56
57        let sample_count = samples.len();
58        let success_count = successful.len();
59        let loss_rate = if sample_count > 0 {
60            1.0 - (success_count as f64 / sample_count as f64)
61        } else {
62            1.0
63        };
64
65        if successful.is_empty() {
66            return Self {
67                min_rtt: Duration::ZERO,
68                max_rtt: Duration::ZERO,
69                mean_rtt: Duration::ZERO,
70                median_rtt: Duration::ZERO,
71                std_dev: Duration::ZERO,
72                jitter: Duration::ZERO,
73                p95_rtt: Duration::ZERO,
74                p99_rtt: Duration::ZERO,
75                loss_rate,
76                sample_count,
77                success_count,
78            };
79        }
80
81        let mut sorted: Vec<u128> = successful.iter()
82            .map(|d| d.as_micros())
83            .collect();
84        sorted.sort_unstable();
85
86        let min = sorted[0];
87        let max = sorted[sorted.len() - 1];
88        let sum: u128 = sorted.iter().sum();
89        let mean = sum / sorted.len() as u128;
90
91        let median = if sorted.len() % 2 == 0 {
92            (sorted[sorted.len() / 2 - 1] + sorted[sorted.len() / 2]) / 2
93        } else {
94            sorted[sorted.len() / 2]
95        };
96
97        // Standard deviation
98        let variance: f64 = sorted.iter()
99            .map(|&x| {
100                let diff = x as f64 - mean as f64;
101                diff * diff
102            })
103            .sum::<f64>() / sorted.len() as f64;
104        let std_dev = variance.sqrt();
105
106        // Jitter (RFC 3550 style - mean of absolute differences)
107        let jitter = if sorted.len() > 1 {
108            let diffs: Vec<u128> = sorted.windows(2)
109                .map(|w| (w[1] as i128 - w[0] as i128).unsigned_abs())
110                .collect();
111            diffs.iter().sum::<u128>() / diffs.len() as u128
112        } else {
113            0
114        };
115
116        // Percentiles
117        let p95_idx = (sorted.len() as f64 * 0.95) as usize;
118        let p99_idx = (sorted.len() as f64 * 0.99) as usize;
119        let p95 = sorted[std::cmp::min(p95_idx, sorted.len() - 1)];
120        let p99 = sorted[std::cmp::min(p99_idx, sorted.len() - 1)];
121
122        Self {
123            min_rtt: Duration::from_micros(min as u64),
124            max_rtt: Duration::from_micros(max as u64),
125            mean_rtt: Duration::from_micros(mean as u64),
126            median_rtt: Duration::from_micros(median as u64),
127            std_dev: Duration::from_micros(std_dev as u64),
128            jitter: Duration::from_micros(jitter as u64),
129            p95_rtt: Duration::from_micros(p95 as u64),
130            p99_rtt: Duration::from_micros(p99 as u64),
131            loss_rate,
132            sample_count,
133            success_count,
134        }
135    }
136
137    /// Check if jitter is considered high (>10% of mean RTT)
138    pub fn has_high_jitter(&self) -> bool {
139        if self.mean_rtt.as_micros() == 0 {
140            return false;
141        }
142        let ratio = self.jitter.as_micros() as f64 / self.mean_rtt.as_micros() as f64;
143        ratio > 0.10
144    }
145
146    /// Check if there's significant packet loss (>1%)
147    pub fn has_packet_loss(&self) -> bool {
148        self.loss_rate > 0.01
149    }
150}
151
152/// Collect latency samples to a target
153pub async fn measure_latency(
154    target: &str,
155    port: u16,
156    sample_count: usize,
157    interval: Duration,
158) -> crate::Result<LatencyStats> {
159    let dns_result = dns::resolve_ipv4(target).await?;
160    let target_ip = match dns_result.ip {
161        IpAddr::V4(ipv4) => ipv4,
162        IpAddr::V6(_) => return Err(Error::InvalidTarget("IPv6 not supported".to_string())),
163    };
164
165    let mut samples = Vec::with_capacity(sample_count);
166
167    for _ in 0..sample_count {
168        let rtt = tcp_ping(target_ip, port, Duration::from_secs(5)).await;
169        samples.push(rtt);
170
171        if interval > Duration::ZERO {
172            tokio::time::sleep(interval).await;
173        }
174    }
175
176    Ok(LatencyStats::from_samples(&samples))
177}
178
179async fn tcp_ping(target: Ipv4Addr, port: u16, timeout: Duration) -> Option<Duration> {
180    let start = Instant::now();
181
182    let result = tokio::time::timeout(timeout, async {
183        tokio::net::TcpStream::connect(SocketAddr::new(IpAddr::V4(target), port)).await
184    }).await;
185
186    match result {
187        Ok(Ok(_)) => Some(start.elapsed()),
188        Ok(Err(_)) => Some(start.elapsed()), // Connection refused still means we reached it
189        Err(_) => None, // Timeout
190    }
191}
192
193// ============================================================================
194// PATH MTU DISCOVERY
195// ============================================================================
196
197/// Result of Path MTU Discovery
198#[derive(Debug, Clone)]
199pub struct PmtudResult {
200    /// Target that was tested
201    pub target: String,
202    /// Discovered Path MTU
203    pub path_mtu: u16,
204    /// Whether PMTUD completed successfully
205    pub success: bool,
206    /// Minimum MTU tested that worked
207    pub min_working: u16,
208    /// Maximum MTU tested that failed
209    pub max_failing: Option<u16>,
210    /// Whether DF (Don't Fragment) is being honored
211    pub df_honored: bool,
212    /// ICMP fragmentation needed messages received
213    pub frag_needed_count: u32,
214}
215
216/// Perform Path MTU Discovery
217///
218/// Uses binary search with ICMP to find the largest packet size
219/// that can traverse the path without fragmentation.
220pub async fn discover_path_mtu(
221    target: &str,
222    options: &PmtudOptions,
223) -> crate::Result<PmtudResult> {
224    let dns_result = dns::resolve_ipv4(target).await?;
225    let target_ip = match dns_result.ip {
226        IpAddr::V4(ipv4) => ipv4,
227        IpAddr::V6(_) => return Err(Error::InvalidTarget("IPv6 not supported".to_string())),
228    };
229
230    // Binary search for MTU
231    let mut low = options.min_mtu;
232    let mut high = options.max_mtu;
233    let mut max_working = low;
234    let mut min_failing: Option<u16> = None;
235    let mut frag_needed_count = 0u32;
236
237    while low <= high {
238        let mid = (low + high) / 2;
239
240        let result = probe_mtu(target_ip, mid, options.timeout).await;
241
242        match result {
243            MtuProbeResult::Success => {
244                max_working = mid;
245                low = mid + 1;
246            }
247            MtuProbeResult::FragmentationNeeded => {
248                min_failing = Some(mid);
249                high = mid - 1;
250                frag_needed_count += 1;
251            }
252            MtuProbeResult::Timeout | MtuProbeResult::Error => {
253                // Could be MTU issue or network problem
254                high = mid - 1;
255            }
256        }
257    }
258
259    Ok(PmtudResult {
260        target: target.to_string(),
261        path_mtu: max_working,
262        success: true,
263        min_working: max_working,
264        max_failing: min_failing,
265        df_honored: frag_needed_count > 0,
266        frag_needed_count,
267    })
268}
269
270/// PMTUD configuration
271#[derive(Debug, Clone)]
272pub struct PmtudOptions {
273    /// Minimum MTU to test
274    pub min_mtu: u16,
275    /// Maximum MTU to test
276    pub max_mtu: u16,
277    /// Timeout per probe
278    pub timeout: Duration,
279}
280
281impl Default for PmtudOptions {
282    fn default() -> Self {
283        Self {
284            min_mtu: 68,    // Minimum IPv4 MTU
285            max_mtu: 1500,  // Standard Ethernet MTU
286            timeout: Duration::from_secs(2),
287        }
288    }
289}
290
291#[derive(Debug)]
292enum MtuProbeResult {
293    Success,
294    FragmentationNeeded,
295    Timeout,
296    Error,
297}
298
299async fn probe_mtu(target: Ipv4Addr, mtu: u16, timeout: Duration) -> MtuProbeResult {
300    let result = tokio::task::spawn_blocking(move || {
301        probe_mtu_sync(target, mtu, timeout)
302    }).await;
303
304    match result {
305        Ok(r) => r,
306        Err(_) => MtuProbeResult::Error,
307    }
308}
309
310fn probe_mtu_sync(target: Ipv4Addr, mtu: u16, timeout: Duration) -> MtuProbeResult {
311    // Create raw ICMP socket
312    let socket = match Socket::new(Domain::IPV4, Type::RAW, Some(Protocol::ICMPV4)) {
313        Ok(s) => s,
314        Err(_) => return MtuProbeResult::Error,
315    };
316
317    // Set Don't Fragment flag
318    #[cfg(target_os = "linux")]
319    {
320        use std::os::unix::io::AsRawFd;
321        let fd = socket.as_raw_fd();
322        let val: libc::c_int = libc::IP_PMTUDISC_DO;
323        unsafe {
324            libc::setsockopt(
325                fd,
326                libc::IPPROTO_IP,
327                libc::IP_MTU_DISCOVER,
328                &val as *const _ as *const libc::c_void,
329                std::mem::size_of::<libc::c_int>() as libc::socklen_t,
330            );
331        }
332    }
333
334    if socket.set_read_timeout(Some(timeout)).is_err() {
335        return MtuProbeResult::Error;
336    }
337
338    // Calculate payload size (MTU - IP header - ICMP header)
339    let payload_size = mtu.saturating_sub(20 + 8) as usize;
340    let packet = build_pmtud_packet(payload_size);
341
342    let dest = SocketAddr::new(IpAddr::V4(target), 0);
343
344    if socket.send_to(&packet, &dest.into()).is_err() {
345        return MtuProbeResult::Error;
346    }
347
348    let mut recv_buf: [MaybeUninit<u8>; 2048] = unsafe { MaybeUninit::uninit().assume_init() };
349
350    match socket.recv_from(&mut recv_buf) {
351        Ok((len, _)) => {
352            let buf: &[u8] = unsafe {
353                std::slice::from_raw_parts(recv_buf.as_ptr() as *const u8, len)
354            };
355
356            if len >= 28 {
357                let ip_header_len = ((buf[0] & 0x0F) * 4) as usize;
358                if len > ip_header_len {
359                    let icmp_type = buf[ip_header_len];
360                    let icmp_code = buf[ip_header_len + 1];
361
362                    // Type 3, Code 4 = Fragmentation Needed
363                    if icmp_type == 3 && icmp_code == 4 {
364                        return MtuProbeResult::FragmentationNeeded;
365                    }
366
367                    // Type 0 = Echo Reply (success)
368                    if icmp_type == 0 {
369                        return MtuProbeResult::Success;
370                    }
371                }
372            }
373            MtuProbeResult::Error
374        }
375        Err(_) => MtuProbeResult::Timeout,
376    }
377}
378
379fn build_pmtud_packet(payload_size: usize) -> Vec<u8> {
380    let mut packet = vec![0u8; 8 + payload_size];
381
382    // ICMP Echo Request
383    packet[0] = 8;  // Type
384    packet[1] = 0;  // Code
385    packet[2] = 0;  // Checksum (computed below)
386    packet[3] = 0;
387    packet[4] = 0;  // Identifier
388    packet[5] = 1;
389    packet[6] = 0;  // Sequence
390    packet[7] = 1;
391
392    // Fill payload with pattern
393    for i in 8..packet.len() {
394        packet[i] = (i % 256) as u8;
395    }
396
397    // Compute checksum
398    let checksum = compute_checksum(&packet);
399    packet[2] = (checksum >> 8) as u8;
400    packet[3] = checksum as u8;
401
402    packet
403}
404
405fn compute_checksum(data: &[u8]) -> u16 {
406    let mut sum: u32 = 0;
407    let mut i = 0;
408
409    while i < data.len() {
410        let word = if i + 1 < data.len() {
411            ((data[i] as u32) << 8) | (data[i + 1] as u32)
412        } else {
413            (data[i] as u32) << 8
414        };
415        sum = sum.wrapping_add(word);
416        i += 2;
417    }
418
419    while sum >> 16 != 0 {
420        sum = (sum & 0xFFFF) + (sum >> 16);
421    }
422
423    !sum as u16
424}
425
426// ============================================================================
427// BUFFERBLOAT DETECTION
428// ============================================================================
429
430/// Result of bufferbloat measurement
431#[derive(Debug, Clone)]
432pub struct BufferbloatResult {
433    /// Target tested
434    pub target: String,
435    /// Baseline latency (no load)
436    pub baseline_latency: Duration,
437    /// Latency under load
438    pub loaded_latency: Duration,
439    /// Latency increase factor
440    pub bloat_factor: f64,
441    /// Bufferbloat grade (A-F)
442    pub grade: BufferbloatGrade,
443    /// Whether bufferbloat was detected
444    pub detected: bool,
445}
446
447/// Bufferbloat severity grade
448#[derive(Debug, Clone, Copy, PartialEq, Eq)]
449pub enum BufferbloatGrade {
450    /// Excellent (< 5ms increase)
451    A,
452    /// Good (5-30ms increase)
453    B,
454    /// Fair (30-60ms increase)
455    C,
456    /// Poor (60-200ms increase)
457    D,
458    /// Bad (> 200ms increase)
459    F,
460}
461
462impl std::fmt::Display for BufferbloatGrade {
463    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
464        match self {
465            BufferbloatGrade::A => write!(f, "A (Excellent)"),
466            BufferbloatGrade::B => write!(f, "B (Good)"),
467            BufferbloatGrade::C => write!(f, "C (Fair)"),
468            BufferbloatGrade::D => write!(f, "D (Poor)"),
469            BufferbloatGrade::F => write!(f, "F (Bad)"),
470        }
471    }
472}
473
474impl BufferbloatGrade {
475    fn from_increase(increase_ms: f64) -> Self {
476        if increase_ms < 5.0 {
477            BufferbloatGrade::A
478        } else if increase_ms < 30.0 {
479            BufferbloatGrade::B
480        } else if increase_ms < 60.0 {
481            BufferbloatGrade::C
482        } else if increase_ms < 200.0 {
483            BufferbloatGrade::D
484        } else {
485            BufferbloatGrade::F
486        }
487    }
488}
489
490/// Options for bufferbloat detection
491#[derive(Debug, Clone)]
492pub struct BufferbloatOptions {
493    /// Number of baseline samples
494    pub baseline_samples: usize,
495    /// Number of loaded samples
496    pub loaded_samples: usize,
497    /// Concurrent connections to create load
498    pub load_connections: usize,
499    /// Port to use for testing
500    pub port: u16,
501}
502
503impl Default for BufferbloatOptions {
504    fn default() -> Self {
505        Self {
506            baseline_samples: 10,
507            loaded_samples: 10,
508            load_connections: 4,
509            port: 443,
510        }
511    }
512}
513
514/// Detect bufferbloat by comparing latency under load vs idle
515///
516/// Note: This is a simplified version. Full bufferbloat testing
517/// requires saturating the connection which may not be appropriate
518/// for all scenarios.
519pub async fn detect_bufferbloat(
520    target: &str,
521    options: &BufferbloatOptions,
522) -> crate::Result<BufferbloatResult> {
523    // Measure baseline latency
524    let baseline = measure_latency(
525        target,
526        options.port,
527        options.baseline_samples,
528        Duration::from_millis(100),
529    ).await?;
530
531    // Create load and measure
532    let loaded = measure_latency_under_load(
533        target,
534        options.port,
535        options.loaded_samples,
536        options.load_connections,
537    ).await?;
538
539    let baseline_ms = baseline.median_rtt.as_secs_f64() * 1000.0;
540    let loaded_ms = loaded.median_rtt.as_secs_f64() * 1000.0;
541    let increase_ms = loaded_ms - baseline_ms;
542    let bloat_factor = if baseline_ms > 0.0 { loaded_ms / baseline_ms } else { 1.0 };
543
544    let grade = BufferbloatGrade::from_increase(increase_ms);
545    let detected = increase_ms > 30.0; // >30ms increase indicates bufferbloat
546
547    Ok(BufferbloatResult {
548        target: target.to_string(),
549        baseline_latency: baseline.median_rtt,
550        loaded_latency: loaded.median_rtt,
551        bloat_factor,
552        grade,
553        detected,
554    })
555}
556
557async fn measure_latency_under_load(
558    target: &str,
559    port: u16,
560    samples: usize,
561    concurrent: usize,
562) -> crate::Result<LatencyStats> {
563    let dns_result = dns::resolve_ipv4(target).await?;
564    let target_ip = match dns_result.ip {
565        IpAddr::V4(ipv4) => ipv4,
566        IpAddr::V6(_) => return Err(Error::InvalidTarget("IPv6 not supported".to_string())),
567    };
568
569    // Start background connections to create load
570    let handles: Vec<_> = (0..concurrent).map(|_| {
571        let addr = SocketAddr::new(IpAddr::V4(target_ip), port);
572        tokio::spawn(async move {
573            // Try to establish and hold connection
574            let _ = tokio::net::TcpStream::connect(addr).await;
575            tokio::time::sleep(Duration::from_secs(5)).await;
576        })
577    }).collect();
578
579    // Give connections time to establish
580    tokio::time::sleep(Duration::from_millis(500)).await;
581
582    // Measure latency while load is active
583    let mut rtts = Vec::with_capacity(samples);
584    for _ in 0..samples {
585        let rtt = tcp_ping(target_ip, port, Duration::from_secs(5)).await;
586        rtts.push(rtt);
587        tokio::time::sleep(Duration::from_millis(50)).await;
588    }
589
590    // Cancel background tasks
591    for h in handles {
592        h.abort();
593    }
594
595    Ok(LatencyStats::from_samples(&rtts))
596}
597
598// ============================================================================
599// PACKET REORDERING
600// ============================================================================
601
602/// Result of packet reordering analysis
603#[derive(Debug, Clone)]
604pub struct ReorderingResult {
605    /// Number of packets sent
606    pub packets_sent: usize,
607    /// Number of packets received
608    pub packets_received: usize,
609    /// Number of out-of-order packets
610    pub out_of_order: usize,
611    /// Reordering rate (0.0 - 1.0)
612    pub reorder_rate: f64,
613    /// Maximum reordering extent (how far out of order)
614    pub max_reorder_extent: usize,
615    /// Duplicate packets received
616    pub duplicates: usize,
617}
618
619impl ReorderingResult {
620    /// Check if significant reordering was detected
621    pub fn has_reordering(&self) -> bool {
622        self.reorder_rate > 0.01 // More than 1% reordered
623    }
624}
625
626/// Analyze packet reordering on a path
627///
628/// Sends numbered packets and tracks arrival order
629pub async fn analyze_reordering(
630    target: &str,
631    port: u16,
632    packet_count: usize,
633) -> crate::Result<ReorderingResult> {
634    let dns_result = dns::resolve_ipv4(target).await?;
635    let target_ip = match dns_result.ip {
636        IpAddr::V4(ipv4) => ipv4,
637        IpAddr::V6(_) => return Err(Error::InvalidTarget("IPv6 not supported".to_string())),
638    };
639
640    // Track sequence numbers
641    let mut received_order: VecDeque<usize> = VecDeque::new();
642    let mut expected_next = 0usize;
643    let mut out_of_order = 0usize;
644    let mut max_extent = 0usize;
645    let mut duplicates = 0usize;
646    let mut received_set = std::collections::HashSet::new();
647
648    for seq in 0..packet_count {
649        // Send probe with sequence number
650        let rtt = tcp_ping(target_ip, port, Duration::from_secs(2)).await;
651
652        if rtt.is_some() {
653            if received_set.contains(&seq) {
654                duplicates += 1;
655            } else {
656                received_set.insert(seq);
657                received_order.push_back(seq);
658
659                if seq != expected_next {
660                    out_of_order += 1;
661                    let extent = seq.abs_diff(expected_next);
662                    max_extent = std::cmp::max(max_extent, extent);
663                }
664                expected_next = seq + 1;
665            }
666        }
667
668        // Small delay between packets
669        tokio::time::sleep(Duration::from_millis(10)).await;
670    }
671
672    let packets_received = received_set.len();
673    let reorder_rate = if packets_received > 0 {
674        out_of_order as f64 / packets_received as f64
675    } else {
676        0.0
677    };
678
679    Ok(ReorderingResult {
680        packets_sent: packet_count,
681        packets_received,
682        out_of_order,
683        reorder_rate,
684        max_reorder_extent: max_extent,
685        duplicates,
686    })
687}
688
689#[cfg(test)]
690mod tests {
691    use super::*;
692
693    #[test]
694    fn test_latency_stats_empty() {
695        let samples: Vec<Option<Duration>> = vec![];
696        let stats = LatencyStats::from_samples(&samples);
697        assert_eq!(stats.sample_count, 0);
698        assert_eq!(stats.loss_rate, 1.0);
699    }
700
701    #[test]
702    fn test_latency_stats_basic() {
703        let samples = vec![
704            Some(Duration::from_millis(10)),
705            Some(Duration::from_millis(20)),
706            Some(Duration::from_millis(15)),
707            None, // Lost packet
708        ];
709        let stats = LatencyStats::from_samples(&samples);
710
711        assert_eq!(stats.sample_count, 4);
712        assert_eq!(stats.success_count, 3);
713        assert_eq!(stats.loss_rate, 0.25);
714        assert_eq!(stats.min_rtt, Duration::from_millis(10));
715        assert_eq!(stats.max_rtt, Duration::from_millis(20));
716    }
717
718    #[test]
719    fn test_bufferbloat_grade() {
720        assert_eq!(BufferbloatGrade::from_increase(2.0), BufferbloatGrade::A);
721        assert_eq!(BufferbloatGrade::from_increase(15.0), BufferbloatGrade::B);
722        assert_eq!(BufferbloatGrade::from_increase(45.0), BufferbloatGrade::C);
723        assert_eq!(BufferbloatGrade::from_increase(100.0), BufferbloatGrade::D);
724        assert_eq!(BufferbloatGrade::from_increase(300.0), BufferbloatGrade::F);
725    }
726
727    #[test]
728    fn test_pmtud_options_default() {
729        let opts = PmtudOptions::default();
730        assert_eq!(opts.min_mtu, 68);
731        assert_eq!(opts.max_mtu, 1500);
732    }
733}