Skip to main content

quiche/recovery/gcongestion/bbr2/
network_model.rs

1// Copyright (c) 2015 The Chromium Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5// Copyright (C) 2023, Cloudflare, Inc.
6// All rights reserved.
7//
8// Redistribution and use in source and binary forms, with or without
9// modification, are permitted provided that the following conditions are
10// met:
11//
12//     * Redistributions of source code must retain the above copyright notice,
13//       this list of conditions and the following disclaimer.
14//
15//     * Redistributions in binary form must reproduce the above copyright
16//       notice, this list of conditions and the following disclaimer in the
17//       documentation and/or other materials provided with the distribution.
18//
19// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS
20// IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO,
21// THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
22// PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
23// CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
24// EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
25// PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
26// PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
27// LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
28// NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
29// SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
30
31use std::ops::Add;
32use std::time::Duration;
33use std::time::Instant;
34
35use crate::recovery::gcongestion::bbr::BandwidthSampler;
36use crate::recovery::gcongestion::bbr2::Params;
37use crate::recovery::gcongestion::Bandwidth;
38use crate::recovery::gcongestion::Lost;
39
40use super::rtt_jump_detector::RttJumpDetector;
41use super::Acked;
42use super::BBRv2CongestionEvent;
43use super::BwLoMode;
44
45pub(super) const DEFAULT_MSS: usize = 1300;
46
47#[derive(Debug)]
48struct RoundTripCounter {
49    round_trip_count: usize,
50    last_sent_packet: u64,
51    // The last sent packet number of the current round trip.
52    end_of_round_trip: Option<u64>,
53}
54
55impl RoundTripCounter {
56    /// Must be called in ascending packet number order.
57    fn on_packet_sent(&mut self, packet_number: u64) {
58        self.last_sent_packet = packet_number;
59    }
60
61    /// Return whether a round trip has just completed.
62    fn on_packets_acked(&mut self, last_acked_packet: u64) -> bool {
63        match self.end_of_round_trip {
64            Some(pkt) if last_acked_packet <= pkt => false,
65            _ => {
66                self.round_trip_count += 1;
67                self.end_of_round_trip = Some(self.last_sent_packet);
68                true
69            },
70        }
71    }
72
73    fn restart_round(&mut self) {
74        self.end_of_round_trip = Some(self.last_sent_packet)
75    }
76}
77
78#[derive(Debug)]
79struct MinRttFilter {
80    min_rtt: Duration,
81    min_rtt_timestamp: Instant,
82}
83
84impl MinRttFilter {
85    fn get(&self) -> Duration {
86        self.min_rtt
87    }
88
89    fn get_timestamps(&self) -> Instant {
90        self.min_rtt_timestamp
91    }
92
93    fn update(&mut self, sample_rtt: Duration, now: Instant) {
94        if sample_rtt < self.min_rtt {
95            self.min_rtt = sample_rtt;
96            self.min_rtt_timestamp = now;
97        }
98    }
99
100    fn force_update(&mut self, sample_rtt: Duration, now: Instant) {
101        self.min_rtt = sample_rtt;
102        self.min_rtt_timestamp = now;
103    }
104}
105
106#[derive(Debug)]
107struct MaxBandwidthFilter {
108    max_bandwidth: [Bandwidth; 2],
109}
110
111impl MaxBandwidthFilter {
112    fn get(&self) -> Bandwidth {
113        self.max_bandwidth[0].max(self.max_bandwidth[1])
114    }
115
116    fn update(&mut self, sample: Bandwidth) {
117        self.max_bandwidth[1] = self.max_bandwidth[1].max(sample);
118    }
119
120    fn advance(&mut self) {
121        if self.max_bandwidth[1] == Bandwidth::zero() {
122            return;
123        }
124
125        self.max_bandwidth[0] = self.max_bandwidth[1];
126        self.max_bandwidth[1] = Bandwidth::zero();
127    }
128}
129
130/// Bbr2NetworkModel takes low level congestion signals(packets sent/acked/lost)
131/// as input and produces BBRv2 model parameters like inflight_(hi|lo),
132/// bandwidth_(hi|lo), bandwidth and rtt estimates, etc.
133#[derive(Debug)]
134pub(super) struct BBRv2NetworkModel {
135    round_trip_counter: RoundTripCounter,
136    /// Bandwidth sampler provides BBR with the bandwidth measurements at
137    /// individual points.
138    bandwidth_sampler: BandwidthSampler,
139    /// The filter that tracks the maximum bandwidth over multiple recent round
140    /// trips.
141    max_bandwidth_filter: MaxBandwidthFilter,
142    min_rtt_filter: MinRttFilter,
143    /// Bytes lost in the current round. Updated once per congestion event.
144    bytes_lost_in_round: usize,
145    /// Number of loss marking events in the current round.
146    loss_events_in_round: usize,
147
148    /// A max of bytes delivered among all congestion events in the current
149    /// round. A congestions event's bytes delivered is the total bytes
150    /// acked between time Ts and Ta, which is the time when the largest
151    /// acked packet(within the congestion event) was sent and acked,
152    /// respectively.
153    max_bytes_delivered_in_round: usize,
154
155    /// The minimum bytes in flight during this round.
156    min_bytes_in_flight_in_round: usize,
157    /// True if sending was limited by inflight_hi anytime in the current round.
158    inflight_hi_limited_in_round: bool,
159
160    /// Max bandwidth in the current round. Updated once per congestion event.
161    bandwidth_latest: Bandwidth,
162    /// Max bandwidth of recent rounds. Updated once per round.
163    bandwidth_lo: Option<Bandwidth>,
164    prior_bandwidth_lo: Option<Bandwidth>,
165
166    /// Max inflight in the current round. Updated once per congestion event.
167    inflight_latest: usize,
168    /// Max inflight of recent rounds. Updated once per round.
169    inflight_lo: usize,
170    inflight_hi: usize,
171
172    cwnd_gain: f32,
173    pacing_gain: f32,
174
175    /// Whether we are cwnd limited prior to the start of the current
176    /// aggregation epoch.
177    cwnd_limited_before_aggregation_epoch: bool,
178
179    /// STARTUP-centric fields which experimentally used by PROBE_UP.
180    full_bandwidth_reached: bool,
181    full_bandwidth_baseline: Bandwidth,
182    rounds_without_bandwidth_growth: usize,
183
184    /// Used by STARTUP and PROBE_UP to decide when to exit.
185    rounds_with_queueing: usize,
186
187    /// Determines whether app limited rounds with no bandwidth growth count
188    /// towards the rounds threshold to exit startup.
189    ignore_app_limited_for_no_bandwidth_growth: bool,
190
191    /// The most recent send rate from the BandwidthSampler.
192    latest_send_rate: Option<Bandwidth>,
193    /// The most recent ack rate from the BandwidthSampler.
194    latest_ack_rate: Option<Bandwidth>,
195
196    /// Detector for persistent RTT jump episodes.
197    rtt_jump_detector: RttJumpDetector,
198}
199
200impl BBRv2NetworkModel {
201    pub(super) fn new(params: &Params, initial_rtt: Duration) -> Self {
202        BBRv2NetworkModel {
203            min_bytes_in_flight_in_round: usize::MAX,
204            inflight_hi_limited_in_round: false,
205            bandwidth_sampler: BandwidthSampler::new(
206                params.initial_max_ack_height_filter_window,
207                params.enable_overestimate_avoidance,
208                params.choose_a0_point_fix,
209            ),
210            round_trip_counter: RoundTripCounter {
211                round_trip_count: 0,
212                last_sent_packet: 0,
213                end_of_round_trip: None,
214            },
215            min_rtt_filter: MinRttFilter {
216                min_rtt: initial_rtt,
217                min_rtt_timestamp: Instant::now(),
218            },
219            max_bandwidth_filter: MaxBandwidthFilter {
220                max_bandwidth: [Bandwidth::zero(), Bandwidth::zero()],
221            },
222            cwnd_limited_before_aggregation_epoch: false,
223            cwnd_gain: params.startup_cwnd_gain,
224            pacing_gain: params.startup_pacing_gain,
225            full_bandwidth_reached: false,
226            bytes_lost_in_round: 0,
227            loss_events_in_round: 0,
228            max_bytes_delivered_in_round: 0,
229            bandwidth_latest: Bandwidth::zero(),
230            bandwidth_lo: None,
231            prior_bandwidth_lo: None,
232            inflight_latest: 0,
233            inflight_lo: usize::MAX,
234            inflight_hi: usize::MAX,
235
236            full_bandwidth_baseline: Bandwidth::zero(),
237            rounds_without_bandwidth_growth: 0,
238            rounds_with_queueing: 0,
239
240            ignore_app_limited_for_no_bandwidth_growth: params
241                .ignore_app_limited_for_no_bandwidth_growth,
242
243            latest_send_rate: None,
244            latest_ack_rate: None,
245
246            rtt_jump_detector: RttJumpDetector::new(params.rtt_jump_detector),
247        }
248    }
249
250    #[cfg(feature = "qlog")]
251    pub(super) fn send_rate(&self) -> Option<Bandwidth> {
252        self.latest_send_rate
253    }
254
255    #[cfg(feature = "qlog")]
256    pub(super) fn ack_rate(&self) -> Option<Bandwidth> {
257        self.latest_ack_rate
258    }
259
260    pub(super) fn max_ack_height(&self) -> usize {
261        self.bandwidth_sampler.max_ack_height().unwrap_or(0)
262    }
263
264    pub(super) fn bandwidth_estimate(&self) -> Bandwidth {
265        match (self.bandwidth_lo, self.max_bandwidth()) {
266            (None, b) => b,
267            (Some(a), b) => a.min(b),
268        }
269    }
270
271    pub(super) fn bdp(&self, bandwidth: Bandwidth, gain: f32) -> usize {
272        (bandwidth * gain).to_bytes_per_period(self.min_rtt()) as usize
273    }
274
275    pub(super) fn bdp1(&self, bandwidth: Bandwidth) -> usize {
276        self.bdp(bandwidth, 1.0)
277    }
278
279    pub(super) fn bdp0(&self) -> usize {
280        self.bdp1(self.max_bandwidth())
281    }
282
283    pub(super) fn min_rtt(&self) -> Duration {
284        self.min_rtt_filter.get()
285    }
286
287    pub(super) fn min_rtt_timestamp(&self) -> Instant {
288        self.min_rtt_filter.get_timestamps()
289    }
290
291    pub(super) fn max_bandwidth(&self) -> Bandwidth {
292        self.max_bandwidth_filter.get()
293    }
294
295    pub(super) fn on_packet_sent(
296        &mut self, sent_time: Instant, bytes_in_flight: usize,
297        packet_number: u64, bytes: usize, is_retransmissible: bool,
298    ) {
299        // Updating the min here ensures a more realistic (0) value when flows
300        // exit quiescence.
301        self.min_bytes_in_flight_in_round =
302            self.min_bytes_in_flight_in_round.min(bytes_in_flight);
303
304        if bytes_in_flight + bytes >= self.inflight_hi {
305            self.inflight_hi_limited_in_round = true;
306        }
307        self.round_trip_counter.on_packet_sent(packet_number);
308
309        self.bandwidth_sampler.on_packet_sent(
310            sent_time,
311            packet_number,
312            bytes,
313            bytes_in_flight,
314            is_retransmissible,
315        );
316    }
317
318    pub(super) fn on_congestion_event_start(
319        &mut self, acked_packets: &[Acked], lost_packets: &[Lost],
320        congestion_event: &mut BBRv2CongestionEvent, params: &Params,
321    ) {
322        let prior_bytes_acked = self.total_bytes_acked();
323        let prior_bytes_lost = self.total_bytes_lost();
324
325        let event_time = congestion_event.event_time;
326
327        congestion_event.end_of_round_trip =
328            if let Some(largest_acked) = acked_packets.last() {
329                self.round_trip_counter
330                    .on_packets_acked(largest_acked.pkt_num)
331            } else {
332                false
333            };
334
335        let sample = self.bandwidth_sampler.on_congestion_event(
336            event_time,
337            acked_packets,
338            lost_packets,
339            Some(self.max_bandwidth()),
340            self.bandwidth_lo.unwrap_or(Bandwidth::infinite()),
341            self.round_trip_count(),
342        );
343
344        if sample.extra_acked == 0 {
345            self.cwnd_limited_before_aggregation_epoch = congestion_event
346                .prior_bytes_in_flight >=
347                congestion_event.prior_cwnd;
348        }
349
350        if sample.last_packet_send_state.is_valid {
351            congestion_event.last_packet_send_state =
352                sample.last_packet_send_state;
353        }
354
355        // Do not update `max_bandwidth_filter` for a loss-only event or when no
356        // acknowledged packet produced a valid sample. In either case,
357        // `total_bytes_acked()` remains unchanged.
358        if let Some(sample_max) = sample.sample_max_bandwidth {
359            if prior_bytes_acked != self.total_bytes_acked() {
360                congestion_event.sample_max_bandwidth = Some(sample_max);
361                if !sample.sample_is_app_limited ||
362                    sample_max > self.max_bandwidth()
363                {
364                    self.max_bandwidth_filter.update(sample_max);
365                }
366            }
367        }
368
369        if let Some(rtt_sample) = sample.sample_rtt {
370            congestion_event.sample_min_rtt = Some(rtt_sample);
371
372            self.rtt_jump_detector.on_rtt_sample(
373                rtt_sample,
374                event_time,
375                self.full_bandwidth_reached,
376            );
377
378            self.min_rtt_filter.update(rtt_sample, event_time);
379        }
380
381        self.latest_send_rate = sample.sample_max_send_rate;
382        self.latest_ack_rate = sample.sample_max_ack_rate;
383
384        congestion_event.bytes_acked =
385            self.total_bytes_acked() - prior_bytes_acked;
386        congestion_event.bytes_lost = self.total_bytes_lost() - prior_bytes_lost;
387
388        congestion_event.bytes_in_flight = congestion_event
389            .prior_bytes_in_flight
390            .saturating_sub(congestion_event.bytes_acked)
391            .saturating_sub(congestion_event.bytes_lost);
392
393        if congestion_event.bytes_lost > 0 {
394            self.bytes_lost_in_round += congestion_event.bytes_lost;
395            self.loss_events_in_round += 1;
396        }
397
398        if congestion_event.bytes_acked > 0 &&
399            congestion_event.last_packet_send_state.is_valid &&
400            self.total_bytes_acked() >
401                congestion_event.last_packet_send_state.total_bytes_acked
402        {
403            let bytes_delivered = self.total_bytes_acked() -
404                congestion_event.last_packet_send_state.total_bytes_acked;
405            self.max_bytes_delivered_in_round =
406                self.max_bytes_delivered_in_round.max(bytes_delivered);
407        }
408
409        self.min_bytes_in_flight_in_round = self
410            .min_bytes_in_flight_in_round
411            .min(congestion_event.bytes_in_flight);
412
413        // `bandwidth_latest` and `inflight_latest` only increased within a
414        // round.
415        if sample.sample_max_bandwidth > Some(self.bandwidth_latest) {
416            self.bandwidth_latest = sample.sample_max_bandwidth.unwrap();
417        }
418
419        if sample.sample_max_inflight > self.inflight_latest {
420            self.inflight_latest = sample.sample_max_inflight;
421        }
422
423        // Adapt lower bounds(bandwidth_lo and inflight_lo).
424        self.adapt_lower_bounds(congestion_event, params);
425
426        if !congestion_event.end_of_round_trip {
427            return;
428        }
429
430        if let Some(bandwidth) = sample.sample_max_bandwidth {
431            self.bandwidth_latest = bandwidth;
432        }
433
434        if sample.sample_max_inflight > 0 {
435            self.inflight_latest = sample.sample_max_inflight;
436        }
437    }
438
439    pub(super) fn on_packet_neutered(&mut self, packet_number: u64) {
440        self.bandwidth_sampler.on_packet_neutered(packet_number)
441    }
442
443    fn adapt_lower_bounds(
444        &mut self, congestion_event: &BBRv2CongestionEvent, params: &Params,
445    ) {
446        if params.bw_lo_mode == BwLoMode::Default {
447            if !congestion_event.end_of_round_trip ||
448                congestion_event.is_probing_for_bandwidth
449            {
450                return;
451            }
452
453            if self.bytes_lost_in_round > 0 {
454                if self.bandwidth_lo.is_none() {
455                    self.bandwidth_lo = Some(self.max_bandwidth());
456                }
457
458                self.bandwidth_lo = Some(
459                    self.bandwidth_latest
460                        .max(self.bandwidth_lo.unwrap() * (1.0 - params.beta)),
461                );
462
463                if self.inflight_lo == usize::MAX {
464                    self.inflight_lo = congestion_event.prior_cwnd;
465                }
466
467                let inflight_lo_new =
468                    (self.inflight_lo as f32 * (1.0 - params.beta)) as usize;
469                self.inflight_lo = self.inflight_latest.max(inflight_lo_new);
470            }
471            return;
472        }
473
474        if congestion_event.bytes_lost == 0 {
475            return;
476        }
477
478        // Ignore losses from packets sent when probing for more bandwidth in
479        // STARTUP or PROBE_UP when they're lost in DRAIN or PROBE_DOWN.
480        if self.pacing_gain() < 1. {
481            return;
482        }
483
484        // Decrease bandwidth_lo whenever there is loss.
485        // Set `bandwidth_lo`if it is not yet set.
486        if self.bandwidth_lo.is_none() {
487            self.bandwidth_lo = Some(self.max_bandwidth());
488        }
489
490        // Save `bandwidth_lo` if it hasn't already been saved.
491        if self.prior_bandwidth_lo.is_none() {
492            self.prior_bandwidth_lo = self.bandwidth_lo;
493        }
494
495        match params.bw_lo_mode {
496            BwLoMode::Default => unreachable!("Handled above"),
497            BwLoMode::MinRttReduction => {
498                let reduction = Bandwidth::from_bytes_and_time_delta(
499                    congestion_event.bytes_lost,
500                    self.min_rtt(),
501                );
502
503                self.bandwidth_lo = self
504                    .bandwidth_lo
505                    .map(|b| (b - reduction).unwrap_or(Bandwidth::zero()));
506            },
507            BwLoMode::InflightReduction => {
508                // Use a max of BDP and inflight to avoid starving app-limited
509                // flows.
510                let effective_inflight =
511                    self.bdp0().max(congestion_event.prior_bytes_in_flight);
512                // This could use bytes_lost_in_round if the bandwidth_lo_ was
513                // saved when entering 'recovery', but this BBRv2
514                // implementation doesn't have recovery defined.
515                self.bandwidth_lo = self.bandwidth_lo.map(|b| {
516                    b * ((effective_inflight as f64 -
517                        congestion_event.bytes_lost as f64) /
518                        effective_inflight as f64)
519                });
520            },
521            BwLoMode::CwndReduction => {
522                self.bandwidth_lo = self.bandwidth_lo.map(|b| {
523                    b * ((congestion_event.prior_cwnd as f64 -
524                        congestion_event.bytes_lost as f64) /
525                        congestion_event.prior_cwnd as f64)
526                });
527            },
528        }
529
530        let mut last_bandwidth = self.bandwidth_latest;
531        // sample_max_bandwidth will be None if the loss is triggered by a timer
532        // expiring. Ideally we'd use the most recent bandwidth sample,
533        // but bandwidth_latest is safer than None.
534        if let Some(sample_max_bandwidth) = congestion_event.sample_max_bandwidth
535        {
536            // bandwidth_latest is the max bandwidth for the round, but to allow
537            // fast, conservation style response to loss, use the last sample.
538            last_bandwidth = sample_max_bandwidth;
539        }
540        if self.pacing_gain > params.full_bw_threshold {
541            // STARTUP applies `pacing_gain` to `bandwidth_lo`. Remove that
542            // factor so pacing can decrease without falling below
543            // the threshold.
544            self.bandwidth_lo = self.bandwidth_lo.max(Some(
545                last_bandwidth * (params.full_bw_threshold / self.pacing_gain),
546            ));
547        } else {
548            // Ensure bandwidth_lo isn't lower than last_bandwidth.
549            self.bandwidth_lo = self.bandwidth_lo.max(Some(last_bandwidth))
550        }
551        // At the end of a round, ensure `bandwidth_lo` does not decrease by
552        // more than `beta`.
553        if congestion_event.end_of_round_trip {
554            self.bandwidth_lo = self.bandwidth_lo.max(
555                self.prior_bandwidth_lo
556                    .take()
557                    .map(|b| b * (1.0 - params.beta)),
558            )
559        }
560        // These modes ignore inflight_lo as well.
561    }
562
563    pub(super) fn on_congestion_event_finish(
564        &mut self, least_unacked_packet: u64,
565        congestion_event: &BBRv2CongestionEvent,
566    ) {
567        if congestion_event.end_of_round_trip {
568            self.on_new_round();
569        }
570
571        self.bandwidth_sampler
572            .remove_obsolete_packets(least_unacked_packet);
573    }
574
575    pub(super) fn maybe_expire_min_rtt(
576        &mut self, congestion_event: &BBRv2CongestionEvent, params: &Params,
577    ) -> bool {
578        if congestion_event.sample_min_rtt.is_none() {
579            return false;
580        }
581
582        if congestion_event.event_time <
583            self.min_rtt_filter.min_rtt_timestamp + params.probe_rtt_period
584        {
585            return false;
586        }
587
588        self.min_rtt_filter.force_update(
589            congestion_event.sample_min_rtt.unwrap(),
590            congestion_event.event_time,
591        );
592
593        true
594    }
595
596    pub(super) fn is_inflight_too_high(
597        &self, congestion_event: &BBRv2CongestionEvent, max_loss_events: usize,
598        params: &Params,
599    ) -> bool {
600        let send_state = &congestion_event.last_packet_send_state;
601
602        if !send_state.is_valid {
603            // Not enough information.
604            return false;
605        }
606
607        if self.loss_events_in_round < max_loss_events {
608            return false;
609        }
610
611        // TODO(vlad): BytesInFlight(send_state);
612        let inflight_at_send = send_state.bytes_in_flight;
613
614        let bytes_lost_in_round = self.bytes_lost_in_round;
615
616        if inflight_at_send > 0 && bytes_lost_in_round > 0 {
617            let lost_in_round_threshold =
618                (inflight_at_send as f32 * params.loss_threshold) as usize;
619            if bytes_lost_in_round > lost_in_round_threshold {
620                return true;
621            }
622        }
623
624        false
625    }
626
627    pub(super) fn restart_round_early(&mut self) {
628        self.on_new_round();
629        self.round_trip_counter.restart_round();
630        self.rounds_with_queueing = 0;
631    }
632
633    fn on_new_round(&mut self) {
634        self.bytes_lost_in_round = 0;
635        self.loss_events_in_round = 0;
636        self.max_bytes_delivered_in_round = 0;
637        self.min_bytes_in_flight_in_round = usize::MAX;
638        self.inflight_hi_limited_in_round = false;
639    }
640
641    pub(super) fn has_bandwidth_growth(
642        &mut self, congestion_event: &BBRv2CongestionEvent, params: &Params,
643    ) -> bool {
644        let threshold = self.full_bandwidth_baseline * params.full_bw_threshold;
645
646        if self.max_bandwidth() >= threshold {
647            self.full_bandwidth_baseline = self.max_bandwidth();
648            self.rounds_without_bandwidth_growth = 0;
649            return true;
650        }
651
652        if !congestion_event.last_packet_send_state.is_valid {
653            // last_packet_send_state not available because the
654            // congestion event did not contain any non-ACK frames.
655            return false;
656        }
657
658        let ignore_round = self.ignore_app_limited_for_no_bandwidth_growth &&
659            congestion_event.last_packet_send_state.is_app_limited;
660
661        if !ignore_round {
662            self.rounds_without_bandwidth_growth += 1;
663        }
664
665        // full_bandwidth_reached is only set to true when not app-limited
666        if self.rounds_without_bandwidth_growth >= params.startup_full_bw_rounds &&
667            !congestion_event.last_packet_send_state.is_app_limited
668        {
669            self.full_bandwidth_reached = true;
670        }
671
672        false
673    }
674
675    pub(super) fn queueing_threshold_extra_bytes(&self) -> usize {
676        // TODO(vlad): 2 * mss
677        2 * DEFAULT_MSS
678    }
679
680    pub(super) fn check_persistent_queue(
681        &mut self, target_gain: f32, params: &Params,
682    ) {
683        let target = self
684            .bdp(self.max_bandwidth(), target_gain)
685            .max(self.bdp0() + self.queueing_threshold_extra_bytes());
686
687        if self.min_bytes_in_flight_in_round < target {
688            self.rounds_with_queueing = 0;
689            return;
690        }
691
692        self.rounds_with_queueing += 1;
693        #[allow(clippy::absurd_extreme_comparisons)]
694        if self.rounds_with_queueing >= params.max_startup_queue_rounds {
695            self.full_bandwidth_reached = true;
696        }
697    }
698
699    pub(super) fn max_bytes_delivered_in_round(&self) -> usize {
700        self.max_bytes_delivered_in_round
701    }
702
703    pub(super) fn total_bytes_acked(&self) -> usize {
704        self.bandwidth_sampler.total_bytes_acked()
705    }
706
707    pub(super) fn total_bytes_lost(&self) -> usize {
708        self.bandwidth_sampler.total_bytes_lost()
709    }
710
711    /// Total number of confirmed persistent RTT jump episodes over the
712    /// lifetime of the connection.
713    pub(super) fn rtt_persistent_jump_count(&self) -> u64 {
714        self.rtt_jump_detector.rtt_persistent_jump_count()
715    }
716
717    /// The start time of the most recently confirmed persistent RTT jump
718    /// episode, if any.
719    #[cfg(test)]
720    pub(super) fn last_persistent_jump_time(&self) -> Option<Instant> {
721        self.rtt_jump_detector.last_persistent_jump_time()
722    }
723
724    /// Whether an RTT jump episode is currently active (elevated but not yet
725    /// resolved), regardless of whether it has been confirmed persistent.
726    #[cfg(test)]
727    pub(super) fn is_rtt_jump_active(&self) -> bool {
728        self.rtt_jump_detector.is_rtt_jump_active()
729    }
730
731    /// Whether the current RTT jump episode has been confirmed as a persistent
732    /// network condition.
733    #[cfg(test)]
734    pub(super) fn is_rtt_jump_persistent(&self) -> bool {
735        self.rtt_jump_detector.is_rtt_jump_persistent()
736    }
737
738    fn round_trip_count(&self) -> usize {
739        self.round_trip_counter.round_trip_count
740    }
741
742    pub(super) fn full_bandwidth_reached(&self) -> bool {
743        self.full_bandwidth_reached
744    }
745
746    pub(super) fn set_full_bandwidth_reached(&mut self) {
747        self.full_bandwidth_reached = true
748    }
749
750    pub(super) fn pacing_gain(&self) -> f32 {
751        self.pacing_gain
752    }
753
754    pub(super) fn set_pacing_gain(&mut self, pacing_gain: f32) {
755        self.pacing_gain = pacing_gain
756    }
757
758    pub(super) fn cwnd_gain(&self) -> f32 {
759        self.cwnd_gain
760    }
761
762    pub(super) fn set_cwnd_gain(&mut self, cwnd_gain: f32) {
763        self.cwnd_gain = cwnd_gain
764    }
765
766    pub(super) fn inflight_hi(&self) -> usize {
767        self.inflight_hi
768    }
769
770    pub(super) fn inflight_hi_with_headroom(&self, params: &Params) -> usize {
771        let headroom =
772            (self.inflight_hi as f32 * params.inflight_hi_headroom) as usize;
773        self.inflight_hi.saturating_sub(headroom)
774    }
775
776    pub(super) fn set_inflight_hi(&mut self, new_inflight_hi: usize) {
777        self.inflight_hi = new_inflight_hi
778    }
779
780    pub(super) fn inflight_hi_default(&self) -> usize {
781        usize::MAX
782    }
783
784    pub(super) fn inflight_lo(&self) -> usize {
785        self.inflight_lo
786    }
787
788    pub(super) fn clear_inflight_lo(&mut self) {
789        self.inflight_lo = usize::MAX
790    }
791
792    pub(super) fn cap_inflight_lo(&mut self, cap: usize) {
793        if self.inflight_lo != usize::MAX {
794            self.inflight_lo = cap.min(self.inflight_lo)
795        }
796    }
797
798    pub(super) fn clear_bandwidth_lo(&mut self) {
799        self.bandwidth_lo = None
800    }
801
802    pub(super) fn advance_max_bandwidth_filter(&mut self) {
803        self.max_bandwidth_filter.advance()
804    }
805
806    pub(super) fn postpone_min_rtt_timestamp(&mut self, duration: Duration) {
807        self.min_rtt_filter
808            .force_update(self.min_rtt(), self.min_rtt_timestamp().add(duration));
809    }
810
811    pub(super) fn on_app_limited(&mut self) {
812        self.bandwidth_sampler.on_app_limited()
813    }
814
815    pub(super) fn loss_events_in_round(&self) -> usize {
816        self.loss_events_in_round
817    }
818
819    pub(super) fn rounds_with_queueing(&self) -> usize {
820        self.rounds_with_queueing
821    }
822}
823
824#[cfg(test)]
825mod tests {
826    use super::*;
827    use crate::recovery::gcongestion::bbr2::DEFAULT_PARAMS;
828    use crate::recovery::gcongestion::BbrRttJumpDetector;
829
830    fn ms(millis: u64) -> Duration {
831        Duration::from_millis(millis)
832    }
833
834    const RTT: Duration = Duration::from_millis(50);
835    const RTT_3X: Duration = Duration::from_millis(150);
836    const RTT_JUMP: Duration = Duration::from_millis(151);
837
838    /// Ack a packet with the given RTT through the real congestion-event path.
839    fn ack_with_rtt(
840        model: &mut BBRv2NetworkModel, params: &Params, pkt_num: u64,
841        base: Instant, sent_offset: Duration, rtt: Duration,
842    ) -> BBRv2CongestionEvent {
843        ack_with_rtt_util(
844            model,
845            params,
846            pkt_num,
847            base,
848            sent_offset,
849            rtt,
850            100_000,
851            1200,
852        )
853    }
854
855    /// As `ack_with_rtt`, but with explicit cwnd and inflight inputs.
856    // Test harness: the extra cwnd/inflight knobs push this one over the
857    // argument-count lint, which is not worth a builder struct in tests.
858    #[allow(clippy::too_many_arguments)]
859    fn ack_with_rtt_util(
860        model: &mut BBRv2NetworkModel, params: &Params, pkt_num: u64,
861        base: Instant, sent_offset: Duration, rtt: Duration, prior_cwnd: usize,
862        prior_in_flight: usize,
863    ) -> BBRv2CongestionEvent {
864        let bytes = 1200;
865        let sent_time = base + sent_offset;
866        model.on_packet_sent(sent_time, 0, pkt_num, bytes, true);
867
868        let ack_time = sent_time + rtt;
869        let acked = [Acked {
870            pkt_num,
871            time_sent: sent_time,
872        }];
873        let mut event = BBRv2CongestionEvent::new(
874            ack_time,
875            prior_cwnd,
876            prior_in_flight,
877            false,
878        );
879        model.on_congestion_event_start(&acked, &[], &mut event, params);
880        event
881    }
882
883    fn rtt_jump_params(detector: BbrRttJumpDetector) -> Params {
884        Params {
885            rtt_jump_detector: detector,
886            ..DEFAULT_PARAMS
887        }
888    }
889
890    #[test]
891    fn rtt_jump_detector_is_disabled_by_default() {
892        let params = &DEFAULT_PARAMS;
893        let mut model = BBRv2NetworkModel::new(params, RTT);
894        let base = Instant::now();
895
896        for pkt in 1..5 {
897            ack_with_rtt(&mut model, params, pkt, base, ms(pkt * 10), RTT);
898        }
899
900        let mut offset = 100;
901        for pkt in 5..20 {
902            ack_with_rtt(&mut model, params, pkt, base, ms(offset), RTT_3X);
903            offset += 100;
904        }
905
906        assert_eq!(model.rtt_persistent_jump_count(), 0);
907        assert!(!model.is_rtt_jump_active());
908        assert_eq!(model.last_persistent_jump_time(), None);
909    }
910
911    #[test]
912    fn global_min_detector_can_be_enabled() {
913        let params = &rtt_jump_params(BbrRttJumpDetector::GlobalMin);
914        let mut model = BBRv2NetworkModel::new(params, RTT);
915        let base = Instant::now();
916
917        for pkt in 1..5 {
918            ack_with_rtt(&mut model, params, pkt, base, ms(pkt * 10), RTT);
919        }
920        model.set_full_bandwidth_reached();
921
922        ack_with_rtt(&mut model, params, 5, base, ms(100), RTT_JUMP);
923        ack_with_rtt(&mut model, params, 6, base, ms(110), RTT_JUMP);
924        ack_with_rtt(&mut model, params, 7, base, ms(300), RTT_JUMP);
925
926        assert_eq!(model.rtt_persistent_jump_count(), 1);
927        assert!(model.is_rtt_jump_persistent());
928        assert!(model.last_persistent_jump_time().is_some());
929    }
930
931    #[test]
932    fn hmm_detector_can_be_enabled() {
933        let params = &rtt_jump_params(BbrRttJumpDetector::Hmm);
934        let mut model = BBRv2NetworkModel::new(params, RTT);
935        let base = Instant::now();
936
937        for pkt in 1..9u64 {
938            ack_with_rtt(&mut model, params, pkt, base, ms(pkt * 10), RTT);
939        }
940        model.set_full_bandwidth_reached();
941
942        let mut offset = 200u64;
943        for pkt in 9u64..40 {
944            ack_with_rtt(&mut model, params, pkt, base, ms(offset), RTT_3X);
945            offset += 100;
946        }
947
948        assert_eq!(model.rtt_persistent_jump_count(), 1);
949        assert!(model.is_rtt_jump_persistent());
950        assert!(model.last_persistent_jump_time().is_some());
951    }
952
953    #[test]
954    fn global_min_detector_sustained_step_becomes_persistent() {
955        let params = &rtt_jump_params(BbrRttJumpDetector::GlobalMin);
956        let mut model = BBRv2NetworkModel::new(params, RTT);
957        let base = Instant::now();
958
959        for pkt in 1..5 {
960            ack_with_rtt(&mut model, params, pkt, base, ms(pkt * 10), RTT);
961        }
962        model.set_full_bandwidth_reached();
963
964        ack_with_rtt(&mut model, params, 5, base, ms(100), RTT_JUMP);
965        assert!(model.is_rtt_jump_active());
966        assert!(!model.is_rtt_jump_persistent());
967
968        ack_with_rtt(&mut model, params, 6, base, ms(110), RTT_JUMP);
969        assert!(!model.is_rtt_jump_persistent());
970
971        ack_with_rtt(&mut model, params, 7, base, ms(300), RTT_JUMP);
972        assert!(model.is_rtt_jump_persistent());
973        assert_eq!(model.rtt_persistent_jump_count(), 1);
974        assert!(model.last_persistent_jump_time().is_some());
975    }
976
977    #[test]
978    fn global_min_detector_uses_strict_3x_threshold() {
979        let params = &rtt_jump_params(BbrRttJumpDetector::GlobalMin);
980        let base = Instant::now();
981
982        let mut at_edge = BBRv2NetworkModel::new(params, RTT);
983        ack_with_rtt(&mut at_edge, params, 1, base, ms(10), RTT);
984        at_edge.set_full_bandwidth_reached();
985        ack_with_rtt(&mut at_edge, params, 2, base, ms(20), RTT_3X);
986        assert!(!at_edge.is_rtt_jump_active());
987
988        let mut just_above = BBRv2NetworkModel::new(params, RTT);
989        ack_with_rtt(&mut just_above, params, 1, base, ms(10), RTT);
990        just_above.set_full_bandwidth_reached();
991        ack_with_rtt(&mut just_above, params, 2, base, ms(20), RTT_JUMP);
992        assert!(just_above.is_rtt_jump_active());
993    }
994
995    #[test]
996    fn global_min_detector_tracks_downward_baseline() {
997        let params = &rtt_jump_params(BbrRttJumpDetector::GlobalMin);
998        let mut model = BBRv2NetworkModel::new(params, RTT);
999        let base = Instant::now();
1000
1001        ack_with_rtt(&mut model, params, 1, base, ms(100), ms(100));
1002        ack_with_rtt(&mut model, params, 2, base, ms(200), ms(299));
1003        assert!(!model.is_rtt_jump_active());
1004
1005        ack_with_rtt(&mut model, params, 3, base, ms(300), RTT);
1006        assert!(!model.is_rtt_jump_active());
1007
1008        model.set_full_bandwidth_reached();
1009        ack_with_rtt(&mut model, params, 4, base, ms(400), RTT_JUMP);
1010        assert!(model.is_rtt_jump_active());
1011    }
1012
1013    #[test]
1014    fn global_min_detector_ignores_startup_rtt_jump_until_full_bandwidth() {
1015        let params = &rtt_jump_params(BbrRttJumpDetector::GlobalMin);
1016        let mut model = BBRv2NetworkModel::new(params, RTT);
1017        let base = Instant::now();
1018
1019        ack_with_rtt(&mut model, params, 1, base, ms(10), RTT);
1020        ack_with_rtt(&mut model, params, 2, base, ms(20), RTT_JUMP);
1021        ack_with_rtt(&mut model, params, 3, base, ms(30), RTT_JUMP);
1022        ack_with_rtt(&mut model, params, 4, base, ms(200), RTT_JUMP);
1023
1024        assert_eq!(model.rtt_persistent_jump_count(), 0);
1025        assert!(!model.is_rtt_jump_active());
1026
1027        model.set_full_bandwidth_reached();
1028        ack_with_rtt(&mut model, params, 5, base, ms(210), RTT_JUMP);
1029        ack_with_rtt(&mut model, params, 6, base, ms(220), RTT_JUMP);
1030        ack_with_rtt(&mut model, params, 7, base, ms(360), RTT_JUMP);
1031
1032        assert_eq!(model.rtt_persistent_jump_count(), 1);
1033        assert!(model.is_rtt_jump_persistent());
1034    }
1035
1036    #[test]
1037    fn hmm_detector_ignores_startup_rtt_jump_until_full_bandwidth() {
1038        let params = &rtt_jump_params(BbrRttJumpDetector::Hmm);
1039        let mut model = BBRv2NetworkModel::new(params, RTT);
1040        let base = Instant::now();
1041
1042        for pkt in 1..9u64 {
1043            ack_with_rtt(&mut model, params, pkt, base, ms(pkt * 10), RTT);
1044        }
1045
1046        for pkt in 9u64..20 {
1047            ack_with_rtt(&mut model, params, pkt, base, ms(pkt * 100), RTT_3X);
1048        }
1049
1050        assert_eq!(model.rtt_persistent_jump_count(), 0);
1051        assert!(!model.is_rtt_jump_active());
1052
1053        model.set_full_bandwidth_reached();
1054        ack_with_rtt(&mut model, params, 20, base, ms(2100), RTT_3X);
1055
1056        assert_eq!(model.rtt_persistent_jump_count(), 0);
1057        assert!(!model.is_rtt_jump_persistent());
1058    }
1059}