Skip to main content

quiche/recovery/
mod.rs

1// Copyright (C) 2018-2019, Cloudflare, Inc.
2// All rights reserved.
3//
4// Redistribution and use in source and binary forms, with or without
5// modification, are permitted provided that the following conditions are
6// met:
7//
8//     * Redistributions of source code must retain the above copyright notice,
9//       this list of conditions and the following disclaimer.
10//
11//     * Redistributions in binary form must reproduce the above copyright
12//       notice, this list of conditions and the following disclaimer in the
13//       documentation and/or other materials provided with the distribution.
14//
15// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS
16// IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO,
17// THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
18// PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
19// CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
20// EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
21// PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
22// PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
23// LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
24// NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
25// SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
26
27use std::str::FromStr;
28use std::time::Duration;
29use std::time::Instant;
30
31use crate::frame;
32use crate::packet;
33use crate::ranges::RangeSet;
34pub(crate) use crate::recovery::bandwidth::Bandwidth;
35use crate::Config;
36use crate::Result;
37
38#[cfg(feature = "qlog")]
39use qlog::events::EventData;
40#[cfg(feature = "qlog")]
41use serde::Serialize;
42
43use smallvec::SmallVec;
44
45use self::congestion::recovery::LegacyRecovery;
46use self::gcongestion::GRecovery;
47pub use gcongestion::BbrBwLoReductionStrategy;
48pub use gcongestion::BbrParams;
49#[cfg(feature = "internal")]
50pub use gcongestion::BbrRttJumpDetector;
51
52// Loss Recovery
53const INITIAL_PACKET_THRESHOLD: u64 = 3;
54
55const MAX_PACKET_THRESHOLD: u64 = 20;
56
57// Time threshold used to calculate the loss time.
58//
59// https://www.rfc-editor.org/rfc/rfc9002.html#section-6.1.2
60const INITIAL_TIME_THRESHOLD: f64 = 9.0 / 8.0;
61
62// Reduce the sensitivity to packet reordering after the first reordering event.
63//
64// Packet reorder is not a real loss event so quickly reduce the sensitivity to
65// avoid penializing subsequent packet reordering.
66//
67// https://www.rfc-editor.org/rfc/rfc9002.html#section-6.1.2
68//
69// Implementations MAY experiment with absolute thresholds, thresholds from
70// previous connections, adaptive thresholds, or the including of RTT variation.
71// Smaller thresholds reduce reordering resilience and increase spurious
72// retransmissions, and larger thresholds increase loss detection delay.
73const PACKET_REORDER_TIME_THRESHOLD: f64 = 5.0 / 4.0;
74
75// # Experiment: enable_relaxed_loss_threshold
76//
77// Time threshold overhead used to calculate the loss time.
78//
79// The actual threshold is calcualted as 1 + INITIAL_TIME_THRESHOLD_OVERHEAD and
80// equivalent to INITIAL_TIME_THRESHOLD.
81const INITIAL_TIME_THRESHOLD_OVERHEAD: f64 = 1.0 / 8.0;
82// # Experiment: enable_relaxed_loss_threshold
83//
84// The factor by which to increase the time threshold on spurious loss.
85const TIME_THRESHOLD_OVERHEAD_MULTIPLIER: f64 = 2.0;
86
87const GRANULARITY: Duration = Duration::from_millis(1);
88
89const MAX_PTO_PROBES_COUNT: usize = 2;
90
91const MINIMUM_WINDOW_PACKETS: usize = 2;
92
93const LOSS_REDUCTION_FACTOR: f64 = 0.5;
94
95// How many non ACK eliciting packets we send before including a PING to solicit
96// an ACK.
97pub(super) const MAX_OUTSTANDING_NON_ACK_ELICITING: usize = 24;
98
99#[derive(Default)]
100struct LossDetectionTimer {
101    time: Option<Instant>,
102}
103
104impl LossDetectionTimer {
105    fn update(&mut self, timeout: Instant) {
106        self.time = Some(timeout);
107    }
108
109    fn clear(&mut self) {
110        self.time = None;
111    }
112}
113
114impl std::fmt::Debug for LossDetectionTimer {
115    fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
116        match self.time {
117            Some(v) => {
118                let now = Instant::now();
119                if v > now {
120                    let d = v.duration_since(now);
121                    write!(f, "{d:?}")
122                } else {
123                    write!(f, "exp")
124                }
125            },
126            None => write!(f, "none"),
127        }
128    }
129}
130
131#[derive(Clone, Copy, PartialEq)]
132pub struct RecoveryConfig {
133    pub initial_rtt: Duration,
134    pub max_send_udp_payload_size: usize,
135    pub max_ack_delay: Duration,
136    pub cc_algorithm: CongestionControlAlgorithm,
137    pub custom_bbr_params: Option<BbrParams>,
138    pub hystart: bool,
139    pub pacing: bool,
140    pub max_pacing_rate: Option<u64>,
141    pub initial_congestion_window_packets: usize,
142    pub enable_relaxed_loss_threshold: bool,
143    pub enable_cubic_idle_restart_fix: bool,
144}
145
146impl RecoveryConfig {
147    pub fn from_config(config: &Config) -> Self {
148        Self {
149            initial_rtt: config.initial_rtt,
150            max_send_udp_payload_size: config.max_send_udp_payload_size,
151            max_ack_delay: Duration::ZERO,
152            cc_algorithm: config.cc_algorithm,
153            custom_bbr_params: config.custom_bbr_params,
154            hystart: config.hystart,
155            pacing: config.pacing,
156            max_pacing_rate: config.max_pacing_rate,
157            initial_congestion_window_packets: config
158                .initial_congestion_window_packets,
159            enable_relaxed_loss_threshold: config.enable_relaxed_loss_threshold,
160            enable_cubic_idle_restart_fix: config.enable_cubic_idle_restart_fix,
161        }
162    }
163}
164
165#[enum_dispatch::enum_dispatch(RecoveryOps)]
166#[allow(clippy::large_enum_variant)]
167#[derive(Debug)]
168pub(crate) enum Recovery {
169    Legacy(LegacyRecovery),
170    GCongestion(GRecovery),
171}
172
173#[derive(Debug, Default, PartialEq)]
174pub struct OnAckReceivedOutcome {
175    pub lost_packets: usize,
176    pub lost_bytes: usize,
177    pub acked_bytes: usize,
178    pub spurious_losses: usize,
179}
180
181#[derive(Debug, Default)]
182pub struct OnLossDetectionTimeoutOutcome {
183    pub lost_packets: usize,
184    pub lost_bytes: usize,
185}
186
187#[enum_dispatch::enum_dispatch]
188/// Api for the Recovery implementation
189pub trait RecoveryOps {
190    fn lost_count(&self) -> usize;
191    fn bytes_lost(&self) -> u64;
192
193    /// Returns whether or not we should elicit an ACK even if we wouldn't
194    /// otherwise have constructed an ACK eliciting packet.
195    fn should_elicit_ack(&self, epoch: packet::Epoch) -> bool;
196
197    fn next_acked_frame(&mut self, epoch: packet::Epoch) -> Option<frame::Frame>;
198
199    fn next_lost_frame(&mut self, epoch: packet::Epoch) -> Option<frame::Frame>;
200
201    fn get_largest_acked_on_epoch(&self, epoch: packet::Epoch) -> Option<u64>;
202    fn has_lost_frames(&self, epoch: packet::Epoch) -> bool;
203    fn loss_probes(&self, epoch: packet::Epoch) -> usize;
204    #[cfg(test)]
205    fn inc_loss_probes(&mut self, epoch: packet::Epoch);
206    #[cfg(test)]
207    fn lost_frames_count(&self, epoch: packet::Epoch) -> usize;
208
209    fn ping_sent(&mut self, epoch: packet::Epoch);
210
211    fn on_packet_sent(
212        &mut self, pkt: Sent, epoch: packet::Epoch,
213        handshake_status: HandshakeStatus, now: Instant, trace_id: &str,
214    );
215    fn get_packet_send_time(&self, now: Instant) -> Instant;
216
217    #[allow(clippy::too_many_arguments)]
218    fn on_ack_received(
219        &mut self, ranges: &RangeSet, ack_delay: u64, epoch: packet::Epoch,
220        handshake_status: HandshakeStatus, now: Instant, skip_pn: Option<u64>,
221        trace_id: &str,
222    ) -> Result<OnAckReceivedOutcome>;
223
224    fn on_loss_detection_timeout(
225        &mut self, handshake_status: HandshakeStatus, now: Instant,
226        trace_id: &str,
227    ) -> OnLossDetectionTimeoutOutcome;
228    fn on_pkt_num_space_discarded(
229        &mut self, epoch: packet::Epoch, handshake_status: HandshakeStatus,
230        now: Instant,
231    );
232    fn on_path_change(
233        &mut self, epoch: packet::Epoch, now: Instant, _trace_id: &str,
234    ) -> (usize, usize);
235    fn loss_detection_timer(&self) -> Option<Instant>;
236    fn cwnd(&self) -> usize;
237    fn cwnd_available(&self) -> usize;
238    fn rtt(&self) -> Duration;
239
240    fn min_rtt(&self) -> Option<Duration>;
241
242    fn max_rtt(&self) -> Option<Duration>;
243
244    fn rttvar(&self) -> Duration;
245
246    fn pto(&self) -> Duration;
247
248    /// The most recent data delivery rate estimate.
249    fn delivery_rate(&self) -> Bandwidth;
250
251    /// Maximum bandwidth estimate, if one is available.
252    fn max_bandwidth(&self) -> Option<Bandwidth>;
253
254    /// Total number of confirmed persistent RTT jump episodes.
255    fn rtt_persistent_jump_count(&self) -> u64;
256
257    /// Statistics from when a CCA first exited the startup phase.
258    fn startup_exit(&self) -> Option<StartupExit>;
259
260    fn max_datagram_size(&self) -> usize;
261
262    fn pmtud_update_max_datagram_size(&mut self, new_max_datagram_size: usize);
263
264    fn update_max_datagram_size(&mut self, new_max_datagram_size: usize);
265
266    fn on_app_limited(&mut self);
267
268    // Since a recovery module is path specific, this tracks the largest packet
269    // sent per path.
270    #[cfg(test)]
271    fn largest_sent_pkt_num_on_path(&self, epoch: packet::Epoch) -> Option<u64>;
272
273    #[cfg(any(test, feature = "qlog"))]
274    fn app_limited(&self) -> bool;
275
276    #[cfg(test)]
277    fn sent_packets_len(&self, epoch: packet::Epoch) -> usize;
278
279    fn bytes_in_flight(&self) -> usize;
280
281    fn bytes_in_flight_duration(&self) -> Duration;
282
283    #[cfg(test)]
284    fn in_flight_count(&self, epoch: packet::Epoch) -> usize;
285
286    #[cfg(test)]
287    fn pacing_rate(&self) -> u64;
288
289    #[cfg(test)]
290    fn pto_count(&self) -> u32;
291
292    // This value might be `None` when the `enable_relaxed_loss_threshold`
293    // experiment is enabled for gcongestion.
294    #[cfg(test)]
295    fn pkt_thresh(&self) -> Option<u64>;
296
297    #[cfg(test)]
298    fn time_thresh(&self) -> f64;
299
300    #[cfg(test)]
301    fn lost_spurious_count(&self) -> usize;
302
303    #[cfg(test)]
304    fn detect_lost_packets_for_test(
305        &mut self, epoch: packet::Epoch, now: Instant,
306    ) -> (usize, usize);
307
308    fn update_app_limited(&mut self, v: bool);
309
310    fn delivery_rate_update_app_limited(&mut self, v: bool);
311
312    fn update_max_ack_delay(&mut self, max_ack_delay: Duration);
313
314    #[cfg(feature = "qlog")]
315    fn state_str(&self, now: Instant) -> &'static str;
316
317    #[cfg(feature = "qlog")]
318    fn get_updated_qlog_event_data(&mut self) -> Option<EventData>;
319
320    #[cfg(feature = "qlog")]
321    fn get_updated_qlog_cc_state(&mut self, now: Instant)
322        -> Option<&'static str>;
323
324    fn send_quantum(&self) -> usize;
325
326    fn get_next_release_time(&self) -> ReleaseDecision;
327
328    fn gcongestion_enabled(&self) -> bool;
329}
330
331impl Recovery {
332    pub fn new_with_config(recovery_config: &RecoveryConfig) -> Self {
333        let grecovery = GRecovery::new(recovery_config);
334        if let Some(grecovery) = grecovery {
335            Recovery::from(grecovery)
336        } else {
337            Recovery::from(LegacyRecovery::new_with_config(recovery_config))
338        }
339    }
340
341    #[cfg(feature = "qlog")]
342    pub fn maybe_qlog(
343        &mut self, qlog: &mut qlog::streamer::QlogStreamer, now: Instant,
344    ) {
345        if let Some(ev_data) = self.get_updated_qlog_event_data() {
346            qlog.add_event_data_with_instant(ev_data, now).ok();
347        }
348
349        if let Some(cc_state) = self.get_updated_qlog_cc_state(now) {
350            let ev_data = EventData::QuicCongestionStateUpdated(
351                qlog::events::quic::CongestionStateUpdated {
352                    old: None,
353                    new: cc_state.to_string(),
354                    trigger: None,
355                },
356            );
357
358            qlog.add_event_data_with_instant(ev_data, now).ok();
359        }
360    }
361
362    #[cfg(test)]
363    pub fn new(config: &Config) -> Self {
364        Self::new_with_config(&RecoveryConfig::from_config(config))
365    }
366}
367
368/// Available congestion control algorithms.
369///
370/// This enum provides currently available list of congestion control
371/// algorithms.
372#[derive(Debug, Copy, Clone, PartialEq, Eq)]
373#[repr(C)]
374pub enum CongestionControlAlgorithm {
375    /// Reno congestion control algorithm. `reno` in a string form.
376    Reno            = 0,
377    /// CUBIC congestion control algorithm (default). `cubic` in a string form.
378    CUBIC           = 1,
379    /// BBRv2 congestion control algorithm implementation from gcongestion
380    /// branch. `bbr2_gcongestion` in a string form.
381    Bbr2Gcongestion = 4,
382}
383
384impl FromStr for CongestionControlAlgorithm {
385    type Err = crate::Error;
386
387    /// Converts a string to `CongestionControlAlgorithm`.
388    ///
389    /// If `name` is not valid, `Error::CongestionControl` is returned.
390    fn from_str(name: &str) -> std::result::Result<Self, Self::Err> {
391        match name {
392            "reno" => Ok(CongestionControlAlgorithm::Reno),
393            "cubic" => Ok(CongestionControlAlgorithm::CUBIC),
394            "bbr" => Ok(CongestionControlAlgorithm::Bbr2Gcongestion),
395            "bbr2" => Ok(CongestionControlAlgorithm::Bbr2Gcongestion),
396            "bbr2_gcongestion" => Ok(CongestionControlAlgorithm::Bbr2Gcongestion),
397            _ => Err(crate::Error::CongestionControl),
398        }
399    }
400}
401
402#[derive(Clone)]
403pub struct Sent {
404    pub pkt_num: u64,
405
406    pub frames: SmallVec<[frame::Frame; 1]>,
407
408    pub time_sent: Instant,
409
410    pub time_acked: Option<Instant>,
411
412    pub time_lost: Option<Instant>,
413
414    pub size: usize,
415
416    pub ack_eliciting: bool,
417
418    pub in_flight: bool,
419
420    pub delivered: usize,
421
422    pub delivered_time: Instant,
423
424    pub first_sent_time: Instant,
425
426    pub is_app_limited: bool,
427
428    pub tx_in_flight: usize,
429
430    pub lost: u64,
431
432    pub has_data: bool,
433
434    pub is_pmtud_probe: bool,
435}
436
437impl std::fmt::Debug for Sent {
438    fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
439        write!(f, "pkt_num={:?} ", self.pkt_num)?;
440        write!(f, "pkt_sent_time={:?} ", self.time_sent)?;
441        write!(f, "pkt_size={:?} ", self.size)?;
442        write!(f, "delivered={:?} ", self.delivered)?;
443        write!(f, "delivered_time={:?} ", self.delivered_time)?;
444        write!(f, "first_sent_time={:?} ", self.first_sent_time)?;
445        write!(f, "is_app_limited={} ", self.is_app_limited)?;
446        write!(f, "tx_in_flight={} ", self.tx_in_flight)?;
447        write!(f, "lost={} ", self.lost)?;
448        write!(f, "has_data={} ", self.has_data)?;
449        write!(f, "is_pmtud_probe={}", self.is_pmtud_probe)?;
450
451        Ok(())
452    }
453}
454
455#[derive(Clone, Copy, Debug)]
456pub struct HandshakeStatus {
457    pub has_handshake_keys: bool,
458
459    pub peer_verified_address: bool,
460
461    pub completed: bool,
462}
463
464#[cfg(test)]
465impl Default for HandshakeStatus {
466    fn default() -> HandshakeStatus {
467        HandshakeStatus {
468            has_handshake_keys: true,
469
470            peer_verified_address: true,
471
472            completed: true,
473        }
474    }
475}
476
477// We don't need to log all qlog metrics every time there is a recovery event.
478// Instead, we can log only the MetricsUpdated event data fields that we care
479// about, only when they change. To support this, the QLogMetrics structure
480// keeps a running picture of the fields.
481#[derive(Default)]
482#[cfg(feature = "qlog")]
483struct QlogMetrics {
484    min_rtt: Duration,
485    smoothed_rtt: Duration,
486    latest_rtt: Duration,
487    rttvar: Duration,
488    cwnd: u64,
489    bytes_in_flight: u64,
490    ssthresh: Option<u64>,
491    pacing_rate: Option<u64>,
492    delivery_rate: Option<u64>,
493    send_rate: Option<u64>,
494    ack_rate: Option<u64>,
495    lost_packets: Option<u64>,
496    lost_bytes: Option<u64>,
497    pto_count: Option<u32>,
498    app_limited: Option<bool>,
499}
500
501#[cfg(feature = "qlog")]
502trait CustomCfQlogField {
503    fn name(&self) -> &'static str;
504    fn as_json_value(&self) -> serde_json::Value;
505}
506
507#[cfg(feature = "qlog")]
508#[serde_with::skip_serializing_none]
509#[derive(Serialize)]
510struct TotalAndDelta {
511    total: Option<u64>,
512    delta: Option<u64>,
513}
514
515#[cfg(feature = "qlog")]
516struct CustomQlogField<T> {
517    name: &'static str,
518    value: T,
519}
520
521#[cfg(feature = "qlog")]
522impl<T> CustomQlogField<T> {
523    fn new(name: &'static str, value: T) -> Self {
524        Self { name, value }
525    }
526}
527
528#[cfg(feature = "qlog")]
529impl<T: Serialize> CustomCfQlogField for CustomQlogField<T> {
530    fn name(&self) -> &'static str {
531        self.name
532    }
533
534    fn as_json_value(&self) -> serde_json::Value {
535        serde_json::json!(&self.value)
536    }
537}
538
539#[cfg(feature = "qlog")]
540struct CfExData(qlog::events::ExData);
541
542#[cfg(feature = "qlog")]
543impl CfExData {
544    fn new() -> Self {
545        Self(qlog::events::ExData::new())
546    }
547
548    fn insert<T: Serialize>(&mut self, name: &'static str, value: T) {
549        let field = CustomQlogField::new(name, value);
550        self.0
551            .insert(field.name().to_string(), field.as_json_value());
552    }
553
554    fn into_inner(self) -> qlog::events::ExData {
555        self.0
556    }
557}
558
559#[cfg(feature = "qlog")]
560impl QlogMetrics {
561    // Make a qlog event if the latest instance of QlogMetrics is different.
562    //
563    // This function diffs each of the fields. A qlog MetricsUpdated event is
564    // only generated if at least one field is different. Where fields are
565    // different, the qlog event contains the latest value.
566    fn maybe_update(&mut self, latest: Self) -> Option<EventData> {
567        let mut emit_event = false;
568
569        let new_min_rtt = if self.min_rtt != latest.min_rtt {
570            self.min_rtt = latest.min_rtt;
571            emit_event = true;
572            Some(latest.min_rtt.as_secs_f32() * 1000.0)
573        } else {
574            None
575        };
576
577        let new_smoothed_rtt = if self.smoothed_rtt != latest.smoothed_rtt {
578            self.smoothed_rtt = latest.smoothed_rtt;
579            emit_event = true;
580            Some(latest.smoothed_rtt.as_secs_f32() * 1000.0)
581        } else {
582            None
583        };
584
585        let new_latest_rtt = if self.latest_rtt != latest.latest_rtt {
586            self.latest_rtt = latest.latest_rtt;
587            emit_event = true;
588            Some(latest.latest_rtt.as_secs_f32() * 1000.0)
589        } else {
590            None
591        };
592
593        let new_rttvar = if self.rttvar != latest.rttvar {
594            self.rttvar = latest.rttvar;
595            emit_event = true;
596            Some(latest.rttvar.as_secs_f32() * 1000.0)
597        } else {
598            None
599        };
600
601        let new_cwnd = if self.cwnd != latest.cwnd {
602            self.cwnd = latest.cwnd;
603            emit_event = true;
604            Some(latest.cwnd)
605        } else {
606            None
607        };
608
609        let new_bytes_in_flight =
610            if self.bytes_in_flight != latest.bytes_in_flight {
611                self.bytes_in_flight = latest.bytes_in_flight;
612                emit_event = true;
613                Some(latest.bytes_in_flight)
614            } else {
615                None
616            };
617
618        let new_ssthresh = if self.ssthresh != latest.ssthresh {
619            self.ssthresh = latest.ssthresh;
620            emit_event = true;
621            latest.ssthresh
622        } else {
623            None
624        };
625
626        let new_pacing_rate = if self.pacing_rate != latest.pacing_rate {
627            self.pacing_rate = latest.pacing_rate;
628            emit_event = true;
629            latest.pacing_rate
630        } else {
631            None
632        };
633
634        let new_pto_count =
635            if latest.pto_count.is_some() && self.pto_count != latest.pto_count {
636                self.pto_count = latest.pto_count;
637                emit_event = true;
638                latest.pto_count.map(|v| v as u16)
639            } else {
640                None
641            };
642
643        // Build ex_data for rate metrics
644        let mut ex_data = CfExData::new();
645        if self.app_limited != latest.app_limited {
646            if let Some(app_limited) = latest.app_limited {
647                self.app_limited = latest.app_limited;
648                emit_event = true;
649                ex_data.insert("cf_app_limited", app_limited);
650            }
651        }
652        if self.delivery_rate != latest.delivery_rate {
653            if let Some(rate) = latest.delivery_rate {
654                self.delivery_rate = latest.delivery_rate;
655                emit_event = true;
656                ex_data.insert("cf_delivery_rate", rate);
657            }
658        }
659        if self.send_rate != latest.send_rate {
660            if let Some(rate) = latest.send_rate {
661                self.send_rate = latest.send_rate;
662                emit_event = true;
663                ex_data.insert("cf_send_rate", rate);
664            }
665        }
666        if self.ack_rate != latest.ack_rate {
667            if let Some(rate) = latest.ack_rate {
668                self.ack_rate = latest.ack_rate;
669                emit_event = true;
670                ex_data.insert("cf_ack_rate", rate);
671            }
672        }
673
674        if self.lost_packets != latest.lost_packets {
675            if let Some(val) = latest.lost_packets {
676                emit_event = true;
677                ex_data.insert("cf_lost_packets", TotalAndDelta {
678                    total: latest.lost_packets,
679                    delta: Some(val - self.lost_packets.unwrap_or(0)),
680                });
681                self.lost_packets = latest.lost_packets;
682            }
683        }
684        if self.lost_bytes != latest.lost_bytes {
685            if let Some(val) = latest.lost_bytes {
686                emit_event = true;
687                ex_data.insert("cf_lost_bytes", TotalAndDelta {
688                    total: latest.lost_bytes,
689                    delta: Some(val - self.lost_bytes.unwrap_or(0)),
690                });
691                self.lost_bytes = latest.lost_bytes;
692            }
693        }
694
695        if emit_event {
696            return Some(EventData::QuicMetricsUpdated(
697                qlog::events::quic::RecoveryMetricsUpdated {
698                    min_rtt: new_min_rtt,
699                    smoothed_rtt: new_smoothed_rtt,
700                    latest_rtt: new_latest_rtt,
701                    rtt_variance: new_rttvar,
702                    congestion_window: new_cwnd,
703                    bytes_in_flight: new_bytes_in_flight,
704                    ssthresh: new_ssthresh,
705                    pacing_rate: new_pacing_rate,
706                    pto_count: new_pto_count,
707                    ex_data: ex_data.into_inner(),
708                    ..Default::default()
709                },
710            ));
711        }
712
713        None
714    }
715}
716
717/// When the pacer thinks is a good time to release the next packet
718#[derive(Debug, Clone, Copy, PartialEq, Eq)]
719pub enum ReleaseTime {
720    Immediate,
721    At(Instant),
722}
723
724/// When the next packet should be release and if it can be part of a burst
725#[derive(Clone, Copy, Debug, PartialEq, Eq)]
726pub struct ReleaseDecision {
727    time: ReleaseTime,
728    allow_burst: bool,
729}
730
731impl ReleaseTime {
732    /// Add the specific delay to the current time
733    fn inc(&mut self, delay: Duration) {
734        match self {
735            ReleaseTime::Immediate => {},
736            ReleaseTime::At(time) => *time += delay,
737        }
738    }
739
740    /// Set the time to the later of two times
741    fn set_max(&mut self, other: Instant) {
742        match self {
743            ReleaseTime::Immediate => *self = ReleaseTime::At(other),
744            ReleaseTime::At(time) => *self = ReleaseTime::At(other.max(*time)),
745        }
746    }
747}
748
749impl ReleaseDecision {
750    pub(crate) const EQUAL_THRESHOLD: Duration = Duration::from_micros(50);
751
752    /// Get the [`Instant`] the next packet should be released. It will never be
753    /// in the past.
754    #[inline]
755    pub fn time(&self, now: Instant) -> Option<Instant> {
756        match self.time {
757            ReleaseTime::Immediate => None,
758            ReleaseTime::At(other) => other.gt(&now).then_some(other),
759        }
760    }
761
762    /// Can this packet be appended to a previous burst
763    #[inline]
764    pub fn can_burst(&self) -> bool {
765        self.allow_burst
766    }
767
768    /// Check if the two packets can be released at the same time
769    #[inline]
770    pub fn time_eq(&self, other: &Self, now: Instant) -> bool {
771        let delta = match (self.time(now), other.time(now)) {
772            (None, None) => Duration::ZERO,
773            (Some(t), None) | (None, Some(t)) => t.duration_since(now),
774            (Some(t1), Some(t2)) if t1 < t2 => t2.duration_since(t1),
775            (Some(t1), Some(t2)) => t1.duration_since(t2),
776        };
777
778        delta <= Self::EQUAL_THRESHOLD
779    }
780}
781
782/// Recovery statistics
783#[derive(Default, Debug)]
784pub struct RecoveryStats {
785    startup_exit: Option<StartupExit>,
786}
787
788impl RecoveryStats {
789    // Record statistics when a CCA first exits startup.
790    pub fn set_startup_exit(&mut self, startup_exit: StartupExit) {
791        if self.startup_exit.is_none() {
792            self.startup_exit = Some(startup_exit);
793        }
794    }
795}
796
797/// Statistics from when a CCA first exited the startup phase.
798#[derive(Debug, Clone, Copy, PartialEq)]
799pub struct StartupExit {
800    /// The congestion_window recorded at Startup exit.
801    pub cwnd: usize,
802
803    /// The bandwidth estimate recorded at Startup exit.
804    pub bandwidth: Option<u64>,
805
806    /// The reason a CCA exited the startup phase.
807    pub reason: StartupExitReason,
808}
809
810impl StartupExit {
811    fn new(
812        cwnd: usize, bandwidth: Option<Bandwidth>, reason: StartupExitReason,
813    ) -> Self {
814        let bandwidth = bandwidth.map(Bandwidth::to_bytes_per_second);
815        Self {
816            cwnd,
817            bandwidth,
818            reason,
819        }
820    }
821}
822
823/// The reason a CCA exited the startup phase.
824#[derive(Debug, Clone, Copy, PartialEq)]
825pub enum StartupExitReason {
826    /// Exit slow start or BBR startup due to excessive loss
827    Loss,
828
829    /// Exit BBR startup due to bandwidth plateau.
830    BandwidthPlateau,
831
832    /// Exit BBR startup due to persistent queue.
833    PersistentQueue,
834
835    /// Exit HyStart++ conservative slow start after the max rounds allowed.
836    ConservativeSlowStartRounds,
837}
838
839#[cfg(test)]
840mod tests {
841    use super::*;
842    use crate::packet;
843    use crate::range_buf::RangeBuf;
844    use crate::test_utils;
845    use crate::CongestionControlAlgorithm;
846    use crate::DEFAULT_INITIAL_RTT;
847    use rstest::rstest;
848    use smallvec::smallvec;
849    use std::str::FromStr;
850
851    fn recovery_for_alg(algo: CongestionControlAlgorithm) -> Recovery {
852        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
853        cfg.set_cc_algorithm(algo);
854        Recovery::new(&cfg)
855    }
856
857    #[cfg(feature = "qlog")]
858    fn app_limited_value(event: EventData) -> Option<serde_json::Value> {
859        let EventData::QuicMetricsUpdated(metrics) = event else {
860            panic!("expected recovery metrics updated event");
861        };
862
863        metrics.ex_data.get("cf_app_limited").cloned()
864    }
865
866    #[cfg(feature = "qlog")]
867    #[test]
868    fn qlog_app_limited_emits_initial_false_and_transitions() {
869        let mut metrics = QlogMetrics::default();
870
871        let event = metrics
872            .maybe_update(QlogMetrics {
873                app_limited: Some(false),
874                ..Default::default()
875            })
876            .unwrap();
877        assert_eq!(app_limited_value(event), Some(false.into()));
878
879        let event = metrics
880            .maybe_update(QlogMetrics {
881                app_limited: Some(true),
882                ..Default::default()
883            })
884            .unwrap();
885        assert_eq!(app_limited_value(event), Some(true.into()));
886
887        let event = metrics
888            .maybe_update(QlogMetrics {
889                app_limited: Some(false),
890                ..Default::default()
891            })
892            .unwrap();
893        assert_eq!(app_limited_value(event), Some(false.into()));
894    }
895
896    #[cfg(feature = "qlog")]
897    #[test]
898    fn qlog_app_limited_suppresses_unchanged_values() {
899        let mut metrics = QlogMetrics::default();
900        metrics
901            .maybe_update(QlogMetrics {
902                app_limited: Some(false),
903                ..Default::default()
904            })
905            .unwrap();
906
907        assert!(metrics
908            .maybe_update(QlogMetrics {
909                app_limited: Some(false),
910                ..Default::default()
911            })
912            .is_none());
913
914        let event = metrics
915            .maybe_update(QlogMetrics {
916                cwnd: 1,
917                app_limited: Some(false),
918                ..Default::default()
919            })
920            .unwrap();
921        assert_eq!(app_limited_value(event), None);
922    }
923
924    #[test]
925    fn lookup_cc_algo_ok() {
926        let algo = CongestionControlAlgorithm::from_str("reno").unwrap();
927        assert_eq!(algo, CongestionControlAlgorithm::Reno);
928        assert!(!recovery_for_alg(algo).gcongestion_enabled());
929
930        let algo = CongestionControlAlgorithm::from_str("cubic").unwrap();
931        assert_eq!(algo, CongestionControlAlgorithm::CUBIC);
932        assert!(!recovery_for_alg(algo).gcongestion_enabled());
933
934        let algo = CongestionControlAlgorithm::from_str("bbr").unwrap();
935        assert_eq!(algo, CongestionControlAlgorithm::Bbr2Gcongestion);
936        assert!(recovery_for_alg(algo).gcongestion_enabled());
937
938        let algo = CongestionControlAlgorithm::from_str("bbr2").unwrap();
939        assert_eq!(algo, CongestionControlAlgorithm::Bbr2Gcongestion);
940        assert!(recovery_for_alg(algo).gcongestion_enabled());
941
942        let algo =
943            CongestionControlAlgorithm::from_str("bbr2_gcongestion").unwrap();
944        assert_eq!(algo, CongestionControlAlgorithm::Bbr2Gcongestion);
945        assert!(recovery_for_alg(algo).gcongestion_enabled());
946    }
947
948    #[test]
949    fn lookup_cc_algo_bad() {
950        assert_eq!(
951            CongestionControlAlgorithm::from_str("???"),
952            Err(crate::Error::CongestionControl)
953        );
954    }
955
956    #[rstest]
957    fn loss_on_pto(
958        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
959    ) {
960        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
961        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
962
963        let mut r = Recovery::new(&cfg);
964
965        let mut now = Instant::now();
966
967        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
968
969        // Start by sending a few packets.
970        let p = Sent {
971            pkt_num: 0,
972            frames: smallvec![],
973            time_sent: now,
974            time_acked: None,
975            time_lost: None,
976            size: 1000,
977            ack_eliciting: true,
978            in_flight: true,
979            delivered: 0,
980            delivered_time: now,
981            first_sent_time: now,
982            is_app_limited: false,
983            tx_in_flight: 0,
984            lost: 0,
985            has_data: false,
986            is_pmtud_probe: false,
987        };
988
989        r.on_packet_sent(
990            p,
991            packet::Epoch::Application,
992            HandshakeStatus::default(),
993            now,
994            "",
995        );
996
997        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 1);
998        assert_eq!(r.bytes_in_flight(), 1000);
999        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1000
1001        let p = Sent {
1002            pkt_num: 1,
1003            frames: smallvec![],
1004            time_sent: now,
1005            time_acked: None,
1006            time_lost: None,
1007            size: 1000,
1008            ack_eliciting: true,
1009            in_flight: true,
1010            delivered: 0,
1011            delivered_time: now,
1012            first_sent_time: now,
1013            is_app_limited: false,
1014            tx_in_flight: 0,
1015            lost: 0,
1016            has_data: false,
1017            is_pmtud_probe: false,
1018        };
1019
1020        r.on_packet_sent(
1021            p,
1022            packet::Epoch::Application,
1023            HandshakeStatus::default(),
1024            now,
1025            "",
1026        );
1027
1028        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
1029        assert_eq!(r.bytes_in_flight(), 2000);
1030        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1031
1032        let p = Sent {
1033            pkt_num: 2,
1034            frames: smallvec![],
1035            time_sent: now,
1036            time_acked: None,
1037            time_lost: None,
1038            size: 1000,
1039            ack_eliciting: true,
1040            in_flight: true,
1041            delivered: 0,
1042            delivered_time: now,
1043            first_sent_time: now,
1044            is_app_limited: false,
1045            tx_in_flight: 0,
1046            lost: 0,
1047            has_data: false,
1048            is_pmtud_probe: false,
1049        };
1050
1051        r.on_packet_sent(
1052            p,
1053            packet::Epoch::Application,
1054            HandshakeStatus::default(),
1055            now,
1056            "",
1057        );
1058        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 3);
1059        assert_eq!(r.bytes_in_flight(), 3000);
1060        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1061
1062        let p = Sent {
1063            pkt_num: 3,
1064            frames: smallvec![],
1065            time_sent: now,
1066            time_acked: None,
1067            time_lost: None,
1068            size: 1000,
1069            ack_eliciting: true,
1070            in_flight: true,
1071            delivered: 0,
1072            delivered_time: now,
1073            first_sent_time: now,
1074            is_app_limited: false,
1075            tx_in_flight: 0,
1076            lost: 0,
1077            has_data: false,
1078            is_pmtud_probe: false,
1079        };
1080
1081        r.on_packet_sent(
1082            p,
1083            packet::Epoch::Application,
1084            HandshakeStatus::default(),
1085            now,
1086            "",
1087        );
1088        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 4);
1089        assert_eq!(r.bytes_in_flight(), 4000);
1090        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1091
1092        // Wait for 10ms.
1093        now += Duration::from_millis(10);
1094
1095        // Only the first 2 packets are acked.
1096        let mut acked = RangeSet::default();
1097        acked.insert(0..2);
1098
1099        assert_eq!(
1100            r.on_ack_received(
1101                &acked,
1102                25,
1103                packet::Epoch::Application,
1104                HandshakeStatus::default(),
1105                now,
1106                None,
1107                "",
1108            )
1109            .unwrap(),
1110            OnAckReceivedOutcome {
1111                lost_packets: 0,
1112                lost_bytes: 0,
1113                acked_bytes: 2 * 1000,
1114                spurious_losses: 0,
1115            }
1116        );
1117
1118        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
1119        assert_eq!(r.bytes_in_flight(), 2000);
1120        assert_eq!(r.bytes_in_flight_duration(), Duration::from_millis(10));
1121        assert_eq!(r.lost_count(), 0);
1122
1123        // Wait until loss detection timer expires.
1124        now = r.loss_detection_timer().unwrap();
1125
1126        // PTO.
1127        r.on_loss_detection_timeout(HandshakeStatus::default(), now, "");
1128        assert_eq!(r.loss_probes(packet::Epoch::Application), 1);
1129        assert_eq!(r.lost_count(), 0);
1130        assert_eq!(r.pto_count(), 1);
1131
1132        let p = Sent {
1133            pkt_num: 4,
1134            frames: smallvec![],
1135            time_sent: now,
1136            time_acked: None,
1137            time_lost: None,
1138            size: 1000,
1139            ack_eliciting: true,
1140            in_flight: true,
1141            delivered: 0,
1142            delivered_time: now,
1143            first_sent_time: now,
1144            is_app_limited: false,
1145            tx_in_flight: 0,
1146            lost: 0,
1147            has_data: false,
1148            is_pmtud_probe: false,
1149        };
1150
1151        r.on_packet_sent(
1152            p,
1153            packet::Epoch::Application,
1154            HandshakeStatus::default(),
1155            now,
1156            "",
1157        );
1158        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 3);
1159        assert_eq!(r.bytes_in_flight(), 3000);
1160        assert_eq!(r.bytes_in_flight_duration(), Duration::from_millis(30));
1161
1162        let p = Sent {
1163            pkt_num: 5,
1164            frames: smallvec![],
1165            time_sent: now,
1166            time_acked: None,
1167            time_lost: None,
1168            size: 1000,
1169            ack_eliciting: true,
1170            in_flight: true,
1171            delivered: 0,
1172            delivered_time: now,
1173            first_sent_time: now,
1174            is_app_limited: false,
1175            tx_in_flight: 0,
1176            lost: 0,
1177            has_data: false,
1178            is_pmtud_probe: false,
1179        };
1180
1181        r.on_packet_sent(
1182            p,
1183            packet::Epoch::Application,
1184            HandshakeStatus::default(),
1185            now,
1186            "",
1187        );
1188        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 4);
1189        assert_eq!(r.bytes_in_flight(), 4000);
1190        assert_eq!(r.bytes_in_flight_duration(), Duration::from_millis(30));
1191        assert_eq!(r.lost_count(), 0);
1192
1193        // Wait for 10ms.
1194        now += Duration::from_millis(10);
1195
1196        // PTO packets are acked.
1197        let mut acked = RangeSet::default();
1198        acked.insert(4..6);
1199
1200        assert_eq!(
1201            r.on_ack_received(
1202                &acked,
1203                25,
1204                packet::Epoch::Application,
1205                HandshakeStatus::default(),
1206                now,
1207                None,
1208                "",
1209            )
1210            .unwrap(),
1211            OnAckReceivedOutcome {
1212                lost_packets: 2,
1213                lost_bytes: 2000,
1214                acked_bytes: 2 * 1000,
1215                spurious_losses: 0,
1216            }
1217        );
1218
1219        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 4);
1220        assert_eq!(r.bytes_in_flight(), 0);
1221        assert_eq!(r.bytes_in_flight_duration(), Duration::from_millis(40));
1222
1223        assert_eq!(r.lost_count(), 2);
1224
1225        // Wait 1 RTT.
1226        now += r.rtt();
1227
1228        assert_eq!(
1229            r.detect_lost_packets_for_test(packet::Epoch::Application, now),
1230            (0, 0)
1231        );
1232
1233        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
1234        if cc_algorithm_name == "reno" || cc_algorithm_name == "cubic" {
1235            assert!(r.startup_exit().is_some());
1236            assert_eq!(r.startup_exit().unwrap().reason, StartupExitReason::Loss);
1237        } else {
1238            assert_eq!(r.startup_exit(), None);
1239        }
1240    }
1241
1242    #[rstest]
1243    fn loss_on_timer(
1244        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1245    ) {
1246        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
1247        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
1248
1249        let mut r = Recovery::new(&cfg);
1250
1251        let mut now = Instant::now();
1252
1253        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
1254
1255        // Start by sending a few packets.
1256        let p = Sent {
1257            pkt_num: 0,
1258            frames: smallvec![],
1259            time_sent: now,
1260            time_acked: None,
1261            time_lost: None,
1262            size: 1000,
1263            ack_eliciting: true,
1264            in_flight: true,
1265            delivered: 0,
1266            delivered_time: now,
1267            first_sent_time: now,
1268            is_app_limited: false,
1269            tx_in_flight: 0,
1270            lost: 0,
1271            has_data: false,
1272            is_pmtud_probe: false,
1273        };
1274
1275        r.on_packet_sent(
1276            p,
1277            packet::Epoch::Application,
1278            HandshakeStatus::default(),
1279            now,
1280            "",
1281        );
1282        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 1);
1283        assert_eq!(r.bytes_in_flight(), 1000);
1284        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1285
1286        let p = Sent {
1287            pkt_num: 1,
1288            frames: smallvec![],
1289            time_sent: now,
1290            time_acked: None,
1291            time_lost: None,
1292            size: 1000,
1293            ack_eliciting: true,
1294            in_flight: true,
1295            delivered: 0,
1296            delivered_time: now,
1297            first_sent_time: now,
1298            is_app_limited: false,
1299            tx_in_flight: 0,
1300            lost: 0,
1301            has_data: false,
1302            is_pmtud_probe: false,
1303        };
1304
1305        r.on_packet_sent(
1306            p,
1307            packet::Epoch::Application,
1308            HandshakeStatus::default(),
1309            now,
1310            "",
1311        );
1312        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
1313        assert_eq!(r.bytes_in_flight(), 2000);
1314        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1315
1316        let p = Sent {
1317            pkt_num: 2,
1318            frames: smallvec![],
1319            time_sent: now,
1320            time_acked: None,
1321            time_lost: None,
1322            size: 1000,
1323            ack_eliciting: true,
1324            in_flight: true,
1325            delivered: 0,
1326            delivered_time: now,
1327            first_sent_time: now,
1328            is_app_limited: false,
1329            tx_in_flight: 0,
1330            lost: 0,
1331            has_data: false,
1332            is_pmtud_probe: false,
1333        };
1334
1335        r.on_packet_sent(
1336            p,
1337            packet::Epoch::Application,
1338            HandshakeStatus::default(),
1339            now,
1340            "",
1341        );
1342        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 3);
1343        assert_eq!(r.bytes_in_flight(), 3000);
1344        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1345
1346        let p = Sent {
1347            pkt_num: 3,
1348            frames: smallvec![],
1349            time_sent: now,
1350            time_acked: None,
1351            time_lost: None,
1352            size: 1000,
1353            ack_eliciting: true,
1354            in_flight: true,
1355            delivered: 0,
1356            delivered_time: now,
1357            first_sent_time: now,
1358            is_app_limited: false,
1359            tx_in_flight: 0,
1360            lost: 0,
1361            has_data: false,
1362            is_pmtud_probe: false,
1363        };
1364
1365        r.on_packet_sent(
1366            p,
1367            packet::Epoch::Application,
1368            HandshakeStatus::default(),
1369            now,
1370            "",
1371        );
1372        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 4);
1373        assert_eq!(r.bytes_in_flight(), 4000);
1374        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1375
1376        // Wait for 10ms.
1377        now += Duration::from_millis(10);
1378
1379        // Only the first 2 packets and the last one are acked.
1380        let mut acked = RangeSet::default();
1381        acked.insert(0..2);
1382        acked.insert(3..4);
1383
1384        assert_eq!(
1385            r.on_ack_received(
1386                &acked,
1387                25,
1388                packet::Epoch::Application,
1389                HandshakeStatus::default(),
1390                now,
1391                None,
1392                "",
1393            )
1394            .unwrap(),
1395            OnAckReceivedOutcome {
1396                lost_packets: 0,
1397                lost_bytes: 0,
1398                acked_bytes: 3 * 1000,
1399                spurious_losses: 0,
1400            }
1401        );
1402
1403        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
1404        assert_eq!(r.bytes_in_flight(), 1000);
1405        assert_eq!(r.bytes_in_flight_duration(), Duration::from_millis(10));
1406        assert_eq!(r.lost_count(), 0);
1407
1408        // Wait until loss detection timer expires.
1409        now = r.loss_detection_timer().unwrap();
1410
1411        // Packet is declared lost.
1412        r.on_loss_detection_timeout(HandshakeStatus::default(), now, "");
1413        assert_eq!(r.loss_probes(packet::Epoch::Application), 0);
1414
1415        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
1416        assert_eq!(r.bytes_in_flight(), 0);
1417        assert_eq!(r.bytes_in_flight_duration(), Duration::from_micros(11250));
1418
1419        assert_eq!(r.lost_count(), 1);
1420
1421        // Wait 1 RTT.
1422        now += r.rtt();
1423
1424        assert_eq!(
1425            r.detect_lost_packets_for_test(packet::Epoch::Application, now),
1426            (0, 0)
1427        );
1428
1429        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
1430        if cc_algorithm_name == "reno" || cc_algorithm_name == "cubic" {
1431            assert!(r.startup_exit().is_some());
1432            assert_eq!(r.startup_exit().unwrap().reason, StartupExitReason::Loss);
1433        } else {
1434            assert_eq!(r.startup_exit(), None);
1435        }
1436    }
1437
1438    #[rstest]
1439    fn loss_on_reordering(
1440        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1441    ) {
1442        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
1443        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
1444
1445        let mut r = Recovery::new(&cfg);
1446
1447        let mut now = Instant::now();
1448
1449        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
1450
1451        // Start by sending a few packets.
1452        //
1453        // pkt number: [0, 1, 2, 3]
1454        for i in 0..4 {
1455            let p = test_utils::helper_packet_sent(i, now, 1000);
1456            r.on_packet_sent(
1457                p,
1458                packet::Epoch::Application,
1459                HandshakeStatus::default(),
1460                now,
1461                "",
1462            );
1463
1464            let pkt_count = (i + 1) as usize;
1465            assert_eq!(r.sent_packets_len(packet::Epoch::Application), pkt_count);
1466            assert_eq!(r.bytes_in_flight(), pkt_count * 1000);
1467            assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1468        }
1469
1470        // Wait for 10ms after sending.
1471        now += Duration::from_millis(10);
1472
1473        // Recieve reordered ACKs, i.e. pkt_num [2, 3]
1474        let mut acked = RangeSet::default();
1475        acked.insert(2..4);
1476        assert_eq!(
1477            r.on_ack_received(
1478                &acked,
1479                25,
1480                packet::Epoch::Application,
1481                HandshakeStatus::default(),
1482                now,
1483                None,
1484                "",
1485            )
1486            .unwrap(),
1487            OnAckReceivedOutcome {
1488                lost_packets: 1,
1489                lost_bytes: 1000,
1490                acked_bytes: 1000 * 2,
1491                spurious_losses: 0,
1492            }
1493        );
1494        // Since we only remove packets from the back to avoid compaction, the
1495        // send length remains the same after receiving reordered ACKs
1496        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 4);
1497        assert_eq!(r.pkt_thresh().unwrap(), INITIAL_PACKET_THRESHOLD);
1498        assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1499
1500        // Wait for 10ms after receiving first set of ACKs.
1501        now += Duration::from_millis(10);
1502
1503        // Recieve remaining ACKs, i.e. pkt_num [0, 1]
1504        let mut acked = RangeSet::default();
1505        acked.insert(0..2);
1506        assert_eq!(
1507            r.on_ack_received(
1508                &acked,
1509                25,
1510                packet::Epoch::Application,
1511                HandshakeStatus::default(),
1512                now,
1513                None,
1514                "",
1515            )
1516            .unwrap(),
1517            OnAckReceivedOutcome {
1518                lost_packets: 0,
1519                lost_bytes: 0,
1520                acked_bytes: 1000,
1521                spurious_losses: 1,
1522            }
1523        );
1524        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
1525        assert_eq!(r.bytes_in_flight(), 0);
1526        assert_eq!(r.bytes_in_flight_duration(), Duration::from_millis(20));
1527
1528        // Spurious loss.
1529        assert_eq!(r.lost_count(), 1);
1530        assert_eq!(r.lost_spurious_count(), 1);
1531
1532        // Packet threshold was increased.
1533        assert_eq!(r.pkt_thresh().unwrap(), 4);
1534        assert_eq!(r.time_thresh(), PACKET_REORDER_TIME_THRESHOLD);
1535
1536        // Wait 1 RTT.
1537        now += r.rtt();
1538
1539        // All packets have been ACKed so dont expect additional lost packets
1540        assert_eq!(
1541            r.detect_lost_packets_for_test(packet::Epoch::Application, now),
1542            (0, 0)
1543        );
1544        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
1545
1546        if cc_algorithm_name == "reno" || cc_algorithm_name == "cubic" {
1547            assert!(r.startup_exit().is_some());
1548            assert_eq!(r.startup_exit().unwrap().reason, StartupExitReason::Loss);
1549        } else {
1550            assert_eq!(r.startup_exit(), None);
1551        }
1552    }
1553
1554    // TODO: This should run agains both `congestion` and `gcongestion`.
1555    // `congestion` and `gcongestion` behave differently. That might be ok
1556    // given the different algorithms but it would be ideal to merge and share
1557    // the logic.
1558    #[rstest]
1559    fn time_thresholds_on_reordering(
1560        #[values("bbr2_gcongestion")] cc_algorithm_name: &str,
1561    ) {
1562        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
1563        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
1564
1565        let mut now = Instant::now();
1566        let mut r = Recovery::new(&cfg);
1567        assert_eq!(r.rtt(), DEFAULT_INITIAL_RTT);
1568
1569        // Choose times around the threshold to test its increase.
1570        //
1571        // ```
1572        //              between_thresh_ms
1573        //                         |
1574        //    initial_thresh_ms    |     spurious_thresh_ms
1575        //      v                  v             v
1576        // --------------------------------------------------
1577        //      | ................ | ..................... |
1578        //            THRESH_GAP         THRESH_GAP
1579        // ```
1580        //
1581        // Threshold gap time.
1582        const THRESH_GAP: Duration = Duration::from_millis(30);
1583        // Initial time theshold based on inital RTT.
1584        let initial_thresh_ms =
1585            DEFAULT_INITIAL_RTT.mul_f64(INITIAL_TIME_THRESHOLD);
1586        // The time threshold after spurious loss.
1587        let spurious_thresh_ms: Duration =
1588            DEFAULT_INITIAL_RTT.mul_f64(PACKET_REORDER_TIME_THRESHOLD);
1589        // Time between the two thresholds
1590        let between_thresh_ms = initial_thresh_ms + THRESH_GAP;
1591        assert!(between_thresh_ms > initial_thresh_ms);
1592        assert!(between_thresh_ms < spurious_thresh_ms);
1593        assert!(between_thresh_ms + THRESH_GAP > spurious_thresh_ms);
1594
1595        for i in 0..6 {
1596            let send_time = now + i * between_thresh_ms;
1597
1598            let p = test_utils::helper_packet_sent(i.into(), send_time, 1000);
1599            r.on_packet_sent(
1600                p,
1601                packet::Epoch::Application,
1602                HandshakeStatus::default(),
1603                send_time,
1604                "",
1605            );
1606        }
1607
1608        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 6);
1609        assert_eq!(r.bytes_in_flight(), 6 * 1000);
1610        assert_eq!(r.pkt_thresh().unwrap(), INITIAL_PACKET_THRESHOLD);
1611        assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1612
1613        // Wait for `between_thresh_ms` after sending to trigger loss based on
1614        // loss threshold.
1615        now += between_thresh_ms;
1616
1617        // Ack packet: 1
1618        //
1619        // [0, 1, 2, 3, 4, 5]
1620        //     ^
1621        let mut acked = RangeSet::default();
1622        acked.insert(1..2);
1623        assert_eq!(
1624            r.on_ack_received(
1625                &acked,
1626                25,
1627                packet::Epoch::Application,
1628                HandshakeStatus::default(),
1629                now,
1630                None,
1631                "",
1632            )
1633            .unwrap(),
1634            OnAckReceivedOutcome {
1635                lost_packets: 1,
1636                lost_bytes: 1000,
1637                acked_bytes: 1000,
1638                spurious_losses: 0,
1639            }
1640        );
1641        assert_eq!(r.pkt_thresh().unwrap(), INITIAL_PACKET_THRESHOLD);
1642        assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1643
1644        // Ack packet: 0
1645        //
1646        // [0, 1, 2, 3, 4, 5]
1647        //  ^  x
1648        let mut acked = RangeSet::default();
1649        acked.insert(0..1);
1650        assert_eq!(
1651            r.on_ack_received(
1652                &acked,
1653                25,
1654                packet::Epoch::Application,
1655                HandshakeStatus::default(),
1656                now,
1657                None,
1658                "",
1659            )
1660            .unwrap(),
1661            OnAckReceivedOutcome {
1662                lost_packets: 0,
1663                lost_bytes: 0,
1664                acked_bytes: 0,
1665                spurious_losses: 1,
1666            }
1667        );
1668        // The time_thresh after spurious loss
1669        assert_eq!(r.time_thresh(), PACKET_REORDER_TIME_THRESHOLD);
1670
1671        // Wait for `between_thresh_ms` after sending. However, since the
1672        // threshold has increased, we do not expect loss.
1673        now += between_thresh_ms;
1674
1675        // Ack packet: 3
1676        //
1677        // [2, 3, 4, 5]
1678        //     ^
1679        let mut acked = RangeSet::default();
1680        acked.insert(3..4);
1681        assert_eq!(
1682            r.on_ack_received(
1683                &acked,
1684                25,
1685                packet::Epoch::Application,
1686                HandshakeStatus::default(),
1687                now,
1688                None,
1689                "",
1690            )
1691            .unwrap(),
1692            OnAckReceivedOutcome {
1693                lost_packets: 0,
1694                lost_bytes: 0,
1695                acked_bytes: 1000,
1696                spurious_losses: 0,
1697            }
1698        );
1699
1700        // Wait for and additional `plus_overhead` to trigger loss based on the
1701        // new time threshold.
1702        now += THRESH_GAP;
1703
1704        // Ack packet: 4
1705        //
1706        // [2, 3, 4, 5]
1707        //     x  ^
1708        let mut acked = RangeSet::default();
1709        acked.insert(4..5);
1710        assert_eq!(
1711            r.on_ack_received(
1712                &acked,
1713                25,
1714                packet::Epoch::Application,
1715                HandshakeStatus::default(),
1716                now,
1717                None,
1718                "",
1719            )
1720            .unwrap(),
1721            OnAckReceivedOutcome {
1722                lost_packets: 1,
1723                lost_bytes: 1000,
1724                acked_bytes: 1000,
1725                spurious_losses: 0,
1726            }
1727        );
1728    }
1729
1730    // TODO: Implement `enable_relaxed_loss_threshold` and enable this test for
1731    // the congestion module.
1732    #[rstest]
1733    fn relaxed_thresholds_on_reordering(
1734        #[values("bbr2_gcongestion")] cc_algorithm_name: &str,
1735    ) {
1736        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
1737        cfg.enable_relaxed_loss_threshold = true;
1738        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
1739
1740        let mut now = Instant::now();
1741        let mut r = Recovery::new(&cfg);
1742        assert_eq!(r.rtt(), DEFAULT_INITIAL_RTT);
1743
1744        // Choose times around the threshold to test its increase.
1745        //
1746        // ```
1747        //              between_thresh_ms
1748        //                         |
1749        //    initial_thresh_ms    |     spurious_thresh_ms
1750        //      v                  v             v
1751        // --------------------------------------------------
1752        //      | ................ | ..................... |
1753        //            THRESH_GAP         THRESH_GAP
1754        // ```
1755        // Threshold gap time.
1756        const THRESH_GAP: Duration = Duration::from_millis(30);
1757        // Initial time theshold based on inital RTT.
1758        let initial_thresh_ms =
1759            DEFAULT_INITIAL_RTT.mul_f64(INITIAL_TIME_THRESHOLD);
1760        // The time threshold after spurious loss.
1761        let spurious_thresh_ms: Duration =
1762            DEFAULT_INITIAL_RTT.mul_f64(PACKET_REORDER_TIME_THRESHOLD);
1763        // Time between the two thresholds
1764        let between_thresh_ms = initial_thresh_ms + THRESH_GAP;
1765        assert!(between_thresh_ms > initial_thresh_ms);
1766        assert!(between_thresh_ms < spurious_thresh_ms);
1767        assert!(between_thresh_ms + THRESH_GAP > spurious_thresh_ms);
1768
1769        for i in 0..6 {
1770            let send_time = now + i * between_thresh_ms;
1771
1772            let p = test_utils::helper_packet_sent(i.into(), send_time, 1000);
1773            r.on_packet_sent(
1774                p,
1775                packet::Epoch::Application,
1776                HandshakeStatus::default(),
1777                send_time,
1778                "",
1779            );
1780        }
1781
1782        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 6);
1783        assert_eq!(r.bytes_in_flight(), 6 * 1000);
1784        // Intitial thresholds
1785        assert_eq!(r.pkt_thresh().unwrap(), INITIAL_PACKET_THRESHOLD);
1786        assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1787
1788        // Wait for `between_thresh_ms` after sending to trigger loss based on
1789        // loss threshold.
1790        now += between_thresh_ms;
1791
1792        // Ack packet: 1
1793        //
1794        // [0, 1, 2, 3, 4, 5]
1795        //     ^
1796        let mut acked = RangeSet::default();
1797        acked.insert(1..2);
1798        assert_eq!(
1799            r.on_ack_received(
1800                &acked,
1801                25,
1802                packet::Epoch::Application,
1803                HandshakeStatus::default(),
1804                now,
1805                None,
1806                "",
1807            )
1808            .unwrap(),
1809            OnAckReceivedOutcome {
1810                lost_packets: 1,
1811                lost_bytes: 1000,
1812                acked_bytes: 1000,
1813                spurious_losses: 0,
1814            }
1815        );
1816        // Thresholds after 1st loss
1817        assert_eq!(r.pkt_thresh().unwrap(), INITIAL_PACKET_THRESHOLD);
1818        assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1819
1820        // Ack packet: 0
1821        //
1822        // [0, 1, 2, 3, 4, 5]
1823        //  ^  x
1824        let mut acked = RangeSet::default();
1825        acked.insert(0..1);
1826        assert_eq!(
1827            r.on_ack_received(
1828                &acked,
1829                25,
1830                packet::Epoch::Application,
1831                HandshakeStatus::default(),
1832                now,
1833                None,
1834                "",
1835            )
1836            .unwrap(),
1837            OnAckReceivedOutcome {
1838                lost_packets: 0,
1839                lost_bytes: 0,
1840                acked_bytes: 0,
1841                spurious_losses: 1,
1842            }
1843        );
1844        // Thresholds after 1st spurious loss
1845        //
1846        // Packet threshold should be disabled. Time threshold overhead should
1847        // stay the same.
1848        assert_eq!(r.pkt_thresh(), None);
1849        assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1850
1851        // Set now to send time of packet 2 so we can trigger spurious loss for
1852        // packet 2.
1853        now += between_thresh_ms;
1854        // Then wait for `between_thresh_ms` after sending packet 2 to trigger
1855        // loss. Since the time threshold has NOT increased, expect a
1856        // loss.
1857        now += between_thresh_ms;
1858
1859        // Ack packet: 3
1860        //
1861        // [2, 3, 4, 5]
1862        //     ^
1863        let mut acked = RangeSet::default();
1864        acked.insert(3..4);
1865        assert_eq!(
1866            r.on_ack_received(
1867                &acked,
1868                25,
1869                packet::Epoch::Application,
1870                HandshakeStatus::default(),
1871                now,
1872                None,
1873                "",
1874            )
1875            .unwrap(),
1876            OnAckReceivedOutcome {
1877                lost_packets: 1,
1878                lost_bytes: 1000,
1879                acked_bytes: 1000,
1880                spurious_losses: 0,
1881            }
1882        );
1883        // Thresholds after 2nd loss.
1884        assert_eq!(r.pkt_thresh(), None);
1885        assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1886
1887        // Wait for and additional `plus_overhead` to trigger loss based on the
1888        // new time threshold.
1889        // now += THRESH_GAP;
1890
1891        // Ack packet: 2
1892        //
1893        // [2, 3, 4, 5]
1894        //  ^  x
1895        let mut acked = RangeSet::default();
1896        acked.insert(2..3);
1897        assert_eq!(
1898            r.on_ack_received(
1899                &acked,
1900                25,
1901                packet::Epoch::Application,
1902                HandshakeStatus::default(),
1903                now,
1904                None,
1905                "",
1906            )
1907            .unwrap(),
1908            OnAckReceivedOutcome {
1909                lost_packets: 0,
1910                lost_bytes: 0,
1911                acked_bytes: 0,
1912                spurious_losses: 1,
1913            }
1914        );
1915        // Thresholds after 2nd spurious loss.
1916        //
1917        // Time threshold overhead should double.
1918        assert_eq!(r.pkt_thresh(), None);
1919        let double_time_thresh_overhead =
1920            1.0 + 2.0 * INITIAL_TIME_THRESHOLD_OVERHEAD;
1921        assert_eq!(r.time_thresh(), double_time_thresh_overhead);
1922    }
1923
1924    #[rstest]
1925    fn pacing(
1926        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1927        #[values(false, true)] time_sent_set_to_now: bool,
1928    ) {
1929        let pacing_enabled = cc_algorithm_name == "bbr2" ||
1930            cc_algorithm_name == "bbr2_gcongestion";
1931
1932        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
1933        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
1934
1935        #[cfg(feature = "internal")]
1936        cfg.set_custom_bbr_params(BbrParams {
1937            time_sent_set_to_now: Some(time_sent_set_to_now),
1938            ..Default::default()
1939        });
1940
1941        let mut r = Recovery::new(&cfg);
1942
1943        let mut now = Instant::now();
1944
1945        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
1946
1947        // send out first packet burst (a full initcwnd).
1948        for i in 0..10 {
1949            let p = Sent {
1950                pkt_num: i,
1951                frames: smallvec![],
1952                time_sent: now,
1953                time_acked: None,
1954                time_lost: None,
1955                size: 1200,
1956                ack_eliciting: true,
1957                in_flight: true,
1958                delivered: 0,
1959                delivered_time: now,
1960                first_sent_time: now,
1961                is_app_limited: false,
1962                tx_in_flight: 0,
1963                lost: 0,
1964                has_data: true,
1965                is_pmtud_probe: false,
1966            };
1967
1968            r.on_packet_sent(
1969                p,
1970                packet::Epoch::Application,
1971                HandshakeStatus::default(),
1972                now,
1973                "",
1974            );
1975        }
1976
1977        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 10);
1978        assert_eq!(r.bytes_in_flight(), 12000);
1979        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
1980
1981        if !pacing_enabled {
1982            assert_eq!(r.pacing_rate(), 0);
1983        } else {
1984            assert_eq!(r.pacing_rate(), 103963);
1985        }
1986        assert_eq!(r.get_packet_send_time(now), now);
1987
1988        assert_eq!(r.cwnd(), 12000);
1989        assert_eq!(r.cwnd_available(), 0);
1990
1991        // Wait 50ms for ACK.
1992        let initial_rtt = Duration::from_millis(50);
1993        now += initial_rtt;
1994
1995        let mut acked = RangeSet::default();
1996        acked.insert(0..10);
1997
1998        assert_eq!(
1999            r.on_ack_received(
2000                &acked,
2001                10,
2002                packet::Epoch::Application,
2003                HandshakeStatus::default(),
2004                now,
2005                None,
2006                "",
2007            )
2008            .unwrap(),
2009            OnAckReceivedOutcome {
2010                lost_packets: 0,
2011                lost_bytes: 0,
2012                acked_bytes: 12000,
2013                spurious_losses: 0,
2014            }
2015        );
2016
2017        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
2018        assert_eq!(r.bytes_in_flight(), 0);
2019        assert_eq!(r.bytes_in_flight_duration(), initial_rtt);
2020        assert_eq!(r.min_rtt(), Some(initial_rtt));
2021        assert_eq!(r.rtt(), initial_rtt);
2022
2023        // 10 MSS increased due to acks.
2024        assert_eq!(r.cwnd(), 12000 + 1200 * 10);
2025
2026        // Send the second packet burst.
2027        let p = Sent {
2028            pkt_num: 10,
2029            frames: smallvec![],
2030            time_sent: now,
2031            time_acked: None,
2032            time_lost: None,
2033            size: 6000,
2034            ack_eliciting: true,
2035            in_flight: true,
2036            delivered: 0,
2037            delivered_time: now,
2038            first_sent_time: now,
2039            is_app_limited: false,
2040            tx_in_flight: 0,
2041            lost: 0,
2042            has_data: true,
2043            is_pmtud_probe: false,
2044        };
2045
2046        r.on_packet_sent(
2047            p,
2048            packet::Epoch::Application,
2049            HandshakeStatus::default(),
2050            now,
2051            "",
2052        );
2053
2054        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 1);
2055        assert_eq!(r.bytes_in_flight(), 6000);
2056        assert_eq!(r.bytes_in_flight_duration(), initial_rtt);
2057
2058        if !pacing_enabled {
2059            // Pacing is disabled.
2060            assert_eq!(r.get_packet_send_time(now), now);
2061        } else {
2062            // Pacing is done from the beginning.
2063            assert_ne!(r.get_packet_send_time(now), now);
2064        }
2065
2066        // Send the third and fourth packet bursts together.
2067        let p = Sent {
2068            pkt_num: 11,
2069            frames: smallvec![],
2070            time_sent: now,
2071            time_acked: None,
2072            time_lost: None,
2073            size: 6000,
2074            ack_eliciting: true,
2075            in_flight: true,
2076            delivered: 0,
2077            delivered_time: now,
2078            first_sent_time: now,
2079            is_app_limited: false,
2080            tx_in_flight: 0,
2081            lost: 0,
2082            has_data: true,
2083            is_pmtud_probe: false,
2084        };
2085
2086        r.on_packet_sent(
2087            p,
2088            packet::Epoch::Application,
2089            HandshakeStatus::default(),
2090            now,
2091            "",
2092        );
2093
2094        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
2095        assert_eq!(r.bytes_in_flight(), 12000);
2096        assert_eq!(r.bytes_in_flight_duration(), initial_rtt);
2097
2098        // Send the fourth packet burst.
2099        let p = Sent {
2100            pkt_num: 12,
2101            frames: smallvec![],
2102            time_sent: now,
2103            time_acked: None,
2104            time_lost: None,
2105            size: 1000,
2106            ack_eliciting: true,
2107            in_flight: true,
2108            delivered: 0,
2109            delivered_time: now,
2110            first_sent_time: now,
2111            is_app_limited: false,
2112            tx_in_flight: 0,
2113            lost: 0,
2114            has_data: true,
2115            is_pmtud_probe: false,
2116        };
2117
2118        r.on_packet_sent(
2119            p,
2120            packet::Epoch::Application,
2121            HandshakeStatus::default(),
2122            now,
2123            "",
2124        );
2125
2126        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 3);
2127        assert_eq!(r.bytes_in_flight(), 13000);
2128        assert_eq!(r.bytes_in_flight_duration(), initial_rtt);
2129
2130        // We pace this outgoing packet. as all conditions for pacing
2131        // are passed.
2132        let pacing_rate = if pacing_enabled {
2133            let cwnd_gain: f64 = 2.0;
2134            // Adjust for cwnd_gain.  BW estimate was made before the CWND
2135            // increase.
2136            let bw = r.cwnd() as f64 / cwnd_gain / initial_rtt.as_secs_f64();
2137            bw as u64
2138        } else {
2139            0
2140        };
2141        assert_eq!(r.pacing_rate(), pacing_rate);
2142
2143        let scale_factor = if pacing_enabled {
2144            // For bbr2_gcongestion, send time is almost 13000 / pacing_rate.
2145            // Don't know where 13000 comes from.
2146            1.08333332
2147        } else {
2148            1.0
2149        };
2150        assert_eq!(
2151            r.get_packet_send_time(now) - now,
2152            if pacing_enabled {
2153                Duration::from_secs_f64(
2154                    scale_factor * 12000.0 / pacing_rate as f64,
2155                )
2156            } else {
2157                Duration::ZERO
2158            }
2159        );
2160        assert_eq!(r.startup_exit(), None);
2161
2162        let reduced_rtt = Duration::from_millis(40);
2163        now += reduced_rtt;
2164
2165        let mut acked = RangeSet::default();
2166        acked.insert(10..11);
2167
2168        assert_eq!(
2169            r.on_ack_received(
2170                &acked,
2171                0,
2172                packet::Epoch::Application,
2173                HandshakeStatus::default(),
2174                now,
2175                None,
2176                "",
2177            )
2178            .unwrap(),
2179            OnAckReceivedOutcome {
2180                lost_packets: 0,
2181                lost_bytes: 0,
2182                acked_bytes: 6000,
2183                spurious_losses: 0,
2184            }
2185        );
2186
2187        let expected_srtt = (7 * initial_rtt + reduced_rtt) / 8;
2188        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
2189        assert_eq!(r.bytes_in_flight(), 7000);
2190        assert_eq!(r.bytes_in_flight_duration(), initial_rtt + reduced_rtt);
2191        assert_eq!(r.min_rtt(), Some(reduced_rtt));
2192        assert_eq!(r.rtt(), expected_srtt);
2193
2194        let mut acked = RangeSet::default();
2195        acked.insert(11..12);
2196
2197        assert_eq!(
2198            r.on_ack_received(
2199                &acked,
2200                0,
2201                packet::Epoch::Application,
2202                HandshakeStatus::default(),
2203                now,
2204                None,
2205                "",
2206            )
2207            .unwrap(),
2208            OnAckReceivedOutcome {
2209                lost_packets: 0,
2210                lost_bytes: 0,
2211                acked_bytes: 6000,
2212                spurious_losses: 0,
2213            }
2214        );
2215
2216        // When enabled, the pacer adds a 25msec delay to the packet
2217        // sends which will be applied to the sent times tracked by
2218        // the recovery module, bringing down RTT to 15msec.
2219        let expected_min_rtt = if pacing_enabled &&
2220            !time_sent_set_to_now &&
2221            cfg!(feature = "internal")
2222        {
2223            reduced_rtt - Duration::from_millis(25)
2224        } else {
2225            reduced_rtt
2226        };
2227
2228        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 1);
2229        assert_eq!(r.bytes_in_flight(), 1000);
2230        assert_eq!(r.bytes_in_flight_duration(), initial_rtt + reduced_rtt);
2231        assert_eq!(r.min_rtt(), Some(expected_min_rtt));
2232
2233        let expected_srtt = (7 * expected_srtt + expected_min_rtt) / 8;
2234        assert_eq!(r.rtt(), expected_srtt);
2235
2236        let mut acked = RangeSet::default();
2237        acked.insert(12..13);
2238
2239        assert_eq!(
2240            r.on_ack_received(
2241                &acked,
2242                0,
2243                packet::Epoch::Application,
2244                HandshakeStatus::default(),
2245                now,
2246                None,
2247                "",
2248            )
2249            .unwrap(),
2250            OnAckReceivedOutcome {
2251                lost_packets: 0,
2252                lost_bytes: 0,
2253                acked_bytes: 1000,
2254                spurious_losses: 0,
2255            }
2256        );
2257
2258        // Pacer adds 50msec delay to the second packet, resulting in
2259        // an effective RTT of 0.
2260        let expected_min_rtt = if pacing_enabled &&
2261            !time_sent_set_to_now &&
2262            cfg!(feature = "internal")
2263        {
2264            Duration::from_millis(0)
2265        } else {
2266            reduced_rtt
2267        };
2268        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
2269        assert_eq!(r.bytes_in_flight(), 0);
2270        assert_eq!(r.bytes_in_flight_duration(), initial_rtt + reduced_rtt);
2271        assert_eq!(r.min_rtt(), Some(expected_min_rtt));
2272
2273        let expected_srtt = (7 * expected_srtt + expected_min_rtt) / 8;
2274        assert_eq!(r.rtt(), expected_srtt);
2275    }
2276
2277    #[rstest]
2278    // initial_cwnd / first_rtt == initial_pacing_rate.  Pacing is 1.0 * bw before
2279    // and after.
2280    #[case::bw_estimate_equal_after_first_rtt(1.0, 1.0)]
2281    // initial_cwnd / first_rtt < initial_pacing_rate.  Pacing decreases from 2 *
2282    // bw to 1.0 * bw.
2283    #[case::bw_estimate_decrease_after_first_rtt(2.0, 1.0)]
2284    // initial_cwnd / first_rtt > initial_pacing_rate from 0.5 * bw to 1.0 * bw.
2285    // Initial pacing remains 0.5 * bw because the initial_pacing_rate parameter
2286    // is used an upper bound for the pacing rate after the first RTT.
2287    // Pacing rate after the first ACK should be:
2288    // min(initial_pacing_rate_bytes_per_second, init_cwnd / first_rtt)
2289    #[case::bw_estimate_increase_after_first_rtt(0.5, 0.5)]
2290    #[cfg(feature = "internal")]
2291    fn initial_pacing_rate_override(
2292        #[case] initial_multipler: f64, #[case] expected_multiplier: f64,
2293    ) {
2294        let rtt = Duration::from_millis(50);
2295        let bw = Bandwidth::from_bytes_and_time_delta(12000, rtt);
2296        let initial_pacing_rate_hint = bw * initial_multipler;
2297        let expected_pacing_with_rtt_measurement = bw * expected_multiplier;
2298
2299        let cc_algorithm_name = "bbr2_gcongestion";
2300        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
2301        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
2302        cfg.set_custom_bbr_params(BbrParams {
2303            initial_pacing_rate_bytes_per_second: Some(
2304                initial_pacing_rate_hint.to_bytes_per_second(),
2305            ),
2306            ..Default::default()
2307        });
2308
2309        let mut r = Recovery::new(&cfg);
2310
2311        let mut now = Instant::now();
2312
2313        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
2314
2315        // send some packets.
2316        for i in 0..2 {
2317            let p = test_utils::helper_packet_sent(i, now, 1200);
2318            r.on_packet_sent(
2319                p,
2320                packet::Epoch::Application,
2321                HandshakeStatus::default(),
2322                now,
2323                "",
2324            );
2325        }
2326
2327        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
2328        assert_eq!(r.bytes_in_flight(), 2400);
2329        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
2330
2331        // Initial pacing rate matches the override value.
2332        assert_eq!(
2333            r.pacing_rate(),
2334            initial_pacing_rate_hint.to_bytes_per_second()
2335        );
2336        assert_eq!(r.get_packet_send_time(now), now);
2337
2338        assert_eq!(r.cwnd(), 12000);
2339        assert_eq!(r.cwnd_available(), 9600);
2340
2341        // Wait 1 rtt for ACK.
2342        now += rtt;
2343
2344        let mut acked = RangeSet::default();
2345        acked.insert(0..2);
2346
2347        assert_eq!(
2348            r.on_ack_received(
2349                &acked,
2350                10,
2351                packet::Epoch::Application,
2352                HandshakeStatus::default(),
2353                now,
2354                None,
2355                "",
2356            )
2357            .unwrap(),
2358            OnAckReceivedOutcome {
2359                lost_packets: 0,
2360                lost_bytes: 0,
2361                acked_bytes: 2400,
2362                spurious_losses: 0,
2363            }
2364        );
2365
2366        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
2367        assert_eq!(r.bytes_in_flight(), 0);
2368        assert_eq!(r.bytes_in_flight_duration(), rtt);
2369        assert_eq!(r.rtt(), rtt);
2370
2371        // Pacing rate is recalculated based on initial cwnd when the
2372        // first RTT estimate is available.
2373        assert_eq!(
2374            r.pacing_rate(),
2375            expected_pacing_with_rtt_measurement.to_bytes_per_second()
2376        );
2377    }
2378
2379    #[rstest]
2380    fn validate_ack_range_on_ack_received(
2381        #[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
2382    ) {
2383        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
2384        cfg.set_cc_algorithm_name(cc_algorithm_name).unwrap();
2385
2386        let epoch = packet::Epoch::Application;
2387        let mut r = Recovery::new(&cfg);
2388        let mut now = Instant::now();
2389        assert_eq!(r.sent_packets_len(epoch), 0);
2390
2391        // Send 4 packets
2392        let pkt_size = 1000;
2393        let pkt_count = 4;
2394        for pkt_num in 0..pkt_count {
2395            let sent = test_utils::helper_packet_sent(pkt_num, now, pkt_size);
2396            r.on_packet_sent(sent, epoch, HandshakeStatus::default(), now, "");
2397        }
2398        assert_eq!(r.sent_packets_len(epoch), pkt_count as usize);
2399        assert_eq!(r.bytes_in_flight(), pkt_count as usize * pkt_size);
2400        assert!(r.get_largest_acked_on_epoch(epoch).is_none());
2401        assert_eq!(r.largest_sent_pkt_num_on_path(epoch).unwrap(), 3);
2402
2403        // Wait for 10ms.
2404        now += Duration::from_millis(10);
2405
2406        // ACK 2 packets
2407        let mut acked = RangeSet::default();
2408        acked.insert(0..2);
2409
2410        assert_eq!(
2411            r.on_ack_received(
2412                &acked,
2413                25,
2414                epoch,
2415                HandshakeStatus::default(),
2416                now,
2417                None,
2418                "",
2419            )
2420            .unwrap(),
2421            OnAckReceivedOutcome {
2422                lost_packets: 0,
2423                lost_bytes: 0,
2424                acked_bytes: 2 * 1000,
2425                spurious_losses: 0,
2426            }
2427        );
2428
2429        assert_eq!(r.sent_packets_len(epoch), 2);
2430        assert_eq!(r.bytes_in_flight(), 2 * 1000);
2431
2432        assert_eq!(r.get_largest_acked_on_epoch(epoch).unwrap(), 1);
2433        assert_eq!(r.largest_sent_pkt_num_on_path(epoch).unwrap(), 3);
2434
2435        // ACK large range
2436        let mut acked = RangeSet::default();
2437        acked.insert(0..10);
2438        assert_eq!(
2439            r.on_ack_received(
2440                &acked,
2441                25,
2442                epoch,
2443                HandshakeStatus::default(),
2444                now,
2445                None,
2446                "",
2447            )
2448            .unwrap(),
2449            OnAckReceivedOutcome {
2450                lost_packets: 0,
2451                lost_bytes: 0,
2452                acked_bytes: 2 * 1000,
2453                spurious_losses: 0,
2454            }
2455        );
2456        assert_eq!(r.sent_packets_len(epoch), 0);
2457        assert_eq!(r.bytes_in_flight(), 0);
2458
2459        assert_eq!(r.get_largest_acked_on_epoch(epoch).unwrap(), 3);
2460        assert_eq!(r.largest_sent_pkt_num_on_path(epoch).unwrap(), 3);
2461    }
2462
2463    #[rstest]
2464    fn pmtud_loss_on_timer(
2465        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
2466    ) {
2467        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
2468        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
2469
2470        let mut r = Recovery::new(&cfg);
2471        assert_eq!(r.cwnd(), 12000);
2472
2473        let mut now = Instant::now();
2474
2475        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
2476
2477        // Start by sending a few packets.
2478        let p = Sent {
2479            pkt_num: 0,
2480            frames: smallvec![],
2481            time_sent: now,
2482            time_acked: None,
2483            time_lost: None,
2484            size: 1000,
2485            ack_eliciting: true,
2486            in_flight: true,
2487            delivered: 0,
2488            delivered_time: now,
2489            first_sent_time: now,
2490            is_app_limited: false,
2491            tx_in_flight: 0,
2492            lost: 0,
2493            has_data: false,
2494            is_pmtud_probe: false,
2495        };
2496
2497        r.on_packet_sent(
2498            p,
2499            packet::Epoch::Application,
2500            HandshakeStatus::default(),
2501            now,
2502            "",
2503        );
2504
2505        assert_eq!(r.in_flight_count(packet::Epoch::Application), 1);
2506        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 1);
2507        assert_eq!(r.bytes_in_flight(), 1000);
2508        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
2509
2510        let p = Sent {
2511            pkt_num: 1,
2512            frames: smallvec![],
2513            time_sent: now,
2514            time_acked: None,
2515            time_lost: None,
2516            size: 1000,
2517            ack_eliciting: true,
2518            in_flight: true,
2519            delivered: 0,
2520            delivered_time: now,
2521            first_sent_time: now,
2522            is_app_limited: false,
2523            tx_in_flight: 0,
2524            lost: 0,
2525            has_data: false,
2526            is_pmtud_probe: true,
2527        };
2528
2529        r.on_packet_sent(
2530            p,
2531            packet::Epoch::Application,
2532            HandshakeStatus::default(),
2533            now,
2534            "",
2535        );
2536
2537        assert_eq!(r.in_flight_count(packet::Epoch::Application), 2);
2538
2539        let p = Sent {
2540            pkt_num: 2,
2541            frames: smallvec![],
2542            time_sent: now,
2543            time_acked: None,
2544            time_lost: None,
2545            size: 1000,
2546            ack_eliciting: true,
2547            in_flight: true,
2548            delivered: 0,
2549            delivered_time: now,
2550            first_sent_time: now,
2551            is_app_limited: false,
2552            tx_in_flight: 0,
2553            lost: 0,
2554            has_data: false,
2555            is_pmtud_probe: false,
2556        };
2557
2558        r.on_packet_sent(
2559            p,
2560            packet::Epoch::Application,
2561            HandshakeStatus::default(),
2562            now,
2563            "",
2564        );
2565
2566        assert_eq!(r.in_flight_count(packet::Epoch::Application), 3);
2567
2568        // Wait for 10ms.
2569        now += Duration::from_millis(10);
2570
2571        // Only the first  packets and the last one are acked.
2572        let mut acked = RangeSet::default();
2573        acked.insert(0..1);
2574        acked.insert(2..3);
2575
2576        assert_eq!(
2577            r.on_ack_received(
2578                &acked,
2579                25,
2580                packet::Epoch::Application,
2581                HandshakeStatus::default(),
2582                now,
2583                None,
2584                "",
2585            )
2586            .unwrap(),
2587            OnAckReceivedOutcome {
2588                lost_packets: 0,
2589                lost_bytes: 0,
2590                acked_bytes: 2 * 1000,
2591                spurious_losses: 0,
2592            }
2593        );
2594
2595        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
2596        assert_eq!(r.bytes_in_flight(), 1000);
2597        assert_eq!(r.bytes_in_flight_duration(), Duration::from_millis(10));
2598        assert_eq!(r.lost_count(), 0);
2599
2600        // Wait until loss detection timer expires.
2601        now = r.loss_detection_timer().unwrap();
2602
2603        // Packet is declared lost.
2604        r.on_loss_detection_timeout(HandshakeStatus::default(), now, "");
2605        assert_eq!(r.loss_probes(packet::Epoch::Application), 0);
2606
2607        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 2);
2608        assert_eq!(r.in_flight_count(packet::Epoch::Application), 0);
2609        assert_eq!(r.bytes_in_flight(), 0);
2610        assert_eq!(r.bytes_in_flight_duration(), Duration::from_micros(11250));
2611        assert_eq!(r.cwnd(), 12000);
2612
2613        assert_eq!(r.lost_count(), 0);
2614
2615        // Wait 1 RTT.
2616        now += r.rtt();
2617
2618        assert_eq!(
2619            r.detect_lost_packets_for_test(packet::Epoch::Application, now),
2620            (0, 0)
2621        );
2622
2623        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
2624        assert_eq!(r.in_flight_count(packet::Epoch::Application), 0);
2625        assert_eq!(r.bytes_in_flight(), 0);
2626        assert_eq!(r.bytes_in_flight_duration(), Duration::from_micros(11250));
2627        assert_eq!(r.lost_count(), 0);
2628        assert_eq!(r.startup_exit(), None);
2629    }
2630
2631    // Modeling delivery_rate for gcongestion is non-trivial so we only test the
2632    // congestion specific algorithms.
2633    #[rstest]
2634    fn congestion_delivery_rate(
2635        #[values("reno", "cubic", "bbr2")] cc_algorithm_name: &str,
2636    ) {
2637        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
2638        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
2639
2640        let mut r = Recovery::new(&cfg);
2641        assert_eq!(r.cwnd(), 12000);
2642
2643        let now = Instant::now();
2644
2645        let mut total_bytes_sent = 0;
2646        for pn in 0..10 {
2647            // Start by sending a few packets.
2648            let bytes = 1000;
2649            let sent = test_utils::helper_packet_sent(pn, now, bytes);
2650            r.on_packet_sent(
2651                sent,
2652                packet::Epoch::Application,
2653                HandshakeStatus::default(),
2654                now,
2655                "",
2656            );
2657
2658            total_bytes_sent += bytes;
2659        }
2660
2661        // Ack
2662        let interval = Duration::from_secs(10);
2663        let mut acked = RangeSet::default();
2664        acked.insert(0..10);
2665        assert_eq!(
2666            r.on_ack_received(
2667                &acked,
2668                25,
2669                packet::Epoch::Application,
2670                HandshakeStatus::default(),
2671                now + interval,
2672                None,
2673                "",
2674            )
2675            .unwrap(),
2676            OnAckReceivedOutcome {
2677                lost_packets: 0,
2678                lost_bytes: 0,
2679                acked_bytes: total_bytes_sent,
2680                spurious_losses: 0,
2681            }
2682        );
2683        assert_eq!(r.delivery_rate().to_bytes_per_second(), 1000);
2684        assert_eq!(r.min_rtt().unwrap(), interval);
2685        // delivery rate should be in units bytes/sec
2686        assert_eq!(
2687            total_bytes_sent as u64 / interval.as_secs(),
2688            r.delivery_rate().to_bytes_per_second()
2689        );
2690        assert_eq!(r.startup_exit(), None);
2691    }
2692
2693    #[rstest]
2694    fn acks_with_no_retransmittable_data(
2695        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
2696    ) {
2697        let rtt = Duration::from_millis(100);
2698
2699        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
2700        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
2701
2702        let mut r = Recovery::new(&cfg);
2703
2704        let mut now = Instant::now();
2705
2706        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
2707
2708        let mut next_packet = 0;
2709        // send some packets.
2710        for _ in 0..3 {
2711            let p = test_utils::helper_packet_sent(next_packet, now, 1200);
2712            next_packet += 1;
2713            r.on_packet_sent(
2714                p,
2715                packet::Epoch::Application,
2716                HandshakeStatus::default(),
2717                now,
2718                "",
2719            );
2720        }
2721
2722        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 3);
2723        assert_eq!(r.bytes_in_flight(), 3600);
2724        assert_eq!(r.bytes_in_flight_duration(), Duration::ZERO);
2725
2726        assert_eq!(
2727            r.pacing_rate(),
2728            if cc_algorithm_name == "bbr2_gcongestion" {
2729                103963
2730            } else {
2731                0
2732            },
2733        );
2734        assert_eq!(r.get_packet_send_time(now), now);
2735        assert_eq!(r.cwnd(), 12000);
2736        assert_eq!(r.cwnd_available(), 8400);
2737
2738        // Wait 1 rtt for ACK.
2739        now += rtt;
2740
2741        let mut acked = RangeSet::default();
2742        acked.insert(0..3);
2743
2744        assert_eq!(
2745            r.on_ack_received(
2746                &acked,
2747                10,
2748                packet::Epoch::Application,
2749                HandshakeStatus::default(),
2750                now,
2751                None,
2752                "",
2753            )
2754            .unwrap(),
2755            OnAckReceivedOutcome {
2756                lost_packets: 0,
2757                lost_bytes: 0,
2758                acked_bytes: 3600,
2759                spurious_losses: 0,
2760            }
2761        );
2762
2763        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 0);
2764        assert_eq!(r.bytes_in_flight(), 0);
2765        assert_eq!(r.bytes_in_flight_duration(), rtt);
2766        assert_eq!(r.rtt(), rtt);
2767
2768        // Pacing rate is recalculated based on initial cwnd when the
2769        // first RTT estimate is available.
2770        assert_eq!(
2771            r.pacing_rate(),
2772            if cc_algorithm_name == "bbr2_gcongestion" {
2773                120000
2774            } else {
2775                0
2776            },
2777        );
2778
2779        // Send some no "in_flight" packets
2780        for iter in 3..1000 {
2781            let mut p = test_utils::helper_packet_sent(next_packet, now, 1200);
2782            // `in_flight = false` marks packets as if they only contained ACK
2783            // frames.
2784            p.in_flight = false;
2785            next_packet += 1;
2786            r.on_packet_sent(
2787                p,
2788                packet::Epoch::Application,
2789                HandshakeStatus::default(),
2790                now,
2791                "",
2792            );
2793
2794            now += rtt;
2795
2796            let mut acked = RangeSet::default();
2797            acked.insert(iter..(iter + 1));
2798
2799            assert_eq!(
2800                r.on_ack_received(
2801                    &acked,
2802                    10,
2803                    packet::Epoch::Application,
2804                    HandshakeStatus::default(),
2805                    now,
2806                    None,
2807                    "",
2808                )
2809                .unwrap(),
2810                OnAckReceivedOutcome {
2811                    lost_packets: 0,
2812                    lost_bytes: 0,
2813                    acked_bytes: 0,
2814                    spurious_losses: 0,
2815                }
2816            );
2817
2818            // Verify that connection has not exited startup.
2819            assert_eq!(r.startup_exit(), None, "{iter}");
2820
2821            // Unchanged metrics.
2822            assert_eq!(
2823                r.sent_packets_len(packet::Epoch::Application),
2824                0,
2825                "{iter}"
2826            );
2827            assert_eq!(r.bytes_in_flight(), 0, "{iter}");
2828            assert_eq!(r.bytes_in_flight_duration(), rtt, "{iter}");
2829            assert_eq!(
2830                r.pacing_rate(),
2831                if cc_algorithm_name == "bbr2_gcongestion" ||
2832                    cc_algorithm_name == "bbr2"
2833                {
2834                    120000
2835                } else {
2836                    0
2837                },
2838                "{iter}"
2839            );
2840        }
2841    }
2842    #[rstest]
2843    fn pto_overflow_reproduction(
2844        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
2845    ) {
2846        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
2847        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
2848        let mut r = Recovery::new(&cfg);
2849        let now = Instant::now();
2850
2851        // Scenario: Handshake not completed
2852        let handshake_status = HandshakeStatus {
2853            has_handshake_keys: true,
2854            peer_verified_address: true,
2855            completed: false,
2856        };
2857
2858        // 1. Send Initial packet to arm the timer
2859        let p_initial = Sent {
2860            pkt_num: 0,
2861            frames: smallvec::smallvec![],
2862            time_sent: now,
2863            time_acked: None,
2864            time_lost: None,
2865            size: 1000,
2866            ack_eliciting: true,
2867            in_flight: true,
2868            delivered: 0,
2869            delivered_time: now,
2870            first_sent_time: now,
2871            is_app_limited: false,
2872            tx_in_flight: 0,
2873            lost: 0,
2874            has_data: false,
2875            is_pmtud_probe: false,
2876        };
2877        r.on_packet_sent(
2878            p_initial,
2879            packet::Epoch::Initial,
2880            handshake_status,
2881            now,
2882            "",
2883        );
2884
2885        // The timer should now be set for the Initial packet.
2886        assert!(r.loss_detection_timer().is_some());
2887
2888        // 2. Send Application packet (0-RTT)
2889        let p_app = Sent {
2890            pkt_num: 0, // Application space has its own packet numbers
2891            frames: smallvec::smallvec![],
2892            time_sent: now,
2893            time_acked: None,
2894            time_lost: None,
2895            size: 1000,
2896            ack_eliciting: true,
2897            in_flight: true,
2898            delivered: 0,
2899            delivered_time: now,
2900            first_sent_time: now,
2901            is_app_limited: false,
2902            tx_in_flight: 0,
2903
2904            lost: 0,
2905            has_data: true,
2906            is_pmtud_probe: false,
2907        };
2908        r.on_packet_sent(
2909            p_app,
2910            packet::Epoch::Application,
2911            handshake_status,
2912            now,
2913            "",
2914        );
2915
2916        // 3. Acknowledge the Initial packet.
2917        // This empties the Initial space, but Application space still has data
2918        // in flight.
2919        let mut ranges = RangeSet::default();
2920        ranges.insert(0..1);
2921        r.on_ack_received(
2922            &ranges,
2923            0,
2924            packet::Epoch::Initial,
2925            handshake_status,
2926            now,
2927            None,
2928            "",
2929        )
2930        .unwrap();
2931
2932        // The timer should be cleared at this point.
2933        // Although there is Application data in flight, the handshake is not
2934        // confirmed, so it cannot be used to arm the PTO timer. Since there are
2935        // no packets in flight in Initial or Handshake spaces either, no timer
2936        // should be set.
2937        assert!(r.loss_detection_timer().is_none());
2938    }
2939
2940    // Test that consecutive PTOs don't add duplicate frames to lost_frames.
2941    // This validates the fix: `if epoch.lost_frames.is_empty()` guard.
2942    #[rstest]
2943    fn pto_does_not_duplicate_frames_on_consecutive_timeouts(
2944        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
2945    ) {
2946        let mut cfg = Config::new(crate::PROTOCOL_VERSION).unwrap();
2947        assert_eq!(cfg.set_cc_algorithm_name(cc_algorithm_name), Ok(()));
2948
2949        let mut r = Recovery::new(&cfg);
2950        let mut now = Instant::now();
2951
2952        // Send a packet with a STREAM frame.
2953        let frames = smallvec![frame::Frame::Stream {
2954            stream_id: 4,
2955            data: RangeBuf::from(b"test", 0, false),
2956        },];
2957
2958        let p = Sent {
2959            pkt_num: 0,
2960            frames: frames.clone(),
2961            time_sent: now,
2962            time_acked: None,
2963            time_lost: None,
2964            size: 1000,
2965            ack_eliciting: true,
2966            in_flight: true,
2967            delivered: 0,
2968            delivered_time: now,
2969            first_sent_time: now,
2970            is_app_limited: false,
2971            tx_in_flight: 0,
2972            lost: 0,
2973            has_data: true,
2974            is_pmtud_probe: false,
2975        };
2976
2977        r.on_packet_sent(
2978            p,
2979            packet::Epoch::Application,
2980            HandshakeStatus::default(),
2981            now,
2982            "",
2983        );
2984
2985        // Verify initial state
2986        assert_eq!(r.lost_frames_count(packet::Epoch::Application), 0);
2987        assert_eq!(r.lost_count(), 0);
2988
2989        // First PTO - should add frames when lost_frames is empty.
2990        now = r.loss_detection_timer().unwrap();
2991        r.on_loss_detection_timeout(HandshakeStatus::default(), now, "");
2992
2993        assert_eq!(r.pto_count(), 1);
2994        let frames_after_first_pto =
2995            r.lost_frames_count(packet::Epoch::Application);
2996        assert_eq!(
2997            frames_after_first_pto, 1,
2998            "First PTO should add exactly 1 frame"
2999        );
3000        assert_eq!(
3001            r.lost_count(),
3002            0,
3003            "PTO doesn't declare packets lost (no CC impact)"
3004        );
3005
3006        // Second PTO while lost_frames is still populated.
3007        // WITHOUT the fix: would add duplicate frame (count becomes 2).
3008        // WITH the fix: skips adding (count stays 1).
3009        now = r.loss_detection_timer().unwrap();
3010        r.on_loss_detection_timeout(HandshakeStatus::default(), now, "");
3011
3012        assert_eq!(r.pto_count(), 2);
3013        let frames_after_second_pto =
3014            r.lost_frames_count(packet::Epoch::Application);
3015        assert_eq!(
3016            frames_after_second_pto, frames_after_first_pto,
3017            "Second PTO must NOT add duplicate frames (fix: `if \
3018             lost_frames.is_empty()`)"
3019        );
3020
3021        // Third PTO for extra validation
3022        now = r.loss_detection_timer().unwrap();
3023        r.on_loss_detection_timeout(HandshakeStatus::default(), now, "");
3024
3025        assert_eq!(r.pto_count(), 3);
3026        let frames_after_third_pto =
3027            r.lost_frames_count(packet::Epoch::Application);
3028        assert_eq!(
3029            frames_after_third_pto, frames_after_first_pto,
3030            "Third PTO must NOT add duplicate frames"
3031        );
3032
3033        // Verify packets are still tracked (not removed)
3034        assert_eq!(r.sent_packets_len(packet::Epoch::Application), 1);
3035        // Verify lost_count never increased (PTO doesn't trigger CC)
3036        assert_eq!(r.lost_count(), 0);
3037    }
3038
3039    // Test that send_on_path after PTO timeout properly sends retransmissions
3040    // and doesn't mark packets as lost (lost_count should remain 0).
3041    #[rstest]
3042    fn pto_send_on_path_retransmits_without_loss(
3043        #[values("reno", "cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
3044    ) {
3045        use crate::test_utils;
3046
3047        let mut pipe = test_utils::Pipe::new(cc_algorithm_name).unwrap();
3048
3049        // Complete handshake
3050        assert_eq!(pipe.handshake(), Ok(()));
3051
3052        // Client sends stream data
3053        assert_eq!(pipe.client.stream_send(4, b"hello", false), Ok(5));
3054
3055        let mut buf = [0; 65535];
3056
3057        // Send the packet but drop it (don't deliver to server)
3058        let (len1, _) = pipe.client.send(&mut buf).unwrap();
3059        assert!(len1 > 0);
3060
3061        // Verify lost_count is 0 (no losses yet)
3062        let initial_lost_count = pipe
3063            .client
3064            .paths
3065            .get_active()
3066            .unwrap()
3067            .recovery
3068            .lost_count();
3069        assert_eq!(initial_lost_count, 0, "No packets should be lost initially");
3070
3071        // Verify frames are not yet in lost_frames
3072        let initial_lost_frames = pipe
3073            .client
3074            .paths
3075            .get_active()
3076            .unwrap()
3077            .recovery
3078            .lost_frames_count(packet::Epoch::Application);
3079        assert_eq!(
3080            initial_lost_frames, 0,
3081            "No frames should be in lost_frames initially"
3082        );
3083
3084        // Wait for PTO timeout
3085        let timer = pipe.client.timeout().unwrap();
3086        std::thread::sleep(timer + Duration::from_millis(1));
3087
3088        // Trigger PTO via on_timeout()
3089        pipe.client.on_timeout();
3090
3091        // After PTO, frames should be in lost_frames for retransmission
3092        let lost_frames_after_pto = pipe
3093            .client
3094            .paths
3095            .get_active()
3096            .unwrap()
3097            .recovery
3098            .lost_frames_count(packet::Epoch::Application);
3099        assert!(
3100            lost_frames_after_pto > 0,
3101            "PTO should add frames to lost_frames for retransmission"
3102        );
3103
3104        // But lost_count should still be 0 (PTO doesn't declare packets lost)
3105        let lost_count_after_pto = pipe
3106            .client
3107            .paths
3108            .get_active()
3109            .unwrap()
3110            .recovery
3111            .lost_count();
3112        assert_eq!(
3113            lost_count_after_pto, 0,
3114            "PTO should not increment lost_count"
3115        );
3116
3117        // Now send the retransmission via send_on_path
3118        let (len2, _) = pipe.client.send(&mut buf).unwrap();
3119        assert!(len2 > 0, "Should send PTO probe packet");
3120
3121        // After sending, lost_count should still be 0
3122        let lost_count_after_send = pipe
3123            .client
3124            .paths
3125            .get_active()
3126            .unwrap()
3127            .recovery
3128            .lost_count();
3129        assert_eq!(
3130            lost_count_after_send, 0,
3131            "Sending PTO probe should not increment lost_count"
3132        );
3133
3134        // Deliver the retransmission to server
3135        assert_eq!(pipe.server_recv(&mut buf[..len2]), Ok(len2));
3136
3137        // Server should receive the stream data
3138        let mut recv_buf = [0; 100];
3139        assert_eq!(pipe.server.stream_recv(4, &mut recv_buf), Ok((5, false)));
3140        assert_eq!(&recv_buf[..5], b"hello");
3141
3142        // Final verification: lost_count on client should still be 0
3143        let final_lost_count = pipe
3144            .client
3145            .paths
3146            .get_active()
3147            .unwrap()
3148            .recovery
3149            .lost_count();
3150        assert_eq!(
3151            final_lost_count, 0,
3152            "No packets should be marked as lost - PTO only retransmits"
3153        );
3154    }
3155}
3156
3157mod bandwidth;
3158mod bytes_in_flight;
3159mod congestion;
3160mod gcongestion;
3161mod rtt;