Skip to main content

quiche/recovery/congestion/
delivery_rate.rs

1// Copyright (C) 2020-2022, Cloudflare, Inc.
2// All rights reserved.
3//
4// Redistribution and use in source and binary forms, with or without
5// modification, are permitted provided that the following conditions are
6// met:
7//
8//     * Redistributions of source code must retain the above copyright notice,
9//       this list of conditions and the following disclaimer.
10//
11//     * Redistributions in binary form must reproduce the above copyright
12//       notice, this list of conditions and the following disclaimer in the
13//       documentation and/or other materials provided with the distribution.
14//
15// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS
16// IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO,
17// THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
18// PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
19// CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
20// EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
21// PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
22// PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
23// LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
24// NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
25// SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
26
27//! Delivery rate estimation.
28//!
29//! This implements the algorithm for estimating delivery rate as described in
30//! <https://tools.ietf.org/html/draft-cheng-iccrg-delivery-rate-estimation-01>
31
32use std::time::Duration;
33use std::time::Instant;
34
35use crate::recovery::bandwidth::Bandwidth;
36
37use super::Acked;
38use super::Sent;
39
40#[derive(Debug)]
41pub struct Rate {
42    delivered: usize,
43
44    delivered_time: Instant,
45
46    first_sent_time: Instant,
47
48    // Packet number of the last sent packet with app limited.
49    end_of_app_limited: u64,
50
51    // Packet number of the last sent packet.
52    last_sent_packet: u64,
53
54    // Packet number of the largest acked packet.
55    largest_acked: u64,
56
57    // Sample of rate estimation.
58    rate_sample: RateSample,
59}
60
61impl Default for Rate {
62    fn default() -> Self {
63        let now = Instant::now();
64
65        Rate {
66            delivered: 0,
67
68            delivered_time: now,
69
70            first_sent_time: now,
71
72            end_of_app_limited: 0,
73
74            last_sent_packet: 0,
75
76            largest_acked: 0,
77
78            rate_sample: RateSample::new(),
79        }
80    }
81}
82
83impl Rate {
84    pub fn on_packet_sent(
85        &mut self, pkt: &mut Sent, bytes_in_flight: usize, bytes_lost: u64,
86    ) {
87        // No packets in flight.
88        if bytes_in_flight == 0 {
89            self.first_sent_time = pkt.time_sent;
90            self.delivered_time = pkt.time_sent;
91        }
92
93        pkt.first_sent_time = self.first_sent_time;
94        pkt.delivered_time = self.delivered_time;
95        pkt.delivered = self.delivered;
96        pkt.is_app_limited = self.app_limited();
97        pkt.tx_in_flight = bytes_in_flight;
98        pkt.lost = bytes_lost;
99
100        self.last_sent_packet = pkt.pkt_num;
101    }
102
103    // Update the delivery rate sample when a packet is acked.
104    pub fn update_rate_sample(&mut self, pkt: &Acked, now: Instant) {
105        self.delivered += pkt.size;
106        self.delivered_time = now;
107
108        // Update info using the newest packet. If rate_sample is not yet
109        // initialized, initialize with the first packet.
110        if self.rate_sample.prior_time.is_none() ||
111            pkt.delivered >= self.rate_sample.prior_delivered
112        {
113            self.rate_sample.prior_delivered = pkt.delivered;
114            self.rate_sample.prior_time = Some(pkt.delivered_time);
115            self.rate_sample.is_app_limited = pkt.is_app_limited;
116            self.rate_sample.send_elapsed =
117                pkt.time_sent.saturating_duration_since(pkt.first_sent_time);
118            self.rate_sample.rtt = pkt.rtt;
119            self.rate_sample.ack_elapsed = self
120                .delivered_time
121                .saturating_duration_since(pkt.delivered_time);
122
123            self.first_sent_time = pkt.time_sent;
124        }
125
126        self.largest_acked = self.largest_acked.max(pkt.pkt_num);
127    }
128
129    pub fn generate_rate_sample(&mut self, min_rtt: Duration) {
130        // End app-limited phase if bubble is ACKed and gone.
131        if self.app_limited() && self.largest_acked > self.end_of_app_limited {
132            self.update_app_limited(false);
133        }
134
135        if self.rate_sample.prior_time.is_some() {
136            let interval = self
137                .rate_sample
138                .send_elapsed
139                .max(self.rate_sample.ack_elapsed);
140
141            self.rate_sample.delivered =
142                self.delivered - self.rate_sample.prior_delivered;
143            self.rate_sample.interval = interval;
144
145            if interval < min_rtt {
146                self.rate_sample.interval = Duration::ZERO;
147
148                // No reliable sample.
149                return;
150            }
151
152            if !interval.is_zero() {
153                let rate_sample_bandwidth = {
154                    let rate_sample_bytes_per_second = (self.rate_sample.delivered
155                        as f64 /
156                        interval.as_secs_f64())
157                        as u64;
158
159                    Bandwidth::from_bytes_per_second(rate_sample_bytes_per_second)
160                };
161
162                // Match the [linux] implementation and only generate a new
163                // sample delivery rate if either:
164                // - the sample was not app_limited
165                // - the new rate is higher than the previous value
166                //
167                // [linux] https://github.com/torvalds/linux/commit/eb8329e0a04db0061f714f033b4454326ba147f4
168                if !self.rate_sample.is_app_limited ||
169                    rate_sample_bandwidth > self.rate_sample.bandwidth
170                {
171                    self.update_delivery_rate(rate_sample_bandwidth);
172                }
173            }
174        }
175    }
176
177    fn update_delivery_rate(&mut self, bandwidth: Bandwidth) {
178        self.rate_sample.bandwidth = bandwidth;
179    }
180
181    pub fn update_app_limited(&mut self, v: bool) {
182        self.end_of_app_limited =
183            if v { self.last_sent_packet.max(1) } else { 0 };
184    }
185
186    pub fn app_limited(&mut self) -> bool {
187        self.end_of_app_limited != 0
188    }
189
190    #[cfg(test)]
191    pub fn delivered(&self) -> usize {
192        self.delivered
193    }
194
195    pub fn sample_delivery_rate(&self) -> Bandwidth {
196        self.rate_sample.bandwidth
197    }
198
199    #[cfg(test)]
200    pub fn sample_is_app_limited(&self) -> bool {
201        self.rate_sample.is_app_limited
202    }
203}
204
205#[derive(Debug)]
206struct RateSample {
207    // The sample delivery_rate in bytes/sec
208    bandwidth: Bandwidth,
209
210    is_app_limited: bool,
211
212    interval: Duration,
213
214    delivered: usize,
215
216    prior_delivered: usize,
217
218    prior_time: Option<Instant>,
219
220    send_elapsed: Duration,
221
222    ack_elapsed: Duration,
223
224    rtt: Duration,
225}
226
227impl RateSample {
228    const fn new() -> Self {
229        RateSample {
230            bandwidth: Bandwidth::zero(),
231            is_app_limited: false,
232            interval: Duration::ZERO,
233            delivered: 0,
234            prior_delivered: 0,
235            prior_time: None,
236            send_elapsed: Duration::ZERO,
237            ack_elapsed: Duration::ZERO,
238            rtt: Duration::ZERO,
239        }
240    }
241}
242
243#[cfg(test)]
244mod tests {
245    use super::*;
246
247    use crate::packet;
248    use crate::ranges;
249    use crate::recovery::congestion::recovery::LegacyRecovery;
250    use crate::recovery::HandshakeStatus;
251    use crate::recovery::RecoveryOps;
252    use crate::test_utils;
253    use crate::Config;
254    use crate::OnAckReceivedOutcome;
255    use std::ops::Range;
256
257    // A `RateSample` generated while `Rate` is app-limited inherits that state.
258    //
259    // This test generates samples before and after `Rate` becomes app-limited
260    // and checks each sample's state.
261    #[test]
262    fn sample_is_app_limited() {
263        let config = Config::new(0xbabababa).unwrap();
264        let mut r = LegacyRecovery::new(&config);
265        let mut now = Instant::now();
266        let mss = r.max_datagram_size();
267
268        // `Rate` is not app-limited before any activity.
269        assert!(!r.congestion.delivery_rate.app_limited());
270        assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 0);
271        assert!(!r.congestion.delivery_rate.sample_is_app_limited());
272
273        // Generate a delivery-rate sample from the first batch.
274        let rtt = Duration::from_secs(2);
275        helper_send_and_ack_packets(&mut r, 0..4, now, rtt, mss);
276
277        // Mark `Rate` as app-limited.
278        r.delivery_rate_update_app_limited(true);
279        assert!(r.congestion.delivery_rate.app_limited());
280        assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 3);
281
282        // `Rate` is app-limited.
283        assert!(r.congestion.delivery_rate.app_limited());
284        assert!(!r.congestion.delivery_rate.sample_is_app_limited());
285
286        // Send and acknowledge the second batch to generate another sample.
287        now += rtt;
288        helper_send_and_ack_packets(&mut r, 4..8, now, rtt, mss);
289
290        // `Rate` is no longer app-limited after sending a packet beyond
291        // `end_of_app_limited`.
292        assert!(!r.congestion.delivery_rate.app_limited());
293        assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 0);
294        // The resulting `RateSample` remains app-limited.
295        assert!(r.congestion.delivery_rate.sample_is_app_limited());
296    }
297
298    // A `RateSample` updates the delivery rate only when it is not app-limited
299    // or its rate exceeds the previous value.
300    #[test]
301    fn app_limited_delivery_rate() {
302        // Confirm that a rate sample is not generated when app-limited.
303        let config = Config::new(0xbabababa).unwrap();
304        let mut r = LegacyRecovery::new(&config);
305        let mut now = Instant::now();
306        let mss = r.max_datagram_size();
307
308        // `Rate` is not app-limited before any activity.
309        assert!(!r.congestion.delivery_rate.app_limited());
310        assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 0);
311        assert!(!r.congestion.delivery_rate.sample_is_app_limited());
312
313        // Generate a delivery-rate sample from the first batch.
314        let mut rtt = Duration::from_secs(2);
315        helper_send_and_ack_packets(&mut r, 0..2, now, rtt, mss);
316
317        // Mark `Rate` as app-limited.
318        r.delivery_rate_update_app_limited(true);
319        assert!(r.congestion.delivery_rate.app_limited());
320        assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 1);
321        assert!(!r.congestion.delivery_rate.sample_is_app_limited());
322
323        let first_delivery_rate = r.delivery_rate().to_bytes_per_second();
324        let expected_delivery_rate = (mss * 2) as u64 / rtt.as_secs();
325        assert_eq!(expected_delivery_rate, 1200);
326        assert_eq!(first_delivery_rate, expected_delivery_rate);
327
328        // A larger RTT produces a lower delivery rate, which does not replace
329        // the app-limited sample.
330        now += rtt;
331        rtt = Duration::from_secs(4);
332        helper_send_and_ack_packets(&mut r, 2..4, now, rtt, mss);
333
334        // `Rate` is no longer app-limited after sending a packet beyond
335        // `end_of_app_limited`.
336        assert!(!r.congestion.delivery_rate.app_limited());
337        assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 0);
338        // The resulting `RateSample` remains app-limited.
339        assert!(r.congestion.delivery_rate.sample_is_app_limited());
340
341        // The lower delivery rate does not replace the previous value.
342        let expected_delivery_rate = (mss * 2) as u64 / rtt.as_secs();
343        assert_eq!(expected_delivery_rate, 600);
344        let app_limited_delivery_rate = r.delivery_rate().to_bytes_per_second();
345        assert_eq!(app_limited_delivery_rate, first_delivery_rate);
346
347        // Mark `Rate` as app-limited.
348        r.delivery_rate_update_app_limited(true);
349        assert!(r.congestion.delivery_rate.app_limited());
350        assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 3);
351        // The resulting `RateSample` remains app-limited.
352        assert!(r.congestion.delivery_rate.sample_is_app_limited());
353
354        // A smaller RTT produces a higher delivery rate, which replaces the
355        // previous value even while app-limited.
356        now += rtt;
357        rtt = Duration::from_secs(1);
358        helper_send_and_ack_packets(&mut r, 4..6, now, rtt, mss);
359
360        // `Rate` is no longer app-limited after sending a packet beyond
361        // `end_of_app_limited`.
362        assert!(!r.congestion.delivery_rate.app_limited());
363        assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 0);
364        // The resulting `RateSample` remains app-limited.
365        assert!(r.congestion.delivery_rate.sample_is_app_limited());
366
367        // The higher delivery rate replaces the previous value.
368        let expected_delivery_rate = (mss * 2) as u64 / rtt.as_secs();
369        assert_eq!(expected_delivery_rate, 2400);
370        let app_limited_delivery_rate = r.delivery_rate().to_bytes_per_second();
371        assert_eq!(app_limited_delivery_rate, expected_delivery_rate);
372    }
373
374    #[test]
375    fn rate_check() {
376        let config = Config::new(0xbabababa).unwrap();
377        let mut r = LegacyRecovery::new(&config);
378
379        let now = Instant::now();
380        let mss = r.max_datagram_size();
381
382        // Send 2 packets.
383        for pn in 0..2 {
384            let pkt = test_utils::helper_packet_sent(pn, now, mss);
385
386            r.on_packet_sent(
387                pkt,
388                packet::Epoch::Application,
389                HandshakeStatus::default(),
390                now,
391                "",
392            );
393        }
394
395        let rtt = Duration::from_millis(50);
396        let now = now + rtt;
397
398        // Ack 2 packets.
399        for pn in 0..2 {
400            let acked = Acked {
401                pkt_num: pn,
402                time_sent: now,
403                size: mss,
404                rtt,
405                delivered: 0,
406                delivered_time: now,
407                first_sent_time: now.checked_sub(rtt).unwrap(),
408                is_app_limited: false,
409            };
410
411            r.congestion.delivery_rate.update_rate_sample(&acked, now);
412        }
413
414        // Update rate sample after 1 rtt.
415        r.congestion.delivery_rate.generate_rate_sample(rtt);
416
417        // Bytes acked so far.
418        assert_eq!(r.congestion.delivery_rate.delivered(), 2400);
419
420        // Estimated delivery rate = (1200 x 2) / 0.05s = 48000.
421        assert_eq!(r.delivery_rate().to_bytes_per_second(), 48000);
422    }
423
424    #[test]
425    fn app_limited_cwnd_full() {
426        let config = Config::new(0xbabababa).unwrap();
427        let mut r = LegacyRecovery::new(&config);
428
429        let now = Instant::now();
430        let mss = r.max_datagram_size();
431
432        // Not App Limited prior to any activity
433        assert!(!r.app_limited());
434        assert!(!r.congestion.delivery_rate.sample_is_app_limited());
435
436        // Send 10 packets to fill cwnd.
437        for pn in 0..5 {
438            let pkt = test_utils::helper_packet_sent(pn, now, mss);
439            r.on_packet_sent(
440                pkt,
441                packet::Epoch::Application,
442                HandshakeStatus::default(),
443                now,
444                "",
445            );
446        }
447
448        // App Limited after sending partial cwnd worth of data
449        assert!(r.app_limited());
450        assert!(!r.congestion.delivery_rate.sample_is_app_limited());
451
452        for pn in 5..10 {
453            let pkt = test_utils::helper_packet_sent(pn, now, mss);
454            r.on_packet_sent(
455                pkt,
456                packet::Epoch::Application,
457                HandshakeStatus::default(),
458                now,
459                "",
460            );
461        }
462
463        // Not App Limited after sending full cwnd worth of data
464        assert!(!r.app_limited());
465        assert!(!r.congestion.delivery_rate.sample_is_app_limited());
466    }
467
468    fn helper_send_and_ack_packets(
469        recovery: &mut LegacyRecovery, range: Range<u64>, now: Instant,
470        rtt: Duration, mss: usize,
471    ) {
472        for pn in range.clone() {
473            let pkt = test_utils::helper_packet_sent(pn, now, mss);
474            recovery.on_packet_sent(
475                pkt,
476                packet::Epoch::Application,
477                HandshakeStatus::default(),
478                now,
479                "",
480            );
481        }
482
483        let packet_count = range.clone().count();
484
485        // Ack packets, which generates a new delivery_rate
486        let mut acked = ranges::RangeSet::default();
487        acked.insert(range);
488
489        let ack_outcome = recovery
490            .on_ack_received(
491                &acked,
492                25,
493                packet::Epoch::Application,
494                HandshakeStatus::default(),
495                now + rtt,
496                None,
497                "",
498            )
499            .unwrap();
500
501        assert_eq!(ack_outcome, OnAckReceivedOutcome {
502            lost_packets: 0,
503            lost_bytes: 0,
504            acked_bytes: mss * packet_count,
505            spurious_losses: 0,
506        });
507    }
508}