quiche/recovery/gcongestion/
pacer.rs1use std::time::Instant;
32
33use crate::recovery::gcongestion::bbr2::BBRv2;
34use crate::recovery::gcongestion::Bandwidth;
35use crate::recovery::gcongestion::CongestionControl;
36use crate::recovery::rtt::RttStats;
37use crate::recovery::RecoveryStats;
38use crate::recovery::ReleaseDecision;
39use crate::recovery::ReleaseTime;
40
41use super::Acked;
42use super::Lost;
43
44const LUMPY_PACING_CWND_FRACTION: f64 = 0.25;
47
48const LUMPY_PACING_SIZE: usize = 2;
51
52const LUMPY_PACING_MIN_BANDWIDTH_KBPS: Bandwidth =
55 Bandwidth::from_kbits_per_second(1_200);
56
57const INITIAL_UNPACED_BURST: usize = 10;
60
61#[derive(Debug)]
62pub struct Pacer {
63 enabled: bool,
65 sender: BBRv2,
67 max_pacing_rate: Option<Bandwidth>,
69 burst_tokens: usize,
71 ideal_next_packet_send_time: ReleaseTime,
73 initial_burst_size: usize,
74 lumpy_tokens: usize,
77 pacing_limited: bool,
80}
81
82impl Pacer {
83 pub(crate) fn new(
87 enabled: bool, congestion: BBRv2, max_pacing_rate: Option<Bandwidth>,
88 ) -> Self {
89 Pacer {
90 enabled,
91 sender: congestion,
92 max_pacing_rate,
93 burst_tokens: INITIAL_UNPACED_BURST,
94 ideal_next_packet_send_time: ReleaseTime::Immediate,
95 initial_burst_size: INITIAL_UNPACED_BURST,
96 lumpy_tokens: 0,
97 pacing_limited: false,
98 }
99 }
100
101 pub fn get_next_release_time(&self) -> ReleaseDecision {
102 if !self.enabled {
103 return ReleaseDecision {
104 time: ReleaseTime::Immediate,
105 allow_burst: true,
106 };
107 }
108
109 let allow_burst = self.burst_tokens > 0 || self.lumpy_tokens > 0;
110 ReleaseDecision {
111 time: self.ideal_next_packet_send_time,
112 allow_burst,
113 }
114 }
115
116 #[cfg(feature = "qlog")]
117 pub fn state_str(&self) -> &'static str {
118 self.sender.state_str()
119 }
120
121 pub fn get_congestion_window(&self) -> usize {
122 self.sender.get_congestion_window()
123 }
124
125 pub fn on_packet_sent(
126 &mut self, sent_time: Instant, bytes_in_flight: usize,
127 packet_number: u64, bytes: usize, is_retransmissible: bool,
128 rtt_stats: &RttStats,
129 ) {
130 self.sender.on_packet_sent(
131 sent_time,
132 bytes_in_flight,
133 packet_number,
134 bytes,
135 is_retransmissible,
136 );
137
138 if !self.enabled || !is_retransmissible {
139 return;
140 }
141
142 if bytes_in_flight == 0 && !self.sender.is_in_recovery() {
144 self.burst_tokens = self
147 .initial_burst_size
148 .min(self.sender.get_congestion_window_in_packets());
149 }
150
151 if self.burst_tokens > 0 {
152 self.burst_tokens -= 1;
153 self.ideal_next_packet_send_time = ReleaseTime::Immediate;
154 self.pacing_limited = false;
155 return;
156 }
157
158 let delay = self
162 .pacing_rate(bytes_in_flight + bytes, rtt_stats)
163 .transfer_time(bytes as u64);
164
165 if !self.pacing_limited || self.lumpy_tokens == 0 {
166 self.lumpy_tokens = 1.max(LUMPY_PACING_SIZE.min(
169 (self.sender.get_congestion_window_in_packets() as f64 *
170 LUMPY_PACING_CWND_FRACTION) as usize,
171 ));
172
173 if self.sender.bandwidth_estimate(rtt_stats) <
174 LUMPY_PACING_MIN_BANDWIDTH_KBPS
175 {
176 self.lumpy_tokens = 1;
179 }
180
181 if bytes_in_flight + bytes >= self.sender.get_congestion_window() {
182 self.lumpy_tokens = 1;
185 }
186 }
187
188 self.lumpy_tokens -= 1;
189 self.ideal_next_packet_send_time.set_max(sent_time);
190 self.ideal_next_packet_send_time.inc(delay);
191 self.pacing_limited = self.sender.can_send(bytes_in_flight + bytes);
193 }
194
195 #[allow(clippy::too_many_arguments)]
196 #[inline]
197 pub fn on_congestion_event(
198 &mut self, rtt_updated: bool, prior_in_flight: usize,
199 bytes_in_flight: usize, event_time: Instant, acked_packets: &[Acked],
200 lost_packets: &[Lost], least_unacked: u64, rtt_stats: &RttStats,
201 recovery_stats: &mut RecoveryStats,
202 ) {
203 self.sender.on_congestion_event(
204 rtt_updated,
205 prior_in_flight,
206 bytes_in_flight,
207 event_time,
208 acked_packets,
209 lost_packets,
210 least_unacked,
211 rtt_stats,
212 recovery_stats,
213 );
214
215 if !self.enabled {
216 return;
217 }
218
219 if !lost_packets.is_empty() {
220 self.burst_tokens = 0;
222 }
223
224 if let Some(max_pacing_rate) = self.max_pacing_rate {
225 if rtt_updated {
226 let max_rate = max_pacing_rate * 1.25f32;
227 let max_cwnd =
228 max_rate.to_bytes_per_period(rtt_stats.smoothed_rtt);
229 self.sender.limit_cwnd(max_cwnd as usize);
230 }
231 }
232 }
233
234 pub fn on_packet_neutered(&mut self, packet_number: u64) {
235 self.sender.on_packet_neutered(packet_number);
236 }
237
238 pub fn on_retransmission_timeout(&mut self, packets_retransmitted: bool) {
239 self.sender.on_retransmission_timeout(packets_retransmitted)
240 }
241
242 pub fn pacing_rate(
243 &self, bytes_in_flight: usize, rtt_stats: &RttStats,
244 ) -> Bandwidth {
245 let sender_rate = self.sender.pacing_rate(bytes_in_flight, rtt_stats);
246 match self.max_pacing_rate {
247 Some(rate) if self.enabled => rate.min(sender_rate),
248 _ => sender_rate,
249 }
250 }
251
252 pub fn bandwidth_estimate(&self, rtt_stats: &RttStats) -> Bandwidth {
253 self.sender.bandwidth_estimate(rtt_stats)
254 }
255
256 pub fn max_bandwidth(&self) -> Bandwidth {
257 self.sender.max_bandwidth()
258 }
259
260 pub fn rtt_persistent_jump_count(&self) -> u64 {
261 self.sender.rtt_persistent_jump_count()
262 }
263
264 #[cfg(feature = "qlog")]
265 pub fn send_rate(&self) -> Option<Bandwidth> {
266 self.sender.send_rate()
267 }
268
269 #[cfg(feature = "qlog")]
270 pub fn ack_rate(&self) -> Option<Bandwidth> {
271 self.sender.ack_rate()
272 }
273
274 pub fn on_app_limited(&mut self, bytes_in_flight: usize) {
275 self.pacing_limited = false;
276 self.sender.on_app_limited(bytes_in_flight);
277 }
278
279 pub fn update_mss(&mut self, new_mss: usize) {
280 self.sender.update_mss(new_mss)
281 }
282
283 #[cfg(feature = "qlog")]
284 pub fn ssthresh(&self) -> Option<u64> {
285 self.sender.ssthresh()
286 }
287
288 #[cfg(any(test, feature = "qlog"))]
289 pub fn is_app_limited(&self, bytes_in_flight: usize) -> bool {
290 !self.is_cwnd_limited(bytes_in_flight)
291 }
292
293 #[cfg(any(test, feature = "qlog"))]
294 fn is_cwnd_limited(&self, bytes_in_flight: usize) -> bool {
295 !self.pacing_limited && self.sender.is_cwnd_limited(bytes_in_flight)
296 }
297}