Skip to main content

asupersync/net/atp/loss/
detector.rs

1//! ATP Loss Detection Algorithms
2//!
3//! Advanced loss detection for ATP with improved accuracy and
4//! integration with transfer decision-making.
5
6#![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
16/// Helper macro to extract success value from AtpOutcome or early return with error.
17macro_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
28/// ATP-enhanced loss detector with adaptive algorithms.
29pub struct AtpLossDetector {
30    /// Per-space loss detection state.
31    spaces: [SpaceLossState; 3],
32    /// Global loss detection configuration.
33    config: LossDetectionConfig,
34    /// Loss pattern analyzer.
35    pattern_analyzer: LossPatternAnalyzer,
36    /// Reordering tolerance tracker.
37    reordering_tracker: ReorderingTracker,
38    /// Lost packet numbers already counted as spurious, keyed by packet number space.
39    spurious_loss_packets: HashSet<(usize, u64)>,
40    /// Detection metrics for analysis.
41    metrics: LossDetectionMetrics,
42}
43
44/// Loss detection state for a single packet number space.
45#[derive(Debug, Clone)]
46struct SpaceLossState {
47    /// Sent packets awaiting acknowledgment.
48    sent_packets: VecDeque<SentPacketMeta>,
49    /// Largest acknowledged packet number.
50    largest_acked: Option<u64>,
51    /// Time of largest acked packet.
52    largest_acked_time: Option<u64>,
53    /// Loss detection timer deadline.
54    loss_timer_deadline: Option<u64>,
55    /// Early retransmit timer deadline.
56    early_retransmit_deadline: Option<u64>,
57}
58
59/// Loss detection configuration.
60#[derive(Debug, Clone, Serialize, Deserialize)]
61pub struct LossDetectionConfig {
62    /// Packet threshold for declaring loss (default: 3).
63    pub packet_threshold: u32,
64    /// Time threshold multiplier (default: 9/8).
65    pub time_threshold_multiplier: f64,
66    /// Minimum time threshold in microseconds.
67    pub min_time_threshold_micros: u64,
68    /// Maximum reordering threshold.
69    pub max_reordering_threshold: u32,
70    /// Enable adaptive packet threshold.
71    pub adaptive_packet_threshold: bool,
72    /// Enable early retransmit.
73    pub enable_early_retransmit: bool,
74    /// Early retransmit threshold.
75    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, // 1ms
84            max_reordering_threshold: 10,
85            adaptive_packet_threshold: true,
86            enable_early_retransmit: true,
87            early_retransmit_threshold: 1,
88        }
89    }
90}
91
92/// Loss pattern analysis for adaptive behavior.
93#[derive(Debug, Clone)]
94struct LossPatternAnalyzer {
95    /// Recent loss events.
96    loss_events: VecDeque<LossEvent>,
97    /// Detected loss patterns.
98    patterns: Vec<LossPattern>,
99    /// Pattern confidence scores.
100    pattern_confidence: HashMap<LossPattern, f64>,
101}
102
103/// Loss event for pattern analysis.
104#[derive(Debug, Clone)]
105struct LossEvent {
106    /// Timestamp of loss detection.
107    timestamp: Instant,
108    /// Lost packet numbers.
109    lost_packets: Vec<u64>,
110    /// Detection method used.
111    detection_method: LossDetectionMethod,
112    /// Network conditions at time of loss.
113    conditions: NetworkConditions,
114}
115
116/// Network conditions snapshot.
117#[derive(Debug, Clone)]
118struct NetworkConditions {
119    /// RTT at time of loss.
120    rtt_micros: Option<u64>,
121    /// RTT variance.
122    rttvar_micros: Option<u64>,
123    /// Bytes in flight.
124    bytes_in_flight: u64,
125    /// Congestion window.
126    congestion_window: u64,
127}
128
129#[derive(Debug, Clone, Copy, PartialEq, Eq)]
130struct CanonicalAckRange {
131    smallest: u64,
132    largest: u64,
133}
134
135/// Transport recovery state used by ATP loss analysis.
136///
137/// The detector keeps its own sent-packet view for ATP decisions, but RTT and
138/// congestion context must come from the live transport recovery state so loss
139/// classification sees the same network conditions as QUIC recovery.
140#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
141pub struct LossTransportState {
142    /// Latest RTT sample.
143    pub latest_rtt_micros: Option<u64>,
144    /// Smoothed RTT estimate.
145    pub smoothed_rtt_micros: Option<u64>,
146    /// RTT variance estimate.
147    pub rttvar_micros: Option<u64>,
148    /// Bytes currently in flight according to transport recovery.
149    pub bytes_in_flight: u64,
150    /// Current congestion window in bytes.
151    pub congestion_window: u64,
152}
153
154impl LossTransportState {
155    /// Build a snapshot from the native QUIC transport machine.
156    #[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    /// Build a snapshot from explicit recovery counters and RTT estimator.
166    #[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/// Detected loss patterns.
189#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
190pub enum LossPattern {
191    /// Random sporadic losses.
192    Sporadic,
193    /// Burst losses (multiple consecutive packets).
194    Burst,
195    /// Periodic losses (pattern of losses).
196    Periodic,
197    /// Reordering-induced false losses.
198    Reordering,
199    /// Congestion-induced losses.
200    Congestion,
201    /// Tail losses (end of flight).
202    Tail,
203}
204
205/// Loss detection methods.
206#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
207pub enum LossDetectionMethod {
208    /// Packet threshold exceeded.
209    PacketThreshold,
210    /// Time threshold exceeded.
211    TimeThreshold,
212    /// Early retransmit.
213    EarlyRetransmit,
214    /// Both packet and time thresholds.
215    Combined,
216}
217
218/// Reordering tolerance tracking.
219#[derive(Debug, Clone)]
220struct ReorderingTracker {
221    /// Recent reordering measurements.
222    reordering_measurements: VecDeque<u32>,
223    /// Current reordering threshold.
224    current_threshold: u32,
225    /// Maximum observed reordering.
226    max_reordering: u32,
227    /// Reordering adaptation factor.
228    adaptation_factor: f64,
229}
230
231/// Loss detection metrics.
232#[derive(Debug, Clone, Serialize, Deserialize)]
233pub struct LossDetectionMetrics {
234    /// Total packets declared lost.
235    pub total_lost_packets: u64,
236    /// Packets lost by packet threshold.
237    pub packet_threshold_losses: u64,
238    /// Packets lost by time threshold.
239    pub time_threshold_losses: u64,
240    /// False loss declarations (spurious retransmits).
241    pub false_losses: u64,
242    /// Average packet threshold used.
243    pub avg_packet_threshold: f64,
244    /// Average time threshold used.
245    pub avg_time_threshold_micros: f64,
246    /// Reordering events detected.
247    pub reordering_events: u64,
248    /// Pattern detection accuracy.
249    pub pattern_accuracy: f64,
250}
251
252/// Loss detection result.
253#[derive(Debug, Clone)]
254pub struct LossDetectionResult {
255    /// Newly detected lost packets.
256    pub lost_packets: Vec<LostPacketInfo>,
257    /// Total number of packets/datagrams represented by this result.
258    ///
259    /// `lost_packets` may be sample-capped for bounded memory when the input is
260    /// aggregate datagram loss evidence; this count preserves the actual loss
261    /// pressure for transfer-brain decisions.
262    pub lost_packet_count: u64,
263    /// Total lost bytes.
264    pub lost_bytes: u64,
265    /// Detection method used.
266    pub detection_method: LossDetectionMethod,
267    /// Confidence in the detection (0.0 - 1.0).
268    pub confidence: f64,
269    /// Recommended actions.
270    pub recommendations: Vec<LossRecommendation>,
271}
272
273/// Information about a lost packet.
274#[derive(Debug, Clone, Serialize, Deserialize)]
275pub struct LostPacketInfo {
276    /// Packet number.
277    pub packet_number: u64,
278    /// Packet size in bytes.
279    pub bytes: u64,
280    /// Time when packet was sent.
281    pub sent_time_micros: u64,
282    /// Time when loss was detected.
283    pub detected_time_micros: u64,
284    /// Reason for declaring loss.
285    pub reason: LossReason,
286}
287
288/// Reason for packet loss declaration.
289#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
290pub enum LossReason {
291    /// Packet threshold exceeded (N packets acked beyond this one).
292    PacketThreshold { threshold: u32 },
293    /// Time threshold exceeded (too much time elapsed).
294    TimeThreshold { threshold_micros: u64 },
295    /// Both thresholds exceeded.
296    BothThresholds {
297        packet_threshold: u32,
298        time_threshold_micros: u64,
299    },
300    /// Early retransmit triggered.
301    EarlyRetransmit,
302}
303
304/// Loss-based recommendations.
305#[derive(Debug, Clone, Serialize, Deserialize)]
306pub enum LossRecommendation {
307    /// Reduce congestion window.
308    ReduceCongestionWindow { factor: f64 },
309    /// Increase reordering threshold.
310    IncreaseReorderingThreshold { new_threshold: u32 },
311    /// Enable pacing.
312    EnablePacing { rate: u64 },
313    /// Switch to different congestion control.
314    SwitchCongestionControl { algorithm: String },
315    /// Enable forward error correction.
316    EnableFec { rate: f64 },
317}
318
319/// Compact loss pressure signal consumed by transfer scheduling.
320#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
321pub struct TransferLossSignal {
322    /// Estimated loss pressure in the interval, clamped to `[0.0, 1.0]`.
323    pub loss_pressure: f64,
324    /// Optional FEC/repair rate hint derived from detector recommendations.
325    pub repair_rate_hint: Option<f64>,
326    /// Optional congestion-window reduction factor.
327    pub cwnd_reduction_factor: Option<f64>,
328    /// Optional pacing rate hint in bytes per second.
329    pub pacing_rate_hint: Option<u64>,
330    /// Whether the detector recommended moving to BBR-like control.
331    pub prefer_bbr: bool,
332    /// Detector confidence for this signal, clamped to `[0.0, 1.0]`.
333    pub confidence: f64,
334}
335
336impl AtpLossDetector {
337    /// Create a new ATP loss detector.
338    #[must_use]
339    pub fn new() -> Self {
340        Self::with_config(LossDetectionConfig::default())
341    }
342
343    /// Create with custom configuration.
344    #[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    /// Track sent packet.
361    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        // Limit memory usage
366        if self.spaces[space_idx].sent_packets.len() > 10_000 {
367            self.spaces[space_idx].sent_packets.pop_front();
368        }
369    }
370
371    /// Process acknowledgment and detect losses.
372    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        // Find newly acknowledged packets
384        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        // Process acknowledgments
389        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        // Update largest acked
400        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        // Detect losses
408        let loss_result = try_outcome!(self.detect_losses(space, now_micros, transport_state));
409
410        // Update pattern analysis
411        if !loss_result.lost_packets.is_empty() {
412            self.update_pattern_analysis(&loss_result, transport_state);
413        }
414
415        // Update reordering tracking. Newly acked packets still in the sent
416        // queue catch ordinary reordering; ACK ranges catch packets already
417        // removed after a previous loss declaration.
418        self.update_reordering_tracking(space_idx, &newly_acked, ack_ranges, &loss_result);
419
420        AtpOutcome::ok(loss_result)
421    }
422
423    /// Feed aggregate datagram-round loss evidence into the ATP loss analyzer.
424    ///
425    /// RaptorQ-style transports receive per-round `NeedMore` pressure instead of
426    /// QUIC ACK ranges. This bridge lets them reuse the same pattern analyzer and
427    /// recommendation logic without synthesizing full packet histories or changing
428    /// their wire protocol.
429    #[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    /// Detect losses in a packet number space.
510    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        // Calculate thresholds
526        let packet_threshold = self.get_adaptive_packet_threshold(space);
527        let time_threshold = self.calculate_time_threshold(*transport_state);
528
529        // Check for time threshold losses only after enough time has elapsed.
530        // A saturating subtraction would floor the boundary to zero and mark a
531        // packet sent at t=0 as time-lost before the threshold is reachable.
532        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            // Packet threshold loss
543            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            // Time threshold loss
557            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                    // Both thresholds
562                    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            // Early retransmit
579            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                    // Edge case: packet marked as lost but no specific reason set
603                    // Fall back to time threshold as a safe default
604                    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        // Determine primary detection method
626        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        // Calculate confidence
639        let confidence = self.calculate_detection_confidence(&lost_packets, detection_method);
640
641        // Generate recommendations
642        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        // Use reordering tracker to adapt threshold
660        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        // Early retransmit if only one packet ahead is acked
672        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        // Base confidence by method
688        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        // Adjust based on pattern analysis
696        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            // Many losses suggest congestion
721            recommendations.push(LossRecommendation::ReduceCongestionWindow { factor: 0.5 });
722        }
723
724        if method == LossDetectionMethod::EarlyRetransmit {
725            // Early retransmit might indicate reordering
726            recommendations.push(LossRecommendation::IncreaseReorderingThreshold {
727                new_threshold: self.reordering_tracker.current_threshold.saturating_add(1),
728            });
729        }
730
731        // Check loss patterns
732        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    /// Get current metrics.
947    #[must_use]
948    pub fn metrics(&self) -> &LossDetectionMetrics {
949        &self.metrics
950    }
951
952    /// Export detection log for analysis.
953    #[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, // Start with default
997            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    /// Convert this result into transfer-brain loss pressure.
1067    ///
1068    /// `sent_units` is the packet/datagram denominator for the detector window.
1069    /// A zero denominator deliberately yields zero pressure while still carrying
1070    /// recommendation hints, so fail-closed callers can avoid divide-by-zero
1071    /// amplification.
1072    #[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/// Loss analysis export for external tools.
1147#[derive(Debug, Clone, Serialize, Deserialize)]
1148pub struct LossAnalysisExport {
1149    /// Current metrics.
1150    pub metrics: LossDetectionMetrics,
1151    /// Detected patterns.
1152    pub patterns: Vec<LossPattern>,
1153    /// Pattern confidence scores.
1154    pub pattern_confidence: HashMap<LossPattern, f64>,
1155    /// Current configuration.
1156    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        // Send packets 0-5
1350        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        // ACK packet 5 (should cause 0, 1, 2 to be declared lost via packet threshold)
1359        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); // Packets 0, 1, 2 lost
1371        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        // Default base RTT is 333ms, so the 9/8 time threshold is far beyond
1390        // this 10ms ACK. Packet 0 has timestamp zero and must not be upgraded
1391        // to a combined packet+time loss by a saturated boundary of zero.
1392        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); // 100ms RTT
1435
1436        // Send packets with significant time gaps
1437        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        // ACK packet 1 much later (should cause packet 0 to be lost via time threshold)
1445        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, // 200ms later
1452                &test_transport_state(&rtt),
1453            )
1454            .expect("Should detect losses");
1455
1456        assert_eq!(result.lost_packets.len(), 1); // Packet 0 lost
1457        assert_eq!(result.detection_method, LossDetectionMethod::TimeThreshold);
1458    }
1459
1460    #[test]
1461    fn loss_pattern_analysis() {
1462        let mut detector = AtpLossDetector::new();
1463
1464        // Drive a burst-loss packet pattern.
1465        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            // Lose packets 0-4 (burst)
1476            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        // Should detect burst pattern
1489        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        // Drive reordering-threshold adaptation.
1504        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}