1use 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
52const INITIAL_PACKET_THRESHOLD: u64 = 3;
54
55const MAX_PACKET_THRESHOLD: u64 = 20;
56
57const INITIAL_TIME_THRESHOLD: f64 = 9.0 / 8.0;
61
62const PACKET_REORDER_TIME_THRESHOLD: f64 = 5.0 / 4.0;
74
75const INITIAL_TIME_THRESHOLD_OVERHEAD: f64 = 1.0 / 8.0;
82const 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
95pub(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]
188pub trait RecoveryOps {
190 fn lost_count(&self) -> usize;
191 fn bytes_lost(&self) -> u64;
192
193 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 fn delivery_rate(&self) -> Bandwidth;
250
251 fn max_bandwidth(&self) -> Option<Bandwidth>;
253
254 fn rtt_persistent_jump_count(&self) -> u64;
256
257 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 #[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 #[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#[derive(Debug, Copy, Clone, PartialEq, Eq)]
373#[repr(C)]
374pub enum CongestionControlAlgorithm {
375 Reno = 0,
377 CUBIC = 1,
379 Bbr2Gcongestion = 4,
382}
383
384impl FromStr for CongestionControlAlgorithm {
385 type Err = crate::Error;
386
387 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#[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 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 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
719pub enum ReleaseTime {
720 Immediate,
721 At(Instant),
722}
723
724#[derive(Clone, Copy, Debug, PartialEq, Eq)]
726pub struct ReleaseDecision {
727 time: ReleaseTime,
728 allow_burst: bool,
729}
730
731impl ReleaseTime {
732 fn inc(&mut self, delay: Duration) {
734 match self {
735 ReleaseTime::Immediate => {},
736 ReleaseTime::At(time) => *time += delay,
737 }
738 }
739
740 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 #[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 #[inline]
764 pub fn can_burst(&self) -> bool {
765 self.allow_burst
766 }
767
768 #[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#[derive(Default, Debug)]
784pub struct RecoveryStats {
785 startup_exit: Option<StartupExit>,
786}
787
788impl RecoveryStats {
789 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#[derive(Debug, Clone, Copy, PartialEq)]
799pub struct StartupExit {
800 pub cwnd: usize,
802
803 pub bandwidth: Option<u64>,
805
806 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#[derive(Debug, Clone, Copy, PartialEq)]
825pub enum StartupExitReason {
826 Loss,
828
829 BandwidthPlateau,
831
832 PersistentQueue,
834
835 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 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 now += Duration::from_millis(10);
1094
1095 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 now = r.loss_detection_timer().unwrap();
1125
1126 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 now += Duration::from_millis(10);
1195
1196 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 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 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 now += Duration::from_millis(10);
1378
1379 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 now = r.loss_detection_timer().unwrap();
1410
1411 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 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 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 now += Duration::from_millis(10);
1472
1473 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 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 now += Duration::from_millis(10);
1502
1503 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 assert_eq!(r.lost_count(), 1);
1530 assert_eq!(r.lost_spurious_count(), 1);
1531
1532 assert_eq!(r.pkt_thresh().unwrap(), 4);
1534 assert_eq!(r.time_thresh(), PACKET_REORDER_TIME_THRESHOLD);
1535
1536 now += r.rtt();
1538
1539 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 #[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 const THRESH_GAP: Duration = Duration::from_millis(30);
1583 let initial_thresh_ms =
1585 DEFAULT_INITIAL_RTT.mul_f64(INITIAL_TIME_THRESHOLD);
1586 let spurious_thresh_ms: Duration =
1588 DEFAULT_INITIAL_RTT.mul_f64(PACKET_REORDER_TIME_THRESHOLD);
1589 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 now += between_thresh_ms;
1616
1617 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 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 assert_eq!(r.time_thresh(), PACKET_REORDER_TIME_THRESHOLD);
1670
1671 now += between_thresh_ms;
1674
1675 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 now += THRESH_GAP;
1703
1704 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 #[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 const THRESH_GAP: Duration = Duration::from_millis(30);
1757 let initial_thresh_ms =
1759 DEFAULT_INITIAL_RTT.mul_f64(INITIAL_TIME_THRESHOLD);
1760 let spurious_thresh_ms: Duration =
1762 DEFAULT_INITIAL_RTT.mul_f64(PACKET_REORDER_TIME_THRESHOLD);
1763 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 assert_eq!(r.pkt_thresh().unwrap(), INITIAL_PACKET_THRESHOLD);
1786 assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1787
1788 now += between_thresh_ms;
1791
1792 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 assert_eq!(r.pkt_thresh().unwrap(), INITIAL_PACKET_THRESHOLD);
1818 assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1819
1820 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 assert_eq!(r.pkt_thresh(), None);
1849 assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1850
1851 now += between_thresh_ms;
1854 now += between_thresh_ms;
1858
1859 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 assert_eq!(r.pkt_thresh(), None);
1885 assert_eq!(r.time_thresh(), INITIAL_TIME_THRESHOLD);
1886
1887 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 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 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 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 assert_eq!(r.cwnd(), 12000 + 1200 * 10);
2025
2026 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 assert_eq!(r.get_packet_send_time(now), now);
2061 } else {
2062 assert_ne!(r.get_packet_send_time(now), now);
2064 }
2065
2066 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 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 let pacing_rate = if pacing_enabled {
2133 let cwnd_gain: f64 = 2.0;
2134 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 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 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 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 #[case::bw_estimate_equal_after_first_rtt(1.0, 1.0)]
2281 #[case::bw_estimate_decrease_after_first_rtt(2.0, 1.0)]
2284 #[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 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 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 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 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 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 now += Duration::from_millis(10);
2405
2406 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 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 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 now += Duration::from_millis(10);
2570
2571 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 now = r.loss_detection_timer().unwrap();
2602
2603 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 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 #[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 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 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 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 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 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 assert_eq!(
2771 r.pacing_rate(),
2772 if cc_algorithm_name == "bbr2_gcongestion" {
2773 120000
2774 } else {
2775 0
2776 },
2777 );
2778
2779 for iter in 3..1000 {
2781 let mut p = test_utils::helper_packet_sent(next_packet, now, 1200);
2782 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 assert_eq!(r.startup_exit(), None, "{iter}");
2820
2821 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 let handshake_status = HandshakeStatus {
2853 has_handshake_keys: true,
2854 peer_verified_address: true,
2855 completed: false,
2856 };
2857
2858 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 assert!(r.loss_detection_timer().is_some());
2887
2888 let p_app = Sent {
2890 pkt_num: 0, 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 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 assert!(r.loss_detection_timer().is_none());
2938 }
2939
2940 #[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 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 assert_eq!(r.lost_frames_count(packet::Epoch::Application), 0);
2987 assert_eq!(r.lost_count(), 0);
2988
2989 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 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 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 assert_eq!(r.sent_packets_len(packet::Epoch::Application), 1);
3035 assert_eq!(r.lost_count(), 0);
3037 }
3038
3039 #[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 assert_eq!(pipe.handshake(), Ok(()));
3051
3052 assert_eq!(pipe.client.stream_send(4, b"hello", false), Ok(5));
3054
3055 let mut buf = [0; 65535];
3056
3057 let (len1, _) = pipe.client.send(&mut buf).unwrap();
3059 assert!(len1 > 0);
3060
3061 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 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 let timer = pipe.client.timeout().unwrap();
3086 std::thread::sleep(timer + Duration::from_millis(1));
3087
3088 pipe.client.on_timeout();
3090
3091 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 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 let (len2, _) = pipe.client.send(&mut buf).unwrap();
3119 assert!(len2 > 0, "Should send PTO probe packet");
3120
3121 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 assert_eq!(pipe.server_recv(&mut buf[..len2]), Ok(len2));
3136
3137 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 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;