1#![allow(dead_code)]
7
8use crate::net::atp::protocol::outcome::AtpOutcome;
9use crate::net::quic_native::{
10 AckRange, PacketNumberSpace, QuicTransportMachine, RttEstimator, SentPacketMeta,
11};
12use serde::{Deserialize, Serialize};
13use std::collections::{HashMap, HashSet, VecDeque};
14use std::time::{Duration, Instant};
15
16macro_rules! try_outcome {
18 ($expr:expr) => {
19 match $expr {
20 AtpOutcome::Ok(v) => v,
21 AtpOutcome::Err(e) => return AtpOutcome::Err(e),
22 AtpOutcome::Cancelled(r) => return AtpOutcome::Cancelled(r),
23 AtpOutcome::Panicked(p) => return AtpOutcome::Panicked(p),
24 }
25 };
26}
27
28pub struct AtpLossDetector {
30 spaces: [SpaceLossState; 3],
32 config: LossDetectionConfig,
34 pattern_analyzer: LossPatternAnalyzer,
36 reordering_tracker: ReorderingTracker,
38 spurious_loss_packets: HashSet<(usize, u64)>,
40 metrics: LossDetectionMetrics,
42}
43
44#[derive(Debug, Clone)]
46struct SpaceLossState {
47 sent_packets: VecDeque<SentPacketMeta>,
49 largest_acked: Option<u64>,
51 largest_acked_time: Option<u64>,
53 loss_timer_deadline: Option<u64>,
55 early_retransmit_deadline: Option<u64>,
57}
58
59#[derive(Debug, Clone, Serialize, Deserialize)]
61pub struct LossDetectionConfig {
62 pub packet_threshold: u32,
64 pub time_threshold_multiplier: f64,
66 pub min_time_threshold_micros: u64,
68 pub max_reordering_threshold: u32,
70 pub adaptive_packet_threshold: bool,
72 pub enable_early_retransmit: bool,
74 pub early_retransmit_threshold: u32,
76}
77
78impl Default for LossDetectionConfig {
79 fn default() -> Self {
80 Self {
81 packet_threshold: 3,
82 time_threshold_multiplier: 9.0 / 8.0,
83 min_time_threshold_micros: 1_000, max_reordering_threshold: 10,
85 adaptive_packet_threshold: true,
86 enable_early_retransmit: true,
87 early_retransmit_threshold: 1,
88 }
89 }
90}
91
92#[derive(Debug, Clone)]
94struct LossPatternAnalyzer {
95 loss_events: VecDeque<LossEvent>,
97 patterns: Vec<LossPattern>,
99 pattern_confidence: HashMap<LossPattern, f64>,
101}
102
103#[derive(Debug, Clone)]
105struct LossEvent {
106 timestamp: Instant,
108 lost_packets: Vec<u64>,
110 detection_method: LossDetectionMethod,
112 conditions: NetworkConditions,
114}
115
116#[derive(Debug, Clone)]
118struct NetworkConditions {
119 rtt_micros: Option<u64>,
121 rttvar_micros: Option<u64>,
123 bytes_in_flight: u64,
125 congestion_window: u64,
127}
128
129#[derive(Debug, Clone, Copy, PartialEq, Eq)]
130struct CanonicalAckRange {
131 smallest: u64,
132 largest: u64,
133}
134
135#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
141pub struct LossTransportState {
142 pub latest_rtt_micros: Option<u64>,
144 pub smoothed_rtt_micros: Option<u64>,
146 pub rttvar_micros: Option<u64>,
148 pub bytes_in_flight: u64,
150 pub congestion_window: u64,
152}
153
154impl LossTransportState {
155 #[must_use]
157 pub fn from_transport(transport: &QuicTransportMachine) -> Self {
158 Self::from_rtt_and_recovery(
159 transport.rtt(),
160 transport.bytes_in_flight(),
161 transport.congestion_window_bytes(),
162 )
163 }
164
165 #[must_use]
167 pub fn from_rtt_and_recovery(
168 rtt: &RttEstimator,
169 bytes_in_flight: u64,
170 congestion_window: u64,
171 ) -> Self {
172 Self {
173 latest_rtt_micros: rtt.latest_rtt_micros(),
174 smoothed_rtt_micros: rtt.smoothed_rtt_micros(),
175 rttvar_micros: rtt.rttvar_micros(),
176 bytes_in_flight,
177 congestion_window,
178 }
179 }
180
181 fn base_rtt_micros(self) -> u64 {
182 self.latest_rtt_micros
183 .or(self.smoothed_rtt_micros)
184 .unwrap_or(333_000)
185 }
186}
187
188#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
190pub enum LossPattern {
191 Sporadic,
193 Burst,
195 Periodic,
197 Reordering,
199 Congestion,
201 Tail,
203}
204
205#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
207pub enum LossDetectionMethod {
208 PacketThreshold,
210 TimeThreshold,
212 EarlyRetransmit,
214 Combined,
216}
217
218#[derive(Debug, Clone)]
220struct ReorderingTracker {
221 reordering_measurements: VecDeque<u32>,
223 current_threshold: u32,
225 max_reordering: u32,
227 adaptation_factor: f64,
229}
230
231#[derive(Debug, Clone, Serialize, Deserialize)]
233pub struct LossDetectionMetrics {
234 pub total_lost_packets: u64,
236 pub packet_threshold_losses: u64,
238 pub time_threshold_losses: u64,
240 pub false_losses: u64,
242 pub avg_packet_threshold: f64,
244 pub avg_time_threshold_micros: f64,
246 pub reordering_events: u64,
248 pub pattern_accuracy: f64,
250}
251
252#[derive(Debug, Clone)]
254pub struct LossDetectionResult {
255 pub lost_packets: Vec<LostPacketInfo>,
257 pub lost_packet_count: u64,
263 pub lost_bytes: u64,
265 pub detection_method: LossDetectionMethod,
267 pub confidence: f64,
269 pub recommendations: Vec<LossRecommendation>,
271}
272
273#[derive(Debug, Clone, Serialize, Deserialize)]
275pub struct LostPacketInfo {
276 pub packet_number: u64,
278 pub bytes: u64,
280 pub sent_time_micros: u64,
282 pub detected_time_micros: u64,
284 pub reason: LossReason,
286}
287
288#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
290pub enum LossReason {
291 PacketThreshold { threshold: u32 },
293 TimeThreshold { threshold_micros: u64 },
295 BothThresholds {
297 packet_threshold: u32,
298 time_threshold_micros: u64,
299 },
300 EarlyRetransmit,
302}
303
304#[derive(Debug, Clone, Serialize, Deserialize)]
306pub enum LossRecommendation {
307 ReduceCongestionWindow { factor: f64 },
309 IncreaseReorderingThreshold { new_threshold: u32 },
311 EnablePacing { rate: u64 },
313 SwitchCongestionControl { algorithm: String },
315 EnableFec { rate: f64 },
317}
318
319#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
321pub struct TransferLossSignal {
322 pub loss_pressure: f64,
324 pub repair_rate_hint: Option<f64>,
326 pub cwnd_reduction_factor: Option<f64>,
328 pub pacing_rate_hint: Option<u64>,
330 pub prefer_bbr: bool,
332 pub confidence: f64,
334}
335
336impl AtpLossDetector {
337 #[must_use]
339 pub fn new() -> Self {
340 Self::with_config(LossDetectionConfig::default())
341 }
342
343 #[must_use]
345 pub fn with_config(config: LossDetectionConfig) -> Self {
346 Self {
347 spaces: [
348 SpaceLossState::new(),
349 SpaceLossState::new(),
350 SpaceLossState::new(),
351 ],
352 config,
353 pattern_analyzer: LossPatternAnalyzer::new(),
354 reordering_tracker: ReorderingTracker::new(),
355 spurious_loss_packets: HashSet::new(),
356 metrics: LossDetectionMetrics::default(),
357 }
358 }
359
360 pub fn on_packet_sent(&mut self, packet: SentPacketMeta) {
362 let space_idx = packet.space as usize;
363 self.spaces[space_idx].sent_packets.push_back(packet);
364
365 if self.spaces[space_idx].sent_packets.len() > 10_000 {
367 self.spaces[space_idx].sent_packets.pop_front();
368 }
369 }
370
371 pub fn on_ack_received(
373 &mut self,
374 space: PacketNumberSpace,
375 ack_ranges: &[AckRange],
376 _ack_delay_micros: u64,
377 now_micros: u64,
378 transport_state: &LossTransportState,
379 ) -> AtpOutcome<LossDetectionResult> {
380 let space_idx = space as usize;
381 let state = &mut self.spaces[space_idx];
382
383 let mut newly_acked = Vec::new();
385 let acked_packet_numbers = acked_sent_packet_index(&state.sent_packets, ack_ranges);
386 let largest_newly_acked = acked_packet_numbers.iter().copied().max();
387
388 let mut remaining_packets = VecDeque::new();
390 while let Some(packet) = state.sent_packets.pop_front() {
391 if acked_packet_numbers.contains(&packet.packet_number) {
392 newly_acked.push(packet);
393 } else {
394 remaining_packets.push_back(packet);
395 }
396 }
397 state.sent_packets = remaining_packets;
398
399 if let Some(largest) = largest_newly_acked {
401 if state.largest_acked.is_none_or(|prev| largest > prev) {
402 state.largest_acked = Some(largest);
403 state.largest_acked_time = Some(now_micros);
404 }
405 }
406
407 let loss_result = try_outcome!(self.detect_losses(space, now_micros, transport_state));
409
410 if !loss_result.lost_packets.is_empty() {
412 self.update_pattern_analysis(&loss_result, transport_state);
413 }
414
415 self.update_reordering_tracking(space_idx, &newly_acked, ack_ranges, &loss_result);
419
420 AtpOutcome::ok(loss_result)
421 }
422
423 #[allow(clippy::cast_precision_loss)]
430 pub fn observe_datagram_loss_sample(
431 &mut self,
432 sent_datagrams: u64,
433 lost_datagrams: u64,
434 rtt: Option<Duration>,
435 bytes_in_flight: u64,
436 congestion_window: u64,
437 ) -> LossDetectionResult {
438 if sent_datagrams == 0 {
439 return LossDetectionResult::empty();
440 }
441
442 let lost_datagrams = lost_datagrams.min(sent_datagrams);
443 if lost_datagrams == 0 {
444 if let Some(rtt) = rtt {
445 self.metrics.avg_time_threshold_micros = rtt.as_micros() as f64;
446 }
447 return LossDetectionResult::empty();
448 }
449
450 let sample_cap = lost_datagrams.min(128);
451 let rtt_micros = rtt.and_then(|duration| u64::try_from(duration.as_micros()).ok());
452 let detected_time_micros = rtt_micros.unwrap_or(1);
453 let lost_packets: Vec<LostPacketInfo> = (0..sample_cap)
454 .map(|packet_number| LostPacketInfo {
455 packet_number,
456 bytes: bytes_in_flight
457 .checked_div(sent_datagrams)
458 .unwrap_or(0)
459 .max(1),
460 sent_time_micros: 0,
461 detected_time_micros,
462 reason: LossReason::PacketThreshold {
463 threshold: self.config.packet_threshold,
464 },
465 })
466 .collect();
467
468 self.metrics.total_lost_packets = self
469 .metrics
470 .total_lost_packets
471 .saturating_add(lost_datagrams);
472 self.metrics.packet_threshold_losses = self
473 .metrics
474 .packet_threshold_losses
475 .saturating_add(lost_datagrams);
476 if let Some(rtt) = rtt {
477 self.metrics.avg_time_threshold_micros = rtt.as_micros() as f64;
478 }
479
480 let detection_method = LossDetectionMethod::PacketThreshold;
481 let confidence = self.calculate_detection_confidence(&lost_packets, detection_method);
482 let recommendations = self.generate_recommendations(&lost_packets, detection_method);
483 let result = LossDetectionResult {
484 lost_packets,
485 lost_packet_count: lost_datagrams,
486 lost_bytes: bytes_in_flight.min(
487 bytes_in_flight
488 .checked_mul(lost_datagrams)
489 .and_then(|bytes| bytes.checked_div(sent_datagrams))
490 .unwrap_or(bytes_in_flight),
491 ),
492 detection_method,
493 confidence,
494 recommendations,
495 };
496 self.update_pattern_analysis(
497 &result,
498 &LossTransportState {
499 latest_rtt_micros: rtt_micros,
500 smoothed_rtt_micros: rtt_micros,
501 rttvar_micros: None,
502 bytes_in_flight,
503 congestion_window,
504 },
505 );
506 result
507 }
508
509 fn detect_losses(
511 &mut self,
512 space: PacketNumberSpace,
513 now_micros: u64,
514 transport_state: &LossTransportState,
515 ) -> AtpOutcome<LossDetectionResult> {
516 let space_idx = space as usize;
517 let Some(largest_acked) = self.spaces[space_idx].largest_acked else {
518 return AtpOutcome::ok(LossDetectionResult::empty());
519 };
520
521 let mut lost_packets = Vec::new();
522 let mut lost_bytes: u64 = 0;
523 let mut detection_methods = Vec::new();
524
525 let packet_threshold = self.get_adaptive_packet_threshold(space);
527 let time_threshold = self.calculate_time_threshold(*transport_state);
528
529 let time_threshold_boundary = now_micros.checked_sub(time_threshold);
533
534 let mut remaining_packets = VecDeque::new();
535 let enable_early_retransmit = self.config.enable_early_retransmit;
536 let early_retransmit_threshold = self.config.early_retransmit_threshold;
537 let state = &mut self.spaces[space_idx];
538 while let Some(packet) = state.sent_packets.pop_front() {
539 let mut is_lost = false;
540 let mut loss_reason = None;
541
542 if packet
544 .packet_number
545 .checked_add(u64::from(packet_threshold))
546 .is_some_and(|threshold_packet| threshold_packet <= largest_acked)
547 {
548 is_lost = true;
549 loss_reason = Some(LossReason::PacketThreshold {
550 threshold: packet_threshold,
551 });
552 detection_methods.push(LossDetectionMethod::PacketThreshold);
553 self.metrics.packet_threshold_losses += 1;
554 }
555
556 if time_threshold_boundary.is_some_and(|boundary| packet.time_sent_micros <= boundary)
558 && packet.packet_number <= largest_acked
559 {
560 if is_lost {
561 loss_reason = Some(LossReason::BothThresholds {
563 packet_threshold,
564 time_threshold_micros: time_threshold,
565 });
566 detection_methods.clear();
567 detection_methods.push(LossDetectionMethod::Combined);
568 } else {
569 is_lost = true;
570 loss_reason = Some(LossReason::TimeThreshold {
571 threshold_micros: time_threshold,
572 });
573 detection_methods.push(LossDetectionMethod::TimeThreshold);
574 self.metrics.time_threshold_losses += 1;
575 }
576 }
577
578 if !is_lost && enable_early_retransmit {
580 if packet
581 .packet_number
582 .checked_add(u64::from(early_retransmit_threshold))
583 == Some(largest_acked)
584 {
585 is_lost = true;
586 loss_reason = Some(LossReason::EarlyRetransmit);
587 detection_methods.push(LossDetectionMethod::EarlyRetransmit);
588 }
589 }
590
591 if is_lost {
592 lost_bytes = lost_bytes.saturating_add(packet.bytes);
593 if let Some(reason) = loss_reason {
594 lost_packets.push(LostPacketInfo {
595 packet_number: packet.packet_number,
596 bytes: packet.bytes,
597 sent_time_micros: packet.time_sent_micros,
598 detected_time_micros: now_micros,
599 reason,
600 });
601 } else {
602 lost_packets.push(LostPacketInfo {
605 packet_number: packet.packet_number,
606 bytes: packet.bytes,
607 sent_time_micros: packet.time_sent_micros,
608 detected_time_micros: now_micros,
609 reason: LossReason::TimeThreshold {
610 threshold_micros: time_threshold,
611 },
612 });
613 }
614 } else {
615 remaining_packets.push_back(packet);
616 }
617 }
618
619 state.sent_packets = remaining_packets;
620 self.metrics.total_lost_packets = self
621 .metrics
622 .total_lost_packets
623 .saturating_add(lost_packets.len() as u64);
624
625 let detection_method = if detection_methods.contains(&LossDetectionMethod::Combined) {
627 LossDetectionMethod::Combined
628 } else if detection_methods.contains(&LossDetectionMethod::PacketThreshold) {
629 LossDetectionMethod::PacketThreshold
630 } else if detection_methods.contains(&LossDetectionMethod::TimeThreshold) {
631 LossDetectionMethod::TimeThreshold
632 } else if detection_methods.contains(&LossDetectionMethod::EarlyRetransmit) {
633 LossDetectionMethod::EarlyRetransmit
634 } else {
635 LossDetectionMethod::PacketThreshold
636 };
637
638 let confidence = self.calculate_detection_confidence(&lost_packets, detection_method);
640
641 let recommendations = self.generate_recommendations(&lost_packets, detection_method);
643
644 AtpOutcome::ok(LossDetectionResult {
645 lost_packet_count: lost_packets.len() as u64,
646 lost_packets,
647 lost_bytes,
648 detection_method,
649 confidence,
650 recommendations,
651 })
652 }
653
654 fn get_adaptive_packet_threshold(&mut self, _space: PacketNumberSpace) -> u32 {
655 if !self.config.adaptive_packet_threshold {
656 return self.config.packet_threshold;
657 }
658
659 let current_threshold = self.reordering_tracker.current_threshold;
661 current_threshold.max(self.config.packet_threshold)
662 }
663
664 fn calculate_time_threshold(&self, transport_state: LossTransportState) -> u64 {
665 let threshold = (transport_state.base_rtt_micros() as f64
666 * self.config.time_threshold_multiplier) as u64;
667 threshold.max(self.config.min_time_threshold_micros)
668 }
669
670 fn should_early_retransmit(&self, packet: &SentPacketMeta, largest_acked: u64) -> bool {
671 packet
673 .packet_number
674 .checked_add(u64::from(self.config.early_retransmit_threshold))
675 == Some(largest_acked)
676 }
677
678 fn calculate_detection_confidence(
679 &self,
680 lost_packets: &[LostPacketInfo],
681 method: LossDetectionMethod,
682 ) -> f64 {
683 if lost_packets.is_empty() {
684 return 1.0;
685 }
686
687 let base_confidence = match method {
689 LossDetectionMethod::Combined => 0.95,
690 LossDetectionMethod::PacketThreshold => 0.85,
691 LossDetectionMethod::TimeThreshold => 0.75,
692 LossDetectionMethod::EarlyRetransmit => 0.60,
693 };
694
695 let pattern_bonus = self
697 .pattern_analyzer
698 .patterns
699 .iter()
700 .map(|pattern| {
701 self.pattern_analyzer
702 .pattern_confidence
703 .get(pattern)
704 .unwrap_or(&0.0)
705 })
706 .fold(0.0_f64, |acc, &conf| acc.max(conf))
707 * 0.1;
708
709 (base_confidence + pattern_bonus).min(1.0_f64)
710 }
711
712 fn generate_recommendations(
713 &self,
714 lost_packets: &[LostPacketInfo],
715 method: LossDetectionMethod,
716 ) -> Vec<LossRecommendation> {
717 let mut recommendations = Vec::new();
718
719 if lost_packets.len() > 5 {
720 recommendations.push(LossRecommendation::ReduceCongestionWindow { factor: 0.5 });
722 }
723
724 if method == LossDetectionMethod::EarlyRetransmit {
725 recommendations.push(LossRecommendation::IncreaseReorderingThreshold {
727 new_threshold: self.reordering_tracker.current_threshold.saturating_add(1),
728 });
729 }
730
731 for pattern in &self.pattern_analyzer.patterns {
733 match pattern {
734 LossPattern::Burst => {
735 recommendations.push(LossRecommendation::EnablePacing { rate: 100_000 });
736 }
737 LossPattern::Periodic => {
738 recommendations.push(LossRecommendation::EnableFec { rate: 0.1 });
739 }
740 LossPattern::Congestion => {
741 recommendations.push(LossRecommendation::SwitchCongestionControl {
742 algorithm: "bbr".to_string(),
743 });
744 }
745 _ => {}
746 }
747 }
748
749 recommendations
750 }
751
752 fn update_pattern_analysis(
753 &mut self,
754 result: &LossDetectionResult,
755 transport_state: &LossTransportState,
756 ) {
757 let loss_event = LossEvent {
758 timestamp: Instant::now(),
759 lost_packets: result
760 .lost_packets
761 .iter()
762 .map(|p| p.packet_number)
763 .collect(),
764 detection_method: result.detection_method,
765 conditions: NetworkConditions {
766 rtt_micros: transport_state
767 .latest_rtt_micros
768 .or(transport_state.smoothed_rtt_micros),
769 rttvar_micros: transport_state.rttvar_micros,
770 bytes_in_flight: transport_state.bytes_in_flight,
771 congestion_window: transport_state.congestion_window,
772 },
773 };
774
775 self.pattern_analyzer.loss_events.push_back(loss_event);
776 if self.pattern_analyzer.loss_events.len() > 1000 {
777 self.pattern_analyzer.loss_events.pop_front();
778 }
779
780 self.analyze_loss_patterns();
781 }
782
783 fn analyze_loss_patterns(&mut self) {
784 self.pattern_analyzer.patterns.clear();
785 self.pattern_analyzer.pattern_confidence.clear();
786
787 if self.pattern_analyzer.loss_events.len() < 3 {
788 return;
789 }
790
791 let recent_events: Vec<_> = self
792 .pattern_analyzer
793 .loss_events
794 .iter()
795 .rev()
796 .take(10)
797 .collect();
798
799 let sample_count = recent_events.len() as f64;
800 let burst_events = recent_events
801 .iter()
802 .filter(|event| event.lost_packets.len() > 3)
803 .count();
804 if burst_events > 0 {
805 self.pattern_analyzer.patterns.push(LossPattern::Burst);
806 self.pattern_analyzer
807 .pattern_confidence
808 .insert(LossPattern::Burst, burst_events as f64 / sample_count);
809 }
810
811 let intervals: Vec<_> = recent_events
812 .windows(2)
813 .map(|w| w[0].timestamp.duration_since(w[1].timestamp))
814 .collect();
815
816 if intervals.len() >= 3 {
817 let interval_micros: Vec<f64> = intervals
818 .iter()
819 .map(|interval| interval.as_micros() as f64)
820 .collect();
821 let avg_interval = interval_micros.iter().sum::<f64>() / interval_micros.len() as f64;
822 if avg_interval > 0.0 {
823 let variance = interval_micros
824 .iter()
825 .map(|interval| {
826 let diff = *interval - avg_interval;
827 diff * diff
828 })
829 .sum::<f64>()
830 / interval_micros.len() as f64;
831 let coefficient_of_variation = variance.sqrt() / avg_interval;
832 if coefficient_of_variation <= 0.15 {
833 self.pattern_analyzer.patterns.push(LossPattern::Periodic);
834 self.pattern_analyzer.pattern_confidence.insert(
835 LossPattern::Periodic,
836 (1.0 - coefficient_of_variation / 0.15).clamp(0.0, 1.0),
837 );
838 }
839 }
840 }
841
842 let congestion_events = recent_events
843 .iter()
844 .filter(|event| {
845 event
846 .conditions
847 .rttvar_micros
848 .zip(event.conditions.rtt_micros)
849 .is_some_and(|(rttvar, rtt)| rtt > 0 && rttvar.saturating_mul(4) > rtt)
850 || (event.conditions.congestion_window > 0
851 && event.conditions.bytes_in_flight
852 >= event.conditions.congestion_window.saturating_mul(9) / 10)
853 })
854 .count();
855 if congestion_events > 0 {
856 self.pattern_analyzer.patterns.push(LossPattern::Congestion);
857 self.pattern_analyzer.pattern_confidence.insert(
858 LossPattern::Congestion,
859 congestion_events as f64 / sample_count,
860 );
861 }
862
863 let tail_events = recent_events
864 .iter()
865 .filter(|event| matches!(event.detection_method, LossDetectionMethod::EarlyRetransmit))
866 .count();
867 if tail_events > 0 {
868 self.pattern_analyzer.patterns.push(LossPattern::Tail);
869 self.pattern_analyzer
870 .pattern_confidence
871 .insert(LossPattern::Tail, tail_events as f64 / sample_count);
872 }
873
874 if self.pattern_analyzer.patterns.is_empty() {
875 self.pattern_analyzer.patterns.push(LossPattern::Sporadic);
876 self.pattern_analyzer
877 .pattern_confidence
878 .insert(LossPattern::Sporadic, 1.0);
879 }
880 }
881
882 fn update_reordering_tracking(
883 &mut self,
884 space_idx: usize,
885 acked_packets: &[SentPacketMeta],
886 ack_ranges: &[AckRange],
887 _loss_result: &LossDetectionResult,
888 ) {
889 let Some(last_loss) = self.pattern_analyzer.loss_events.back() else {
890 return;
891 };
892
893 let canonical_ranges = canonical_ack_ranges(ack_ranges);
894 let mut reordered_packets: HashSet<u64> = HashSet::new();
895
896 for acked in acked_packets {
897 if last_loss.lost_packets.contains(&acked.packet_number) {
898 reordered_packets.insert(acked.packet_number);
899 }
900 }
901
902 for lost_packet in &last_loss.lost_packets {
903 if canonical_ranges_contain_packet(&canonical_ranges, *lost_packet) {
904 reordered_packets.insert(*lost_packet);
905 }
906 }
907
908 let mut newly_recorded_reordering = false;
909 for packet_number in reordered_packets {
910 if !self
911 .spurious_loss_packets
912 .insert((space_idx, packet_number))
913 {
914 continue;
915 }
916
917 newly_recorded_reordering = true;
918 self.metrics.false_losses = self.metrics.false_losses.saturating_add(1);
919 let reordering_depth = last_loss
920 .lost_packets
921 .iter()
922 .copied()
923 .filter(|lost| *lost > packet_number)
924 .count()
925 .saturating_add(1);
926 self.reordering_tracker
927 .record_reordering(u32::try_from(reordering_depth).unwrap_or(u32::MAX));
928 }
929
930 if !self
931 .pattern_analyzer
932 .patterns
933 .contains(&LossPattern::Reordering)
934 && newly_recorded_reordering
935 {
936 self.pattern_analyzer.patterns.push(LossPattern::Reordering);
937 }
938 if newly_recorded_reordering {
939 self.pattern_analyzer.pattern_confidence.insert(
940 LossPattern::Reordering,
941 self.reordering_tracker.confidence(),
942 );
943 }
944 }
945
946 #[must_use]
948 pub fn metrics(&self) -> &LossDetectionMetrics {
949 &self.metrics
950 }
951
952 #[must_use]
954 pub fn export_analysis(&self) -> LossAnalysisExport {
955 LossAnalysisExport {
956 metrics: self.metrics.clone(),
957 patterns: self.pattern_analyzer.patterns.clone(),
958 pattern_confidence: self.pattern_analyzer.pattern_confidence.clone(),
959 config: self.config.clone(),
960 }
961 }
962}
963
964impl Default for AtpLossDetector {
965 fn default() -> Self {
966 Self::new()
967 }
968}
969
970impl SpaceLossState {
971 fn new() -> Self {
972 Self {
973 sent_packets: VecDeque::new(),
974 largest_acked: None,
975 largest_acked_time: None,
976 loss_timer_deadline: None,
977 early_retransmit_deadline: None,
978 }
979 }
980}
981
982impl LossPatternAnalyzer {
983 fn new() -> Self {
984 Self {
985 loss_events: VecDeque::new(),
986 patterns: Vec::new(),
987 pattern_confidence: HashMap::new(),
988 }
989 }
990}
991
992impl ReorderingTracker {
993 fn new() -> Self {
994 Self {
995 reordering_measurements: VecDeque::new(),
996 current_threshold: 3, max_reordering: 0,
998 adaptation_factor: 0.1,
999 }
1000 }
1001
1002 fn adapt_threshold(&mut self) {
1003 self.record_reordering(self.current_threshold);
1004 }
1005
1006 fn record_reordering(&mut self, depth: u32) {
1007 let measured_depth = depth.max(1);
1008 self.max_reordering = self.max_reordering.max(measured_depth);
1009 let target_threshold = self.current_threshold.max(measured_depth.saturating_add(1));
1010 let blended = self.current_threshold as f64
1011 + (target_threshold as f64 - self.current_threshold as f64) * self.adaptation_factor;
1012 self.current_threshold = blended.ceil() as u32;
1013 self.current_threshold = self.current_threshold.min(10);
1014 self.reordering_measurements.push_back(measured_depth);
1015 if self.reordering_measurements.len() > 100 {
1016 self.reordering_measurements.pop_front();
1017 }
1018 }
1019
1020 fn confidence(&self) -> f64 {
1021 if self.reordering_measurements.is_empty() {
1022 return 0.0;
1023 }
1024
1025 let recent = self.reordering_measurements.len().min(20);
1026 let recent_sum = self
1027 .reordering_measurements
1028 .iter()
1029 .rev()
1030 .take(recent)
1031 .copied()
1032 .map(f64::from)
1033 .sum::<f64>();
1034 let recent_avg = recent_sum / recent as f64;
1035 (recent_avg / f64::from(self.current_threshold.max(1))).clamp(0.0, 1.0)
1036 }
1037}
1038
1039impl Default for LossDetectionMetrics {
1040 fn default() -> Self {
1041 Self {
1042 total_lost_packets: 0,
1043 packet_threshold_losses: 0,
1044 time_threshold_losses: 0,
1045 false_losses: 0,
1046 avg_packet_threshold: 3.0,
1047 avg_time_threshold_micros: 333_000.0,
1048 reordering_events: 0,
1049 pattern_accuracy: 0.0,
1050 }
1051 }
1052}
1053
1054impl LossDetectionResult {
1055 fn empty() -> Self {
1056 Self {
1057 lost_packets: Vec::new(),
1058 lost_packet_count: 0,
1059 lost_bytes: 0,
1060 detection_method: LossDetectionMethod::PacketThreshold,
1061 confidence: 1.0,
1062 recommendations: Vec::new(),
1063 }
1064 }
1065
1066 #[must_use]
1073 pub fn transfer_signal(&self, sent_units: u64) -> TransferLossSignal {
1074 let observed_pressure = if sent_units == 0 {
1075 0.0
1076 } else {
1077 self.lost_packet_count.min(sent_units) as f64 / sent_units as f64
1078 };
1079
1080 let mut loss_pressure = observed_pressure;
1081 let mut repair_rate_hint = None;
1082 let mut cwnd_reduction_factor = None;
1083 let mut pacing_rate_hint = None;
1084 let mut prefer_bbr = false;
1085
1086 for recommendation in &self.recommendations {
1087 match recommendation {
1088 LossRecommendation::ReduceCongestionWindow { factor } => {
1089 let factor = finite_unit(*factor);
1090 cwnd_reduction_factor = Some(
1091 cwnd_reduction_factor.map_or(factor, |current: f64| current.min(factor)),
1092 );
1093 }
1094 LossRecommendation::IncreaseReorderingThreshold { .. } => {}
1095 LossRecommendation::EnablePacing { rate } => {
1096 pacing_rate_hint =
1097 Some(pacing_rate_hint.map_or(*rate, |current: u64| current.min(*rate)));
1098 }
1099 LossRecommendation::SwitchCongestionControl { algorithm } => {
1100 prefer_bbr |= algorithm.eq_ignore_ascii_case("bbr");
1101 }
1102 LossRecommendation::EnableFec { rate } => {
1103 if let Some(rate) = bounded_repair_rate(*rate) {
1104 repair_rate_hint =
1105 Some(repair_rate_hint.map_or(rate, |current: f64| current.max(rate)));
1106 loss_pressure = loss_pressure.max(rate);
1107 }
1108 }
1109 }
1110 }
1111
1112 let confidence = finite_unit(self.confidence);
1113 let repair_rate_hint = repair_rate_hint.or_else(|| {
1114 (observed_pressure > 0.0)
1115 .then(|| bounded_repair_rate(observed_pressure * 1.5))
1116 .flatten()
1117 });
1118
1119 TransferLossSignal {
1120 loss_pressure: finite_unit(loss_pressure),
1121 repair_rate_hint,
1122 cwnd_reduction_factor,
1123 pacing_rate_hint,
1124 prefer_bbr,
1125 confidence,
1126 }
1127 }
1128}
1129
1130fn finite_unit(value: f64) -> f64 {
1131 if value.is_finite() {
1132 value.clamp(0.0, 1.0)
1133 } else {
1134 0.0
1135 }
1136}
1137
1138fn bounded_repair_rate(rate: f64) -> Option<f64> {
1139 if rate.is_finite() && rate > 0.0 {
1140 Some(rate.clamp(0.05, 0.30))
1141 } else {
1142 None
1143 }
1144}
1145
1146#[derive(Debug, Clone, Serialize, Deserialize)]
1148pub struct LossAnalysisExport {
1149 pub metrics: LossDetectionMetrics,
1151 pub patterns: Vec<LossPattern>,
1153 pub pattern_confidence: HashMap<LossPattern, f64>,
1155 pub config: LossDetectionConfig,
1157}
1158
1159fn canonical_ack_ranges(ack_ranges: &[AckRange]) -> Vec<CanonicalAckRange> {
1160 let mut ranges: Vec<_> = ack_ranges
1161 .iter()
1162 .map(|range| CanonicalAckRange {
1163 smallest: range.smallest,
1164 largest: range.largest,
1165 })
1166 .collect();
1167 ranges.sort_unstable_by_key(|range| (range.smallest, range.largest));
1168
1169 let mut merged: Vec<CanonicalAckRange> = Vec::with_capacity(ranges.len());
1170 for range in ranges {
1171 if let Some(last) = merged.last_mut() {
1172 if range.smallest <= last.largest.saturating_add(1) {
1173 last.largest = last.largest.max(range.largest);
1174 continue;
1175 }
1176 }
1177 merged.push(range);
1178 }
1179 merged
1180}
1181
1182fn acked_sent_packet_index(
1183 sent_packets: &VecDeque<SentPacketMeta>,
1184 ack_ranges: &[AckRange],
1185) -> HashSet<u64> {
1186 let ranges = canonical_ack_ranges(ack_ranges);
1187 let mut acked_packet_numbers = HashSet::with_capacity(sent_packets.len());
1188
1189 if sent_packets_are_packet_number_ordered(sent_packets) {
1190 let mut range_idx = 0;
1191 for packet in sent_packets {
1192 while let Some(range) = ranges.get(range_idx) {
1193 if packet.packet_number <= range.largest {
1194 break;
1195 }
1196 range_idx += 1;
1197 }
1198
1199 let Some(range) = ranges.get(range_idx) else {
1200 break;
1201 };
1202
1203 if packet.packet_number >= range.smallest {
1204 acked_packet_numbers.insert(packet.packet_number);
1205 }
1206 }
1207 } else {
1208 for packet in sent_packets {
1209 if canonical_ranges_contain_packet(&ranges, packet.packet_number) {
1210 acked_packet_numbers.insert(packet.packet_number);
1211 }
1212 }
1213 }
1214
1215 acked_packet_numbers
1216}
1217
1218fn sent_packets_are_packet_number_ordered(sent_packets: &VecDeque<SentPacketMeta>) -> bool {
1219 let mut previous_packet_number = None;
1220 for packet in sent_packets {
1221 if previous_packet_number.is_some_and(|previous| packet.packet_number < previous) {
1222 return false;
1223 }
1224 previous_packet_number = Some(packet.packet_number);
1225 }
1226 true
1227}
1228
1229fn canonical_ranges_contain_packet(ranges: &[CanonicalAckRange], packet_number: u64) -> bool {
1230 let range_idx = ranges.partition_point(|range| range.largest < packet_number);
1231 ranges
1232 .get(range_idx)
1233 .is_some_and(|range| packet_number >= range.smallest)
1234}
1235
1236#[cfg(test)]
1237mod tests {
1238 use super::*;
1239 use crate::net::quic_native::{
1240 AckRange, PacketNumberSpace, QuicTransportMachine, RttEstimator, SentPacketMeta,
1241 };
1242
1243 fn create_test_packet(space: PacketNumberSpace, pn: u64, time: u64) -> SentPacketMeta {
1244 SentPacketMeta {
1245 space,
1246 packet_number: pn,
1247 bytes: 1200,
1248 ack_eliciting: true,
1249 in_flight: true,
1250 time_sent_micros: time,
1251 }
1252 }
1253
1254 fn test_transport_state(rtt: &RttEstimator) -> LossTransportState {
1255 LossTransportState::from_rtt_and_recovery(rtt, 4_800, 12_000)
1256 }
1257
1258 fn packet_threshold_only_detector() -> AtpLossDetector {
1259 AtpLossDetector::with_config(LossDetectionConfig {
1260 adaptive_packet_threshold: false,
1261 enable_early_retransmit: false,
1262 ..LossDetectionConfig::default()
1263 })
1264 }
1265
1266 #[test]
1267 fn datagram_loss_sample_emits_advisory_recommendations() {
1268 let mut detector = AtpLossDetector::new();
1269
1270 let result = detector.observe_datagram_loss_sample(
1271 100,
1272 10,
1273 Some(Duration::from_millis(25)),
1274 120_000,
1275 240_000,
1276 );
1277
1278 assert_eq!(result.lost_packets.len(), 10);
1279 assert_eq!(result.lost_packet_count, 10);
1280 assert_eq!(result.lost_bytes, 12_000);
1281 let signal = result.transfer_signal(100);
1282 assert_eq!(signal.loss_pressure, 0.1);
1283 assert_eq!(signal.cwnd_reduction_factor, Some(0.5));
1284 assert!(
1285 signal
1286 .repair_rate_hint
1287 .is_some_and(|rate| (rate - 0.15).abs() < f64::EPSILON)
1288 );
1289 assert!(result.recommendations.iter().any(|recommendation| {
1290 matches!(
1291 recommendation,
1292 LossRecommendation::ReduceCongestionWindow { .. }
1293 )
1294 }));
1295 }
1296
1297 #[test]
1298 fn transfer_signal_keeps_congestion_hints_out_of_fec_pressure() {
1299 let result = LossDetectionResult {
1300 lost_packets: Vec::new(),
1301 lost_packet_count: 0,
1302 lost_bytes: 0,
1303 detection_method: LossDetectionMethod::PacketThreshold,
1304 confidence: 0.9,
1305 recommendations: vec![
1306 LossRecommendation::ReduceCongestionWindow { factor: 0.5 },
1307 LossRecommendation::EnablePacing { rate: 75_000 },
1308 LossRecommendation::SwitchCongestionControl {
1309 algorithm: "bbr".to_string(),
1310 },
1311 ],
1312 };
1313
1314 let signal = result.transfer_signal(100);
1315
1316 assert_eq!(signal.loss_pressure, 0.0);
1317 assert_eq!(signal.repair_rate_hint, None);
1318 assert_eq!(signal.cwnd_reduction_factor, Some(0.5));
1319 assert_eq!(signal.pacing_rate_hint, Some(75_000));
1320 assert!(signal.prefer_bbr);
1321 }
1322
1323 #[test]
1324 fn transfer_signal_bounds_explicit_fec_pressure() {
1325 let result = LossDetectionResult {
1326 lost_packets: Vec::new(),
1327 lost_packet_count: 0,
1328 lost_bytes: 0,
1329 detection_method: LossDetectionMethod::PacketThreshold,
1330 confidence: 1.0,
1331 recommendations: vec![
1332 LossRecommendation::EnableFec { rate: 0.01 },
1333 LossRecommendation::EnableFec { rate: 0.9 },
1334 LossRecommendation::EnableFec { rate: f64::NAN },
1335 ],
1336 };
1337
1338 let signal = result.transfer_signal(0);
1339
1340 assert_eq!(signal.loss_pressure, 0.3);
1341 assert_eq!(signal.repair_rate_hint, Some(0.3));
1342 }
1343
1344 #[test]
1345 fn loss_detector_packet_threshold() {
1346 let mut detector = packet_threshold_only_detector();
1347 let rtt = RttEstimator::default();
1348
1349 for pn in 0..6 {
1351 detector.on_packet_sent(create_test_packet(
1352 PacketNumberSpace::ApplicationData,
1353 pn,
1354 pn * 1000,
1355 ));
1356 }
1357
1358 let ack_ranges = [AckRange::new(5, 5).unwrap()];
1360 let result = detector
1361 .on_ack_received(
1362 PacketNumberSpace::ApplicationData,
1363 &ack_ranges,
1364 0,
1365 10_000,
1366 &test_transport_state(&rtt),
1367 )
1368 .expect("Should detect losses");
1369
1370 assert_eq!(result.lost_packets.len(), 3); assert_eq!(
1372 result.detection_method,
1373 LossDetectionMethod::PacketThreshold
1374 );
1375 assert!(
1376 result
1377 .lost_packets
1378 .iter()
1379 .all(|packet| matches!(packet.reason, LossReason::PacketThreshold { .. }))
1380 );
1381 assert_eq!(detector.metrics().time_threshold_losses, 0);
1382 }
1383
1384 #[test]
1385 fn time_threshold_waits_until_boundary_is_reachable() {
1386 let mut detector = packet_threshold_only_detector();
1387 let rtt = RttEstimator::default();
1388
1389 for pn in 0..6 {
1393 detector.on_packet_sent(create_test_packet(
1394 PacketNumberSpace::ApplicationData,
1395 pn,
1396 pn * 1000,
1397 ));
1398 }
1399
1400 let ack_ranges = [AckRange::new(5, 5).unwrap()];
1401 let result = detector
1402 .on_ack_received(
1403 PacketNumberSpace::ApplicationData,
1404 &ack_ranges,
1405 0,
1406 10_000,
1407 &test_transport_state(&rtt),
1408 )
1409 .expect("packet-threshold losses should be detected");
1410
1411 assert_eq!(
1412 result
1413 .lost_packets
1414 .iter()
1415 .map(|packet| (packet.packet_number, packet.reason))
1416 .collect::<Vec<_>>(),
1417 vec![
1418 (0, LossReason::PacketThreshold { threshold: 3 }),
1419 (1, LossReason::PacketThreshold { threshold: 3 }),
1420 (2, LossReason::PacketThreshold { threshold: 3 }),
1421 ]
1422 );
1423 assert_eq!(
1424 result.detection_method,
1425 LossDetectionMethod::PacketThreshold
1426 );
1427 assert_eq!(detector.metrics().time_threshold_losses, 0);
1428 }
1429
1430 #[test]
1431 fn loss_detector_time_threshold() {
1432 let mut detector = AtpLossDetector::new();
1433 let mut rtt = RttEstimator::default();
1434 rtt.update(100_000, 0); detector.on_packet_sent(create_test_packet(PacketNumberSpace::ApplicationData, 0, 0));
1438 detector.on_packet_sent(create_test_packet(
1439 PacketNumberSpace::ApplicationData,
1440 1,
1441 1000,
1442 ));
1443
1444 let ack_ranges = [AckRange::new(1, 1).unwrap()];
1446 let result = detector
1447 .on_ack_received(
1448 PacketNumberSpace::ApplicationData,
1449 &ack_ranges,
1450 0,
1451 200_000, &test_transport_state(&rtt),
1453 )
1454 .expect("Should detect losses");
1455
1456 assert_eq!(result.lost_packets.len(), 1); assert_eq!(result.detection_method, LossDetectionMethod::TimeThreshold);
1458 }
1459
1460 #[test]
1461 fn loss_pattern_analysis() {
1462 let mut detector = AtpLossDetector::new();
1463
1464 for _ in 0..5 {
1466 let rtt = RttEstimator::default();
1467 for pn in 0..10 {
1468 detector.on_packet_sent(create_test_packet(
1469 PacketNumberSpace::ApplicationData,
1470 pn,
1471 pn * 1000,
1472 ));
1473 }
1474
1475 let ack_ranges = [AckRange::new(9, 5).unwrap()];
1477 let _result = detector
1478 .on_ack_received(
1479 PacketNumberSpace::ApplicationData,
1480 &ack_ranges,
1481 0,
1482 50_000,
1483 &test_transport_state(&rtt),
1484 )
1485 .unwrap();
1486 }
1487
1488 detector.analyze_loss_patterns();
1490 assert!(
1491 detector
1492 .pattern_analyzer
1493 .patterns
1494 .contains(&LossPattern::Burst)
1495 );
1496 }
1497
1498 #[test]
1499 fn reordering_detection() {
1500 let mut tracker = ReorderingTracker::new();
1501 let initial_threshold = tracker.current_threshold;
1502
1503 tracker.adapt_threshold();
1505
1506 assert!(tracker.current_threshold > initial_threshold);
1507 }
1508
1509 #[test]
1510 fn loss_pattern_analysis_records_transport_recovery_state() {
1511 let mut detector = AtpLossDetector::new();
1512 let mut rtt = RttEstimator::default();
1513 rtt.update(100_000, 0);
1514 let transport_state = LossTransportState::from_rtt_and_recovery(&rtt, 6_000, 24_000);
1515
1516 for pn in 0..6 {
1517 detector.on_packet_sent(create_test_packet(
1518 PacketNumberSpace::ApplicationData,
1519 pn,
1520 pn * 1000,
1521 ));
1522 }
1523
1524 let ack_ranges = [AckRange::new(5, 5).unwrap()];
1525 let result = detector
1526 .on_ack_received(
1527 PacketNumberSpace::ApplicationData,
1528 &ack_ranges,
1529 0,
1530 10_000,
1531 &transport_state,
1532 )
1533 .expect("Should detect losses");
1534
1535 assert!(!result.lost_packets.is_empty());
1536 let event = detector
1537 .pattern_analyzer
1538 .loss_events
1539 .back()
1540 .expect("loss event recorded");
1541 assert_eq!(event.conditions.rtt_micros, Some(100_000));
1542 assert_eq!(event.conditions.rttvar_micros, Some(50_000));
1543 assert_eq!(event.conditions.bytes_in_flight, 6_000);
1544 assert_eq!(event.conditions.congestion_window, 24_000);
1545 }
1546
1547 #[test]
1548 fn transport_state_reads_native_quic_recovery_counters() {
1549 let mut transport = QuicTransportMachine::new();
1550 transport.on_packet_sent(create_test_packet(
1551 PacketNumberSpace::ApplicationData,
1552 0,
1553 10_000,
1554 ));
1555 transport.on_packet_sent(create_test_packet(
1556 PacketNumberSpace::ApplicationData,
1557 1,
1558 20_000,
1559 ));
1560
1561 let initial_state = LossTransportState::from_transport(&transport);
1562 assert_eq!(initial_state.bytes_in_flight, 2_400);
1563 assert_eq!(
1564 initial_state.congestion_window,
1565 transport.congestion_window_bytes()
1566 );
1567 assert_eq!(initial_state.latest_rtt_micros, None);
1568
1569 let _ack = transport.on_ack_received(PacketNumberSpace::ApplicationData, &[1], 0, 50_000);
1570 let acked_state = LossTransportState::from_transport(&transport);
1571 assert_eq!(acked_state.bytes_in_flight, 1_200);
1572 assert_eq!(acked_state.latest_rtt_micros, Some(30_000));
1573 assert_eq!(acked_state.smoothed_rtt_micros, Some(30_000));
1574 assert_eq!(acked_state.rttvar_micros, Some(15_000));
1575 }
1576
1577 #[test]
1578 fn ack_matching_canonicalizes_ranges_before_hash_lookup() {
1579 let mut detector = AtpLossDetector::new();
1580 let rtt = RttEstimator::default();
1581
1582 for pn in 0..13 {
1583 detector.on_packet_sent(create_test_packet(
1584 PacketNumberSpace::ApplicationData,
1585 pn,
1586 pn * 1000,
1587 ));
1588 }
1589
1590 let ack_ranges = [
1591 AckRange::new(9, 7).unwrap(),
1592 AckRange::new(3, 1).unwrap(),
1593 AckRange::new(8, 5).unwrap(),
1594 ];
1595 let result = detector
1596 .on_ack_received(
1597 PacketNumberSpace::ApplicationData,
1598 &ack_ranges,
1599 0,
1600 10_000,
1601 &test_transport_state(&rtt),
1602 )
1603 .expect("ACK ranges should be processed");
1604
1605 assert_eq!(
1606 result
1607 .lost_packets
1608 .iter()
1609 .map(|packet| packet.packet_number)
1610 .collect::<Vec<_>>(),
1611 vec![0, 4]
1612 );
1613 assert_eq!(
1614 detector.spaces[PacketNumberSpace::ApplicationData as usize]
1615 .sent_packets
1616 .iter()
1617 .map(|packet| packet.packet_number)
1618 .collect::<Vec<_>>(),
1619 vec![10, 11, 12]
1620 );
1621 }
1622
1623 #[test]
1624 fn ack_matching_ignores_unsent_packet_numbers() {
1625 let mut detector = AtpLossDetector::new();
1626 let rtt = RttEstimator::default();
1627
1628 for pn in 0..5 {
1629 detector.on_packet_sent(create_test_packet(
1630 PacketNumberSpace::ApplicationData,
1631 pn,
1632 pn * 1000,
1633 ));
1634 }
1635
1636 let ack_ranges = [AckRange::new(1_000_000, 1_000_000).unwrap()];
1637 let result = detector
1638 .on_ack_received(
1639 PacketNumberSpace::ApplicationData,
1640 &ack_ranges,
1641 0,
1642 10_000,
1643 &test_transport_state(&rtt),
1644 )
1645 .expect("Unsent ACK should not fail");
1646
1647 assert!(result.lost_packets.is_empty());
1648 assert_eq!(
1649 detector.spaces[PacketNumberSpace::ApplicationData as usize]
1650 .sent_packets
1651 .iter()
1652 .map(|packet| packet.packet_number)
1653 .collect::<Vec<_>>(),
1654 vec![0, 1, 2, 3, 4]
1655 );
1656 }
1657
1658 #[test]
1659 fn early_retransmit_threshold_does_not_overflow_at_max_packet_number() {
1660 let mut detector = AtpLossDetector::new();
1661 let rtt = RttEstimator::default();
1662
1663 detector.on_packet_sent(create_test_packet(
1664 PacketNumberSpace::ApplicationData,
1665 u64::MAX - 1,
1666 1_000,
1667 ));
1668 detector.on_packet_sent(create_test_packet(
1669 PacketNumberSpace::ApplicationData,
1670 u64::MAX,
1671 1_100,
1672 ));
1673
1674 let ack_ranges = [AckRange::new(u64::MAX - 1, u64::MAX - 1).unwrap()];
1675 let result = detector
1676 .on_ack_received(
1677 PacketNumberSpace::ApplicationData,
1678 &ack_ranges,
1679 0,
1680 2_000,
1681 &test_transport_state(&rtt),
1682 )
1683 .expect("overflow-edge ACK should be processed");
1684
1685 assert!(result.lost_packets.is_empty());
1686 assert_eq!(
1687 detector.spaces[PacketNumberSpace::ApplicationData as usize]
1688 .sent_packets
1689 .iter()
1690 .map(|packet| packet.packet_number)
1691 .collect::<Vec<_>>(),
1692 vec![u64::MAX]
1693 );
1694 }
1695
1696 #[test]
1697 fn late_ack_of_declared_lost_packet_records_one_spurious_loss() {
1698 let mut detector = packet_threshold_only_detector();
1699 let rtt = RttEstimator::default();
1700
1701 for pn in 0..6 {
1702 detector.on_packet_sent(create_test_packet(
1703 PacketNumberSpace::ApplicationData,
1704 pn,
1705 pn * 1000,
1706 ));
1707 }
1708
1709 let initial_ack = [AckRange::new(5, 5).unwrap()];
1710 let loss = detector
1711 .on_ack_received(
1712 PacketNumberSpace::ApplicationData,
1713 &initial_ack,
1714 0,
1715 10_000,
1716 &test_transport_state(&rtt),
1717 )
1718 .expect("initial ACK should detect losses");
1719 assert_eq!(
1720 loss.lost_packets
1721 .iter()
1722 .map(|packet| packet.packet_number)
1723 .collect::<Vec<_>>(),
1724 vec![0, 1, 2]
1725 );
1726 assert_eq!(detector.metrics().false_losses, 0);
1727
1728 let late_ack = [AckRange::new(0, 0).unwrap()];
1729 let late_result = detector
1730 .on_ack_received(
1731 PacketNumberSpace::ApplicationData,
1732 &late_ack,
1733 0,
1734 11_000,
1735 &test_transport_state(&rtt),
1736 )
1737 .expect("late ACK should be processed");
1738 assert!(late_result.lost_packets.is_empty());
1739 assert_eq!(detector.metrics().false_losses, 1);
1740 assert!(
1741 detector
1742 .pattern_analyzer
1743 .patterns
1744 .contains(&LossPattern::Reordering)
1745 );
1746
1747 let duplicate_late_result = detector
1748 .on_ack_received(
1749 PacketNumberSpace::ApplicationData,
1750 &late_ack,
1751 0,
1752 12_000,
1753 &test_transport_state(&rtt),
1754 )
1755 .expect("duplicate late ACK should be processed");
1756 assert!(duplicate_late_result.lost_packets.is_empty());
1757 assert_eq!(detector.metrics().false_losses, 1);
1758 }
1759}