Skip to main content

fips_core/proto/mmp/
metrics.rs

1//! MMP derived metrics.
2//!
3//! `MmpMetrics` processes incoming ReceiverReports (from our peer) and
4//! maintains derived metrics: SRTT, loss rate, goodput, ETX, and dual
5//! EWMA trend indicators. Updated by the sender side when it receives
6//! a ReceiverReport about its own traffic.
7
8use super::algorithms::{DualEwma, SrttEstimator, compute_etx};
9use super::report::ReceiverReport;
10use std::time::Instant;
11use tracing::trace;
12
13/// Derived MMP metrics, updated from incoming ReceiverReports.
14///
15/// This lives on the sender side: when we receive a ReceiverReport from
16/// our peer describing what they observed about our traffic, we process
17/// it here to compute RTT, loss, goodput, and trend indicators.
18pub struct MmpMetrics {
19    /// Smoothed RTT from timestamp echo.
20    pub srtt: SrttEstimator,
21
22    /// Dual EWMA trend detectors.
23    pub rtt_trend: DualEwma,
24    pub loss_trend: DualEwma,
25    pub goodput_trend: DualEwma,
26    pub jitter_trend: DualEwma,
27    pub etx_trend: DualEwma,
28
29    /// Forward delivery ratio (what fraction of our frames the peer received).
30    pub delivery_ratio_forward: f64,
31    /// Reverse delivery ratio (set when we compute from our own receiver state).
32    pub delivery_ratio_reverse: f64,
33    /// ETX computed from bidirectional delivery ratios.
34    pub etx: f64,
35
36    /// Smoothed goodput in bytes/sec (forward direction: what the peer received from us).
37    pub goodput_bps: f64,
38
39    // --- State for delta computation ---
40    /// Previous ReceiverReport's cumulative counters (for computing interval deltas).
41    prev_rr_cum_packets: u64,
42    prev_rr_cum_bytes: u64,
43    prev_rr_highest_counter: u64,
44    prev_rr_ecn_ce: u32,
45    prev_rr_reorder: u32,
46    /// Time of previous ReceiverReport (for goodput rate computation).
47    prev_rr_time: Option<Instant>,
48    /// Time of the most recent accepted RTT sample.
49    last_srtt_update: Option<Instant>,
50    /// Whether we have a previous ReceiverReport for delta computation.
51    has_prev_rr: bool,
52    /// Counter span in the most recent ReceiverReport delta.
53    last_forward_counter_span: u64,
54    /// Loss rate in the most recent ReceiverReport delta.
55    last_forward_loss_rate: Option<f64>,
56    /// Accumulated low-rate forward loss evidence since the last actionable
57    /// route-quality sample was emitted.
58    forward_loss_window_span: u64,
59    forward_loss_window_lost: u64,
60
61    // --- State for reverse delivery ratio delta computation ---
62    /// Previous reverse-side cumulative packets received (our receiver state).
63    prev_reverse_packets: u64,
64    /// Previous reverse-side highest counter (our receiver state).
65    prev_reverse_highest: u64,
66    /// Whether we have a previous reverse-side snapshot for delta computation.
67    has_prev_reverse: bool,
68}
69
70impl MmpMetrics {
71    /// Discard deltas that straddle an outbound carrier change.
72    ///
73    /// ReceiverReport counters are cumulative for the FSP session and do not
74    /// identify which next hop carried each packet. The first report after a
75    /// route change therefore cannot be attributed to either the old or new
76    /// carrier; retain its counters only as the new route's baseline.
77    pub fn reset_forward_report_baseline(&mut self) {
78        self.prev_rr_cum_packets = 0;
79        self.prev_rr_cum_bytes = 0;
80        self.prev_rr_highest_counter = 0;
81        self.prev_rr_ecn_ce = 0;
82        self.prev_rr_reorder = 0;
83        self.prev_rr_time = None;
84        self.has_prev_rr = false;
85        self.last_forward_counter_span = 0;
86        self.last_forward_loss_rate = None;
87        self.forward_loss_window_span = 0;
88        self.forward_loss_window_lost = 0;
89    }
90
91    /// Reset state derived from ReceiverReport counters for rekey cutover.
92    ///
93    /// The new session starts with counter 0, so the prev_rr deltas must
94    /// be reset to avoid computing bogus loss/goodput from the counter
95    /// discontinuity. RTT (SRTT) is preserved since it remains valid.
96    pub fn reset_for_rekey(&mut self) {
97        self.reset_forward_report_baseline();
98        self.delivery_ratio_forward = 1.0;
99        self.prev_reverse_packets = 0;
100        self.prev_reverse_highest = 0;
101        self.has_prev_reverse = false;
102        // Keep srtt, srtt freshness, etx, trends, goodput_bps — they'll refresh from data.
103    }
104
105    pub fn new() -> Self {
106        Self {
107            srtt: SrttEstimator::new(),
108            rtt_trend: DualEwma::new(),
109            loss_trend: DualEwma::new(),
110            goodput_trend: DualEwma::new(),
111            jitter_trend: DualEwma::new(),
112            etx_trend: DualEwma::new(),
113            delivery_ratio_forward: 1.0,
114            delivery_ratio_reverse: 1.0,
115            etx: 1.0,
116            goodput_bps: 0.0,
117            prev_rr_cum_packets: 0,
118            prev_rr_cum_bytes: 0,
119            prev_rr_highest_counter: 0,
120            prev_rr_ecn_ce: 0,
121            prev_rr_reorder: 0,
122            prev_rr_time: None,
123            last_srtt_update: None,
124            has_prev_rr: false,
125            last_forward_counter_span: 0,
126            last_forward_loss_rate: None,
127            forward_loss_window_span: 0,
128            forward_loss_window_lost: 0,
129            prev_reverse_packets: 0,
130            prev_reverse_highest: 0,
131            has_prev_reverse: false,
132        }
133    }
134
135    /// Process an incoming ReceiverReport (from the peer about our traffic).
136    ///
137    /// `our_timestamp_ms` is the current session-relative time in ms (for RTT).
138    /// `now` is the current monotonic time (for goodput rate computation).
139    ///
140    /// Returns `true` if this report produced the first SRTT measurement
141    /// (transition from uninitialized to initialized).
142    pub fn process_receiver_report(
143        &mut self,
144        rr: &ReceiverReport,
145        our_timestamp_ms: u32,
146        now: Instant,
147    ) -> bool {
148        let had_srtt = self.srtt.initialized();
149
150        if self.has_prev_rr {
151            let counters_regressed = rr.highest_counter < self.prev_rr_highest_counter
152                || rr.cumulative_packets_recv < self.prev_rr_cum_packets
153                || rr.cumulative_bytes_recv < self.prev_rr_cum_bytes
154                || rr.ecn_ce_count < self.prev_rr_ecn_ce
155                || rr.cumulative_reorder_count < self.prev_rr_reorder;
156            let duplicate_counters = rr.highest_counter == self.prev_rr_highest_counter
157                && rr.cumulative_packets_recv == self.prev_rr_cum_packets
158                && rr.cumulative_bytes_recv == self.prev_rr_cum_bytes
159                && rr.ecn_ce_count == self.prev_rr_ecn_ce
160                && rr.cumulative_reorder_count == self.prev_rr_reorder;
161            if counters_regressed || duplicate_counters {
162                trace!(
163                    highest_counter = rr.highest_counter,
164                    prev_highest_counter = self.prev_rr_highest_counter,
165                    cumulative_packets_recv = rr.cumulative_packets_recv,
166                    prev_cumulative_packets_recv = self.prev_rr_cum_packets,
167                    cumulative_bytes_recv = rr.cumulative_bytes_recv,
168                    prev_cumulative_bytes_recv = self.prev_rr_cum_bytes,
169                    "Ignoring stale MMP ReceiverReport"
170                );
171                return false;
172            }
173        }
174
175        // --- RTT from timestamp echo ---
176        // RTT = now - echoed_timestamp - dwell_time
177        if rr.timestamp_echo > 0 {
178            let echo_ms = rr.timestamp_echo;
179            let dwell_ms = u32::from(rr.dwell_time);
180            let rtt_sample_ms = echo_ms
181                .checked_add(dwell_ms)
182                .and_then(|send_done_ms| our_timestamp_ms.checked_sub(send_done_ms));
183
184            match rtt_sample_ms {
185                Some(rtt_ms) if rtt_ms > 0 => {
186                    let rtt_us = (rtt_ms as i64) * 1000;
187                    trace!(
188                        our_ts = our_timestamp_ms,
189                        echo = echo_ms,
190                        dwell = dwell_ms,
191                        rtt_ms = rtt_ms,
192                        srtt_ms = self.srtt.srtt_us() as f64 / 1000.0,
193                        "RTT sample from timestamp echo"
194                    );
195                    self.srtt.update(rtt_us);
196                    self.last_srtt_update = Some(now);
197                    self.rtt_trend.update(rtt_us as f64);
198                }
199                _ => {
200                    trace!(
201                        our_ts = our_timestamp_ms,
202                        echo = echo_ms,
203                        dwell = dwell_ms,
204                        "Ignoring invalid MMP RTT sample"
205                    );
206                }
207            }
208        }
209
210        // --- Loss rate from cumulative counters ---
211        // Delta: frames the peer should have received vs. actually received
212        self.last_forward_counter_span = 0;
213        self.last_forward_loss_rate = None;
214        if self.has_prev_rr {
215            let counter_span = rr
216                .highest_counter
217                .saturating_sub(self.prev_rr_highest_counter);
218            let packets_delta = rr
219                .cumulative_packets_recv
220                .saturating_sub(self.prev_rr_cum_packets);
221
222            if counter_span > 0 {
223                let delivery = (packets_delta as f64) / (counter_span as f64);
224                self.delivery_ratio_forward = delivery.clamp(0.0, 1.0);
225                let loss_rate = 1.0 - self.delivery_ratio_forward;
226                self.last_forward_counter_span = counter_span;
227                self.last_forward_loss_rate = Some(loss_rate);
228                self.forward_loss_window_span =
229                    self.forward_loss_window_span.saturating_add(counter_span);
230                self.forward_loss_window_lost = self
231                    .forward_loss_window_lost
232                    .saturating_add(counter_span.saturating_sub(packets_delta));
233                self.loss_trend.update(loss_rate);
234                self.etx = compute_etx(self.delivery_ratio_forward, self.delivery_ratio_reverse);
235                self.etx_trend.update(self.etx);
236            }
237        }
238
239        // --- Goodput from cumulative bytes + time delta ---
240        if self.has_prev_rr {
241            let bytes_delta = rr
242                .cumulative_bytes_recv
243                .saturating_sub(self.prev_rr_cum_bytes);
244            self.goodput_trend.update(bytes_delta as f64);
245
246            // Compute bytes/sec if we have a time reference
247            if let Some(prev_time) = self.prev_rr_time {
248                let elapsed = now.duration_since(prev_time);
249                let secs = elapsed.as_secs_f64();
250                if secs > 0.0 {
251                    let bps = bytes_delta as f64 / secs;
252                    // EWMA smoothing: α = 1/4
253                    if self.goodput_bps == 0.0 {
254                        self.goodput_bps = bps;
255                    } else {
256                        self.goodput_bps += (bps - self.goodput_bps) * 0.25;
257                    }
258                }
259            }
260        }
261
262        // --- Jitter trend ---
263        self.jitter_trend.update(rr.jitter as f64);
264
265        // --- Save for next delta ---
266        self.prev_rr_cum_packets = rr.cumulative_packets_recv;
267        self.prev_rr_cum_bytes = rr.cumulative_bytes_recv;
268        self.prev_rr_highest_counter = rr.highest_counter;
269        self.prev_rr_ecn_ce = rr.ecn_ce_count;
270        self.prev_rr_reorder = rr.cumulative_reorder_count;
271        self.prev_rr_time = Some(now);
272        self.has_prev_rr = true;
273
274        !had_srtt && self.srtt.initialized()
275    }
276
277    /// Update the reverse delivery ratio from our own receiver state.
278    ///
279    /// Computes a per-interval delta (same as forward ratio) rather than
280    /// a lifetime cumulative ratio, so ETX responds to recent conditions.
281    pub fn update_reverse_delivery(&mut self, our_recv_packets: u64, peer_highest: u64) {
282        if self.has_prev_reverse {
283            let counter_span = peer_highest.saturating_sub(self.prev_reverse_highest);
284            let packets_delta = our_recv_packets.saturating_sub(self.prev_reverse_packets);
285
286            if counter_span > 0 {
287                let delivery = (packets_delta as f64) / (counter_span as f64);
288                self.delivery_ratio_reverse = delivery.clamp(0.0, 1.0);
289                self.etx = compute_etx(self.delivery_ratio_forward, self.delivery_ratio_reverse);
290                self.etx_trend.update(self.etx);
291            }
292        }
293
294        self.prev_reverse_packets = our_recv_packets;
295        self.prev_reverse_highest = peer_highest;
296        self.has_prev_reverse = true;
297    }
298
299    /// Current smoothed RTT in milliseconds, or `None` if not yet measured.
300    pub fn srtt_ms(&self) -> Option<f64> {
301        if self.srtt.initialized() {
302            Some(self.srtt.srtt_us() as f64 / 1000.0)
303        } else {
304            None
305        }
306    }
307
308    /// Age of the current SRTT sample in milliseconds.
309    pub fn srtt_age_ms(&self, now: Instant) -> Option<u64> {
310        self.last_srtt_update.map(|updated_at| {
311            now.saturating_duration_since(updated_at)
312                .as_millis()
313                .min(u128::from(u64::MAX)) as u64
314        })
315    }
316
317    /// Current loss rate (0.0 = no loss, 1.0 = total loss).
318    pub fn loss_rate(&self) -> f64 {
319        1.0 - self.delivery_ratio_forward
320    }
321
322    /// Smoothed loss rate (long-term EWMA), or `None` if not yet initialized.
323    pub fn smoothed_loss(&self) -> Option<f64> {
324        if self.loss_trend.initialized() {
325            Some(self.loss_trend.long())
326        } else {
327            None
328        }
329    }
330
331    /// Most recent forward-loss sample from a ReceiverReport delta.
332    pub fn last_forward_loss_sample(&self) -> Option<(u64, f64)> {
333        self.last_forward_loss_rate
334            .map(|loss| (self.last_forward_counter_span, loss))
335    }
336
337    /// Return route-quality loss evidence once enough packets have been
338    /// observed. High-rate reports are emitted as-is; low-rate reports, such as
339    /// interactive pings, accumulate until they have the same weight.
340    pub fn take_forward_loss_evidence(&mut self, min_span: u64) -> Option<(u64, f64)> {
341        if min_span == 0 {
342            return self.last_forward_loss_sample();
343        }
344
345        if self.last_forward_counter_span >= min_span
346            && let Some(loss) = self.last_forward_loss_rate
347        {
348            self.forward_loss_window_span = 0;
349            self.forward_loss_window_lost = 0;
350            return Some((self.last_forward_counter_span, loss));
351        }
352
353        if self.forward_loss_window_span >= min_span {
354            let span = self.forward_loss_window_span;
355            let loss = (self.forward_loss_window_lost as f64 / span as f64).clamp(0.0, 1.0);
356            self.forward_loss_window_span = 0;
357            self.forward_loss_window_lost = 0;
358            return Some((span, loss));
359        }
360
361        None
362    }
363
364    /// Smoothed ETX (long-term EWMA), or `None` if not yet initialized.
365    pub fn smoothed_etx(&self) -> Option<f64> {
366        if self.etx_trend.initialized() {
367            Some(self.etx_trend.long())
368        } else {
369            None
370        }
371    }
372
373    /// Current smoothed goodput in bytes/sec, or 0 if not yet measured.
374    pub fn goodput_bps(&self) -> f64 {
375        self.goodput_bps
376    }
377
378    /// Cumulative ECN CE count from the most recent ReceiverReport.
379    pub fn last_ecn_ce_count(&self) -> u32 {
380        self.prev_rr_ecn_ce
381    }
382}
383
384impl Default for MmpMetrics {
385    fn default() -> Self {
386        Self::new()
387    }
388}
389
390// ============================================================================
391// Tests
392// ============================================================================
393
394#[cfg(test)]
395mod tests {
396    use super::*;
397    use std::time::Duration;
398
399    fn make_rr(
400        highest_counter: u64,
401        cum_packets: u64,
402        cum_bytes: u64,
403        timestamp_echo: u32,
404        dwell: u16,
405        jitter: u32,
406    ) -> ReceiverReport {
407        ReceiverReport {
408            highest_counter,
409            cumulative_packets_recv: cum_packets,
410            cumulative_bytes_recv: cum_bytes,
411            timestamp_echo,
412            dwell_time: dwell,
413            max_burst_loss: 0,
414            mean_burst_loss: 0,
415            jitter,
416            ecn_ce_count: 0,
417            owd_trend: 0,
418            burst_loss_count: 0,
419            cumulative_reorder_count: 0,
420            interval_packets_recv: 0,
421            interval_bytes_recv: 0,
422        }
423    }
424
425    #[test]
426    fn test_rtt_from_echo() {
427        let mut m = MmpMetrics::new();
428        let now = Instant::now();
429        // Peer echoes timestamp 1000ms, dwell=5ms, our current time=1050ms
430        let rr = make_rr(10, 10, 5000, 1000, 5, 0);
431        m.process_receiver_report(&rr, 1050, now);
432
433        assert!(m.srtt.initialized());
434        // RTT = 1050 - 1000 - 5 = 45ms
435        let srtt_ms = m.srtt_ms().unwrap();
436        assert!((srtt_ms - 45.0).abs() < 1.0, "srtt={srtt_ms}, expected ~45");
437    }
438
439    #[test]
440    fn test_ignores_duplicate_receiver_report_after_valid_sample() {
441        let mut m = MmpMetrics::new();
442        let now = Instant::now();
443
444        let valid_rr = make_rr(10, 10, 5000, 1000, 5, 0);
445        m.process_receiver_report(&valid_rr, 1050, now);
446        let baseline_srtt_ms = m.srtt_ms().unwrap();
447        assert_eq!(m.srtt_age_ms(now), Some(0));
448
449        // A duplicate of the same counters arriving later would be a 5s RTT
450        // sample if accepted. It is stale and must not move SRTT.
451        m.process_receiver_report(&valid_rr, 6000, now + Duration::from_secs(5));
452
453        let srtt_ms = m.srtt_ms().unwrap();
454        assert_eq!(srtt_ms, baseline_srtt_ms);
455        assert_eq!(m.srtt_age_ms(now + Duration::from_secs(5)), Some(5000));
456    }
457
458    #[test]
459    fn test_ignores_out_of_order_receiver_report_after_valid_sample() {
460        let mut m = MmpMetrics::new();
461        let now = Instant::now();
462
463        let valid_rr = make_rr(20, 20, 10000, 1000, 5, 0);
464        m.process_receiver_report(&valid_rr, 1050, now);
465        let baseline_srtt_ms = m.srtt_ms().unwrap();
466
467        let old_rr = make_rr(10, 10, 5000, 1000, 0, 0);
468        m.process_receiver_report(&old_rr, 6000, now + Duration::from_secs(5));
469
470        let srtt_ms = m.srtt_ms().unwrap();
471        assert_eq!(srtt_ms, baseline_srtt_ms);
472    }
473
474    #[test]
475    fn test_ignores_wrapped_rtt_sample() {
476        let mut m = MmpMetrics::new();
477        let now = Instant::now();
478
479        let wrapped_rr = make_rr(10, 10, 5000, u32::MAX - 10, 20, 0);
480        m.process_receiver_report(&wrapped_rr, 15, now);
481
482        assert!(m.srtt_ms().is_none());
483    }
484
485    #[test]
486    fn test_loss_rate_computation() {
487        let mut m = MmpMetrics::new();
488        let t0 = Instant::now();
489
490        // First report: baseline
491        let rr1 = make_rr(100, 100, 50000, 0, 0, 0);
492        m.process_receiver_report(&rr1, 0, t0);
493
494        // Second report: 200 counters sent, 190 received (5% loss)
495        let rr2 = make_rr(300, 290, 145000, 0, 0, 0);
496        m.process_receiver_report(&rr2, 0, t0 + Duration::from_secs(1));
497
498        let loss = m.loss_rate();
499        assert!((loss - 0.05).abs() < 0.01, "loss={loss}, expected ~0.05");
500        assert_eq!(m.last_forward_loss_sample(), Some((200, loss)));
501    }
502
503    #[test]
504    fn test_etx_updates() {
505        let mut m = MmpMetrics::new();
506        assert_eq!(m.etx, 1.0); // initial: perfect
507
508        // Simulate some loss via forward ratio
509        m.delivery_ratio_forward = 0.9;
510
511        // First call establishes the baseline (no ETX update yet)
512        m.update_reverse_delivery(100, 100);
513        assert_eq!(m.etx, 1.0); // still perfect — baseline only
514
515        // Second call: 190 of 200 frames received (5% loss)
516        m.update_reverse_delivery(290, 300);
517        assert!(m.etx > 1.0);
518        assert!(m.etx < 2.0);
519    }
520
521    #[test]
522    fn test_no_rtt_without_echo() {
523        let mut m = MmpMetrics::new();
524        let now = Instant::now();
525        let rr = make_rr(10, 10, 5000, 0, 0, 0);
526        m.process_receiver_report(&rr, 1000, now);
527        assert!(m.srtt_ms().is_none());
528    }
529
530    #[test]
531    fn test_jitter_trend() {
532        let mut m = MmpMetrics::new();
533        let t0 = Instant::now();
534        let rr1 = make_rr(10, 10, 5000, 0, 0, 100);
535        m.process_receiver_report(&rr1, 0, t0);
536
537        let rr2 = make_rr(20, 20, 10000, 0, 0, 500);
538        m.process_receiver_report(&rr2, 0, t0 + Duration::from_secs(1));
539
540        assert!(m.jitter_trend.initialized());
541        // Short-term should be closer to 500 than long-term
542        assert!(m.jitter_trend.short() > m.jitter_trend.long());
543    }
544
545    #[test]
546    fn test_goodput_bps() {
547        let mut m = MmpMetrics::new();
548        let t0 = Instant::now();
549
550        // First report: baseline (50KB received)
551        let rr1 = make_rr(100, 100, 50_000, 0, 0, 0);
552        m.process_receiver_report(&rr1, 0, t0);
553        assert_eq!(m.goodput_bps(), 0.0); // no rate yet (first report)
554
555        // Second report 1s later: 150KB total (100KB delta in 1s = 100KB/s)
556        let rr2 = make_rr(300, 290, 150_000, 0, 0, 0);
557        m.process_receiver_report(&rr2, 0, t0 + Duration::from_secs(1));
558        assert!(
559            m.goodput_bps() > 90_000.0,
560            "goodput={}, expected ~100000",
561            m.goodput_bps()
562        );
563        assert!(
564            m.goodput_bps() < 110_000.0,
565            "goodput={}, expected ~100000",
566            m.goodput_bps()
567        );
568    }
569
570    #[test]
571    fn test_reverse_delivery_delta() {
572        let mut m = MmpMetrics::new();
573
574        // First call: baseline only, no ratio update
575        m.update_reverse_delivery(100, 100);
576        assert_eq!(m.delivery_ratio_reverse, 1.0); // unchanged from default
577
578        // Second call: perfect delivery (200 new frames, all received)
579        m.update_reverse_delivery(300, 300);
580        assert!((m.delivery_ratio_reverse - 1.0).abs() < 0.001);
581
582        // Third call: 50% loss (100 frames sent, 50 received)
583        m.update_reverse_delivery(350, 400);
584        assert!(
585            (m.delivery_ratio_reverse - 0.5).abs() < 0.001,
586            "reverse={}, expected 0.5",
587            m.delivery_ratio_reverse
588        );
589    }
590
591    #[test]
592    fn test_reverse_delivery_rekey_reset() {
593        let mut m = MmpMetrics::new();
594
595        // Establish baseline and one measurement
596        m.update_reverse_delivery(100, 100);
597        m.update_reverse_delivery(300, 300);
598        assert!((m.delivery_ratio_reverse - 1.0).abs() < 0.001);
599
600        // Rekey resets reverse state
601        m.reset_for_rekey();
602
603        // First call after rekey: baseline only
604        m.update_reverse_delivery(50, 50);
605        // delivery_ratio_reverse was reset to 1.0 by reset_for_rekey's
606        // clearing of delivery_ratio_forward; reverse is not explicitly
607        // reset — but the delta state is, so next call computes fresh.
608        assert_eq!(m.delivery_ratio_reverse, 1.0);
609
610        // Second call after rekey: 80% delivery
611        m.update_reverse_delivery(90, 100);
612        assert!(
613            (m.delivery_ratio_reverse - 0.8).abs() < 0.001,
614            "reverse={}, expected 0.8",
615            m.delivery_ratio_reverse
616        );
617    }
618}