rtc_interceptor/pacing/
pacer.rs1use std::time::{Duration, Instant};
4
5pub const MIN_BURST_BITS: f64 = 8.0 * 1500.0;
11
12#[derive(Debug, Clone)]
18pub struct Pacer {
19 bitrate: f64,
21 configured_burst_bits: Option<f64>,
29 budget_bits: f64,
31 last_refill: Option<Instant>,
33}
34
35impl Pacer {
36 pub fn new(bits_per_second: f64) -> Self {
41 Self {
42 bitrate: bits_per_second.max(0.0),
43 configured_burst_bits: None,
44 budget_bits: Self::burst_for(bits_per_second),
45 last_refill: None,
46 }
47 }
48
49 pub fn with_burst_bits(mut self, burst_bits: f64) -> Self {
55 self.configured_burst_bits = Some(burst_bits.max(MIN_BURST_BITS));
56 self.budget_bits = self.budget_bits.min(self.burst_bits());
57 self
58 }
59
60 fn burst_for(bits_per_second: f64) -> f64 {
62 (bits_per_second / 10.0).max(MIN_BURST_BITS)
63 }
64
65 pub fn target_bitrate(&self) -> f64 {
67 self.bitrate
68 }
69
70 pub fn set_target_bitrate(&mut self, bits_per_second: f64) {
80 self.bitrate = bits_per_second.max(0.0);
81 self.budget_bits = self.budget_bits.min(self.burst_bits());
82 }
83
84 pub fn refill(&mut self, now: Instant) {
86 if let Some(last) = self.last_refill {
87 let elapsed = now.saturating_duration_since(last).as_secs_f64();
88 self.budget_bits = (self.budget_bits + elapsed * self.bitrate).min(self.burst_bits());
89 }
90 self.last_refill = Some(now);
91 }
92
93 pub fn can_afford(&self, bits: f64) -> bool {
95 self.budget_bits >= bits
96 }
97
98 pub fn consume(&mut self, bits: f64) {
104 self.budget_bits -= bits;
105 }
106
107 pub fn time_until_affordable(&self, bits: f64) -> Option<Duration> {
112 if self.budget_bits >= bits {
113 return Some(Duration::ZERO);
114 }
115 if self.bitrate <= 0.0 {
116 return None;
117 }
118 Some(Duration::from_secs_f64(
119 (bits - self.budget_bits) / self.bitrate,
120 ))
121 }
122
123 pub fn affordable_at(&self, bits: f64) -> Option<Instant> {
125 let last = self.last_refill?;
126 self.time_until_affordable(bits).map(|wait| last + wait)
127 }
128
129 pub fn budget_bits(&self) -> f64 {
131 self.budget_bits
132 }
133
134 pub fn can_release(&self, bits: f64) -> bool {
142 let burst_bits = self.burst_bits();
143 if bits > burst_bits {
144 return self.budget_bits >= burst_bits;
145 }
146 self.budget_bits >= bits
147 }
148
149 pub fn releasable_at(&self, bits: f64) -> Option<Instant> {
151 self.affordable_at(bits.min(self.burst_bits()))
152 }
153
154 pub fn burst_bits(&self) -> f64 {
160 self.configured_burst_bits
161 .unwrap_or_else(|| Self::burst_for(self.bitrate))
162 }
163}
164
165#[cfg(test)]
166mod tests {
167 use super::*;
168
169 fn pacer() -> Pacer {
171 Pacer::new(1_000_000.0).with_burst_bits(MIN_BURST_BITS)
172 }
173
174 #[test]
175 fn a_new_bucket_starts_full() {
176 let pacer = pacer();
177 assert_eq!(MIN_BURST_BITS, pacer.budget_bits());
178 assert!(pacer.can_afford(MIN_BURST_BITS));
179 }
180
181 #[test]
182 fn spending_reduces_the_budget() {
183 let mut pacer = pacer();
184 pacer.consume(1000.0);
185 assert_eq!(MIN_BURST_BITS - 1000.0, pacer.budget_bits());
186 }
187
188 #[test]
190 fn the_budget_refills_from_elapsed_time() {
191 let now = Instant::now();
192 let mut pacer = pacer();
193 pacer.refill(now);
194 pacer.consume(pacer.budget_bits());
195 assert_eq!(0.0, pacer.budget_bits());
196
197 pacer.refill(now + Duration::from_millis(10));
199 assert!(
200 (pacer.budget_bits() - 10_000.0).abs() < 1.0,
201 "{}",
202 pacer.budget_bits()
203 );
204 }
205
206 #[test]
207 fn refilling_in_steps_matches_refilling_at_once() {
208 let now = Instant::now();
209
210 let mut stepped = pacer();
211 stepped.refill(now);
212 stepped.consume(stepped.budget_bits());
213 for step in 1..=10 {
214 stepped.refill(now + Duration::from_millis(step));
215 }
216
217 let mut at_once = pacer();
218 at_once.refill(now);
219 at_once.consume(at_once.budget_bits());
220 at_once.refill(now + Duration::from_millis(10));
221
222 assert!((stepped.budget_bits() - at_once.budget_bits()).abs() < 1.0);
223 }
224
225 #[test]
228 fn the_budget_is_capped_at_the_burst() {
229 let now = Instant::now();
230 let mut pacer = pacer();
231 pacer.refill(now);
232 pacer.refill(now + Duration::from_secs(60));
233
234 assert_eq!(MIN_BURST_BITS, pacer.budget_bits());
235 }
236
237 #[test]
238 fn the_time_until_affordable_is_zero_when_it_already_is() {
239 let pacer = pacer();
240 assert_eq!(Some(Duration::ZERO), pacer.time_until_affordable(100.0));
241 }
242
243 #[test]
244 fn the_time_until_affordable_scales_with_the_shortfall() {
245 let now = Instant::now();
246 let mut pacer = pacer();
247 pacer.refill(now);
248 pacer.consume(pacer.budget_bits());
249
250 let wait = pacer
252 .time_until_affordable(10_000.0)
253 .expect("a finite wait");
254 assert!((wait.as_secs_f64() - 0.010).abs() < 0.0005, "got {wait:?}");
255 assert_eq!(Some(now + wait), pacer.affordable_at(10_000.0));
256 }
257
258 #[test]
261 fn nothing_becomes_affordable_at_a_zero_rate() {
262 let now = Instant::now();
263 let mut pacer = Pacer::new(0.0);
264 pacer.refill(now);
265 pacer.consume(pacer.budget_bits());
266
267 assert_eq!(None, pacer.time_until_affordable(1000.0));
268 assert_eq!(None, pacer.affordable_at(1000.0));
269 }
270
271 #[test]
272 fn changing_the_rate_changes_how_fast_the_budget_refills() {
273 let now = Instant::now();
274 let mut pacer = Pacer::new(1_000_000.0);
276 pacer.refill(now);
277 pacer.consume(pacer.budget_bits());
278
279 pacer.set_target_bitrate(2_000_000.0);
280 pacer.refill(now + Duration::from_millis(10));
281
282 assert!(
284 (pacer.budget_bits() - 20_000.0).abs() < 1.0,
285 "{}",
286 pacer.budget_bits()
287 );
288 assert_eq!(2_000_000.0, pacer.target_bitrate());
289 }
290
291 #[test]
294 fn lowering_the_rate_clamps_the_budget_to_the_new_burst() {
295 let now = Instant::now();
296 let mut pacer = Pacer::new(100_000_000.0);
297 pacer.refill(now);
298 let before = pacer.budget_bits();
299
300 pacer.set_target_bitrate(1000.0);
301
302 assert!(pacer.budget_bits() < before);
303 assert_eq!(
304 MIN_BURST_BITS,
305 pacer.budget_bits(),
306 "clamped to the floor burst"
307 );
308 }
309
310 #[test]
311 fn a_negative_rate_is_treated_as_zero() {
312 let mut pacer = Pacer::new(-5.0);
313 assert_eq!(0.0, pacer.target_bitrate());
314 pacer.set_target_bitrate(-1.0);
315 assert_eq!(0.0, pacer.target_bitrate());
316 }
317
318 #[test]
321 fn a_configured_burst_survives_a_rate_change() {
322 let mut pacer = Pacer::new(1_000_000.0).with_burst_bits(MIN_BURST_BITS);
323 assert_eq!(MIN_BURST_BITS, pacer.burst_bits());
324
325 pacer.set_target_bitrate(100_000_000.0);
326
327 assert_eq!(
328 MIN_BURST_BITS,
329 pacer.burst_bits(),
330 "a rate change must not widen a burst the caller set"
331 );
332 assert_eq!(100_000_000.0, pacer.target_bitrate());
333 }
334
335 #[test]
337 fn a_derived_burst_follows_the_rate() {
338 let mut pacer = Pacer::new(1_000_000.0);
339 assert_eq!(100_000.0, pacer.burst_bits());
340
341 pacer.set_target_bitrate(2_000_000.0);
342
343 assert_eq!(200_000.0, pacer.burst_bits());
344 }
345
346 #[test]
349 fn an_oversized_packet_waits_for_a_full_budget() {
350 let now = Instant::now();
351 let mut pacer = pacer();
352 pacer.refill(now);
353 let oversized = MIN_BURST_BITS * 2.0;
354
355 assert!(pacer.can_release(oversized), "a full budget releases it");
356 pacer.consume(oversized);
357 assert!(
358 !pacer.can_release(oversized),
359 "the next one waits for the debt to be repaid"
360 );
361
362 let at = pacer.releasable_at(oversized).expect("a finite wait");
364 assert_eq!(
365 Some(Duration::from_secs_f64(oversized / 1_000_000.0)),
366 pacer.time_until_affordable(MIN_BURST_BITS)
367 );
368 pacer.refill(at);
369 assert!(pacer.can_release(oversized), "and then it goes");
370 }
371
372 #[test]
375 fn a_packet_larger_than_the_burst_can_still_be_sent() {
376 let now = Instant::now();
377 let mut pacer = pacer();
378 pacer.refill(now);
379
380 let oversized = MIN_BURST_BITS * 2.0;
381 assert!(!pacer.can_afford(oversized));
382
383 pacer.refill(now + Duration::from_secs(10));
385 assert!(!pacer.can_afford(oversized));
386
387 pacer.consume(oversized);
388 assert!(
389 pacer.budget_bits() < 0.0,
390 "the overshoot is paid back over time"
391 );
392
393 pacer.refill(now + Duration::from_secs(20));
394 assert_eq!(
395 MIN_BURST_BITS,
396 pacer.budget_bits(),
397 "and recovers to the burst"
398 );
399 }
400}