1use super::algorithms::{DualEwma, SrttEstimator, compute_etx};
9use super::report::ReceiverReport;
10use std::time::Instant;
11use tracing::trace;
12
13pub struct MmpMetrics {
19 pub srtt: SrttEstimator,
21
22 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 pub delivery_ratio_forward: f64,
31 pub delivery_ratio_reverse: f64,
33 pub etx: f64,
35
36 pub goodput_bps: f64,
38
39 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 prev_rr_time: Option<Instant>,
48 last_srtt_update: Option<Instant>,
50 has_prev_rr: bool,
52 last_forward_counter_span: u64,
54 last_forward_loss_rate: Option<f64>,
56 forward_loss_window_span: u64,
59 forward_loss_window_lost: u64,
60
61 prev_reverse_packets: u64,
64 prev_reverse_highest: u64,
66 has_prev_reverse: bool,
68}
69
70impl MmpMetrics {
71 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 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 }
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 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 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 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 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 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 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 self.jitter_trend.update(rr.jitter as f64);
264
265 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 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 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 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 pub fn loss_rate(&self) -> f64 {
319 1.0 - self.delivery_ratio_forward
320 }
321
322 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 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 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 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 pub fn goodput_bps(&self) -> f64 {
375 self.goodput_bps
376 }
377
378 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#[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 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 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 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 let rr1 = make_rr(100, 100, 50000, 0, 0, 0);
492 m.process_receiver_report(&rr1, 0, t0);
493
494 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); m.delivery_ratio_forward = 0.9;
510
511 m.update_reverse_delivery(100, 100);
513 assert_eq!(m.etx, 1.0); 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 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 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); 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 m.update_reverse_delivery(100, 100);
576 assert_eq!(m.delivery_ratio_reverse, 1.0); m.update_reverse_delivery(300, 300);
580 assert!((m.delivery_ratio_reverse - 1.0).abs() < 0.001);
581
582 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 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 m.reset_for_rekey();
602
603 m.update_reverse_delivery(50, 50);
605 assert_eq!(m.delivery_ratio_reverse, 1.0);
609
610 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}