quiche/recovery/congestion/
delivery_rate.rs1use 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 end_of_app_limited: u64,
50
51 last_sent_packet: u64,
53
54 largest_acked: u64,
56
57 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 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 pub fn update_rate_sample(&mut self, pkt: &Acked, now: Instant) {
105 self.delivered += pkt.size;
106 self.delivered_time = now;
107
108 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 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 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 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 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 #[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 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 let rtt = Duration::from_secs(2);
275 helper_send_and_ack_packets(&mut r, 0..4, now, rtt, mss);
276
277 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 assert!(r.congestion.delivery_rate.app_limited());
284 assert!(!r.congestion.delivery_rate.sample_is_app_limited());
285
286 now += rtt;
288 helper_send_and_ack_packets(&mut r, 4..8, now, rtt, mss);
289
290 assert!(!r.congestion.delivery_rate.app_limited());
293 assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 0);
294 assert!(r.congestion.delivery_rate.sample_is_app_limited());
296 }
297
298 #[test]
301 fn app_limited_delivery_rate() {
302 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 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 let mut rtt = Duration::from_secs(2);
315 helper_send_and_ack_packets(&mut r, 0..2, now, rtt, mss);
316
317 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 now += rtt;
331 rtt = Duration::from_secs(4);
332 helper_send_and_ack_packets(&mut r, 2..4, now, rtt, mss);
333
334 assert!(!r.congestion.delivery_rate.app_limited());
337 assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 0);
338 assert!(r.congestion.delivery_rate.sample_is_app_limited());
340
341 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 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 assert!(r.congestion.delivery_rate.sample_is_app_limited());
353
354 now += rtt;
357 rtt = Duration::from_secs(1);
358 helper_send_and_ack_packets(&mut r, 4..6, now, rtt, mss);
359
360 assert!(!r.congestion.delivery_rate.app_limited());
363 assert_eq!(r.congestion.delivery_rate.end_of_app_limited, 0);
364 assert!(r.congestion.delivery_rate.sample_is_app_limited());
366
367 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 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 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 r.congestion.delivery_rate.generate_rate_sample(rtt);
416
417 assert_eq!(r.congestion.delivery_rate.delivered(), 2400);
419
420 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 assert!(!r.app_limited());
434 assert!(!r.congestion.delivery_rate.sample_is_app_limited());
435
436 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 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 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 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}