Skip to content
File

Blob: firmware/vendor/str0m/src/pacer/leaky.rs

rust1371 lines
1use std::collections::VecDeque;
2use std::time::{Duration, Instant};
3 
4use super::Pacer;
5use super::PaddingRequest;
6use super::QueueState;
7use crate::Reason;
8use crate::bwe_::ProbeClusterConfig;
9use crate::bwe_::ProbeClusterState;
10use crate::bwe_::{log_pacer_media_debt, log_pacer_padding_debt};
11use crate::pacer::PacerReason;
12use crate::rtp_::{Bitrate, DataSize, MidRid, TwccClusterId};
13use crate::util::Soonest;
14 
15const MAX_BITRATE: Bitrate = Bitrate::gbps(10);
16const MAX_DEBT_IN_TIME: Duration = Duration::from_millis(500);
17const PADDING_BURST_INTERVAL: Duration = Duration::from_millis(5);
18const PACING: Duration = Duration::from_millis(40);
19 
20/// A leaky bucket pacer that can overshoot the target bitrate when required.
21pub struct LeakyBucketPacer {
22 /// Pacing bitrate.
23 pacing_bitrate: Bitrate,
24 /// Adjusted pacing bitrate for when we need to drain queues.
25 adjusted_bitrate: Bitrate,
26 /// The bitrate at which to send padding packets when the pacing rate isn't being achieved.
27 padding_bitrate: Bitrate,
28 /// The last time we refreshed media debt and potentially adjusted the bitrate.
29 last_handle_time: Option<Instant>,
30 /// The last time we indicated that a packet should be sent.
31 last_emitted: Option<Instant>,
32 /// The next time we should send a queued packet.
33 next_poll_time: Option<(Instant, PacerReason)>,
34 /// The current media debt.
35 media_debt: DataSize,
36 /// The current padding debt.
37 padding_debt: DataSize,
38 /// The longest the average packet can spend in the queue before we force it to be drained.
39 queue_limit: Duration,
40 /// The queue states given by last handle_timeout.
41 queue_states: Vec<QueueState>,
42 /// The next return value for `poll_queue``
43 next_poll_queue: Option<MidRid>,
44 /// Queue of probe clusters waiting to be executed
45 probe_queue: VecDeque<ProbeClusterState>,
46 /// Last completed probe cluster (to be consumed by check_probe_complete)
47 completed_probe: Option<TwccClusterId>,
48 /// Gates poll_queue() until handle_timeout() is called after packet emission.
49 needs_timeout_before_next_poll: bool,
50 /// Caches whether we have any queue to send padding on (RTX).
51 has_padding_queue: bool,
52}
53 
54impl Pacer for LeakyBucketPacer {
55 fn set_pacing_rate(&mut self, pacing_bitrate: Bitrate) {
56 self.pacing_bitrate = pacing_bitrate;
57 
58 // bitrate will be updated on next handle_timeout().
59 }
60 
61 fn set_padding_rate(&mut self, padding_bitrate: Bitrate) {
62 self.padding_bitrate = padding_bitrate;
63 
64 // bitrate will be updated on next handle_timeout().
65 }
66 
67 fn poll_timeout(&self) -> (Option<Instant>, Reason) {
68 let next_handle_time = self.last_handle_time.map(|lh| lh + PACING);
69 
70 let poll_at = self
71 .next_poll_time
72 .map(|(t, r)| (Some(t), Reason::Pacer(r)))
73 .unwrap_or((None, Reason::NotHappening));
74 
75 (next_handle_time, Reason::Pacer(PacerReason::Handle)).soonest(poll_at)
76 }
77 
78 fn handle_timeout(
79 &mut self,
80 now: Instant,
81 iter: impl Iterator<Item = QueueState>,
82 ) -> Option<PaddingRequest> {
83 // Clear the gate when time advances
84 self.needs_timeout_before_next_poll = false;
85 
86 // Clear the poll time - it will be recalculated below if needed.
87 // This is important because if we return early (e.g., next_poll_queue is already set),
88 // we don't want the old Immediate timeout to keep firing.
89 self.next_poll_time = None;
90 
91 // This is called periodically and whenever packet is queued.
92 self.queue_states.clear();
93 self.queue_states.extend(iter);
94 
95 let elapsed = self.update_handle_time_and_get_elapsed(now);
96 
97 self.clear_debt(elapsed);
98 self.maybe_update_adjusted_bitrate(now);
99 
100 if let Some(request) = self.maybe_create_padding_request(now) {
101 self.next_poll_queue = Some(request.midrid);
102 return Some(request);
103 }
104 
105 if self.next_poll_queue.is_some() {
106 return None;
107 }
108 
109 let (next_poll_time_and_reason, queue) = self.next_poll(now)?;
110 
111 if now < next_poll_time_and_reason.0 {
112 // We don't set this because between now and the neaxt poll, queue state can change such
113 // that we should poll a different queue i.e. media could be queued.
114 self.next_poll_queue = None;
115 } else {
116 self.next_poll_queue = queue.map(|q| q.midrid);
117 }
118 
119 self.next_poll_time = Some(next_poll_time_and_reason);
120 
121 None
122 }
123 
124 fn poll_queue(&mut self) -> Option<(MidRid, Option<TwccClusterId>)> {
125 // GATE: Block if we need timeout first
126 if self.needs_timeout_before_next_poll {
127 return None;
128 }
129 
130 let next = self.next_poll_queue.take()?;
131 
132 // Mark that we need timeout before next poll
133 self.needs_timeout_before_next_poll = true;
134 self.request_immediate_timeout();
135 
136 // Capture the cluster ID at poll time, before register_send() might clear it
137 let cluster_id = self.active_cluster();
138 
139 Some((next, cluster_id))
140 }
141 
142 fn register_send(&mut self, now: Instant, packet_size: DataSize, _from: MidRid) {
143 self.last_emitted = Some(now);
144 
145 self.media_debt += packet_size;
146 self.media_debt = self
147 .media_debt
148 .min(self.adjusted_bitrate * MAX_DEBT_IN_TIME);
149 log_pacer_media_debt!(self.media_debt.as_bytes_usize());
150 self.add_padding_debt(packet_size);
151 
152 // Update active probe state to track this packet
153 // This ensures probe timing advances correctly even when sending media packets
154 if let Some(probe) = self.probe_queue.front_mut() {
155 probe.record_packet(now, packet_size);
156 }
157 
158 // Check if probe is complete and store it for later retrieval
159 if let Some(cluster_id) = self.check_probe_complete_internal(now) {
160 self.completed_probe = Some(cluster_id);
161 }
162 }
163 
164 fn has_padding_queue(&self) -> bool {
165 self.has_padding_queue
166 }
167}
168 
169impl LeakyBucketPacer {
170 pub fn new(initial_pacing_bitrate: Bitrate) -> Self {
171 const DEFAULT_QUEUE_LIMIT: Duration = Duration::from_secs(2);
172 
173 Self {
174 pacing_bitrate: initial_pacing_bitrate,
175 adjusted_bitrate: Bitrate::ZERO,
176 padding_bitrate: Bitrate::ZERO,
177 last_handle_time: None,
178 last_emitted: None,
179 next_poll_time: None,
180 media_debt: DataSize::ZERO,
181 padding_debt: DataSize::ZERO,
182 queue_limit: DEFAULT_QUEUE_LIMIT,
183 queue_states: vec![],
184 next_poll_queue: None,
185 probe_queue: VecDeque::new(),
186 completed_probe: None,
187 needs_timeout_before_next_poll: true,
188 has_padding_queue: false,
189 }
190 }
191 
192 /// Start executing a probe cluster.
193 ///
194 /// The pacer will pace at the probe's target bitrate and track packets sent.
195 /// Probes are queued and executed sequentially.
196 pub(crate) fn start_probe(&mut self, config: ProbeClusterConfig) {
197 trace!(?config, "Probe start");
198 self.probe_queue.push_back(ProbeClusterState::new(config));
199 }
200 
201 /// Get the cluster ID of the active probe, if any.
202 pub(crate) fn active_cluster(&self) -> Option<TwccClusterId> {
203 self.probe_queue.front().map(|p| p.config().cluster())
204 }
205 
206 /// Check if the active probe is complete and should be finished.
207 pub(crate) fn check_probe_complete(&mut self, now: Instant) -> Option<TwccClusterId> {
208 // Check if we have a completed probe from a previous call
209 if let Some(cluster_id) = self.completed_probe.take() {
210 return Some(cluster_id);
211 }
212 
213 // Otherwise check if the active probe just completed
214 self.check_probe_complete_internal(now)
215 }
216 
217 /// Internal method to check if probe is complete (doesn't consume completed_probe)
218 fn check_probe_complete_internal(&mut self, now: Instant) -> Option<TwccClusterId> {
219 let probe = self.probe_queue.front()?;
220 
221 if probe.is_complete(now) {
222 let cluster_id = probe.config().cluster();
223 self.probe_queue.pop_front();
224 return Some(cluster_id);
225 }
226 
227 None
228 }
229 
230 fn update_handle_time_and_get_elapsed(&mut self, now: Instant) -> Duration {
231 // Due the calling code this also happens when a packet is queued in any upstream queue.
232 let Some(previous_handle_time) = self.last_handle_time else {
233 self.last_handle_time = Some(now);
234 return Duration::ZERO;
235 };
236 
237 let elapsed = now - previous_handle_time;
238 self.last_handle_time = Some(now);
239 
240 elapsed
241 }
242 
243 fn clear_debt(&mut self, elapsed: Duration) {
244 self.media_debt = self
245 .media_debt
246 .saturating_sub(self.adjusted_bitrate * elapsed);
247 self.padding_debt = self
248 .padding_debt
249 .saturating_sub(self.padding_bitrate * elapsed);
250 log_pacer_media_debt!(self.media_debt.as_bytes_usize());
251 log_pacer_padding_debt!(self.padding_debt.as_bytes_usize());
252 }
253 
254 fn next_poll(&self, now: Instant) -> Option<((Instant, PacerReason), Option<&QueueState>)> {
255 // If we have never sent before, do so immediately on an arbitrary non-empty queue.
256 if self.last_emitted.is_none() {
257 let mut queues = self
258 .queue_states
259 .iter()
260 .filter(|q| q.snapshot.packet_count > 0);
261 
262 return queues
263 .next()
264 .map(|q| ((now, PacerReason::FirstEver), Some(q)));
265 };
266 
267 let unpaced = self
268 .queue_states
269 .iter()
270 .filter(|qs| qs.unpaced)
271 .filter_map(|qs| qs.snapshot.first_unsent.map(|t| (t, qs)))
272 .min_by_key(|(t, _)| *t);
273 
274 // Unpaced packets (such as audio by default) are sent immediately.
275 if let Some((queued_at, qs)) = unpaced {
276 return Some(((queued_at, PacerReason::Unpaced), Some(qs)));
277 }
278 
279 let non_empty_queue = {
280 let non_empty_queues = self
281 .queue_states
282 .iter()
283 .filter(|q| q.snapshot.packet_count > 0);
284 
285 // Send on the non-empty queue with the lowest priority that, was least recently
286 // sent on.
287 non_empty_queues.min_by_key(|q| (q.snapshot.priority, q.snapshot.last_emitted))
288 };
289 
290 if let Some(queue) = non_empty_queue {
291 if self.adjusted_bitrate > Bitrate::ZERO {
292 // Check if we're actively probing and should use probe-specific timing
293 let poll_at = if let Some(probe) = self.probe_queue.front() {
294 // During probe: use absolute time directly from probe state
295 (probe.next_probe_time(), PacerReason::Probe1)
296 } else {
297 // Normal pacing: use relative offset based on debt
298 let drain_debt_time = self.media_debt / self.adjusted_bitrate;
299 let next_send_offset = if drain_debt_time > PACING {
300 // If we have incurred too much debt we need to wait to let it clear out before sending
301 // again.
302 drain_debt_time
303 } else {
304 Duration::ZERO
305 };
306 
307 let time = self
308 .last_handle_time
309 .map(|h| h + next_send_offset)
310 .unwrap_or(now);
311 
312 (time, PacerReason::Paced)
313 };
314 
315 return Some((poll_at, Some(queue)));
316 }
317 }
318 
319 let any_queue_for_padding = self.queue_states.iter().any(|q| q.use_for_padding);
320 let padding_possible = self.padding_bitrate > Bitrate::ZERO && any_queue_for_padding;
321 
322 if !padding_possible {
323 return None;
324 }
325 
326 // If we're actively probing, use probe timing for padding
327 if let Some(probe) = self.probe_queue.front() {
328 let next_probe_time = probe.next_probe_time();
329 // We explicitly don't return a queue to poll here. We need another call to
330 // handle_timeout to request the padding before we can poll the selected queue.
331 return Some(((next_probe_time, PacerReason::Probe2), None));
332 }
333 
334 // If all queues are empty and we have a padding rate, wait until we have drained
335 // both the media debt and padding debt to send some padding.
336 let mut drain_debt_time =
337 (self.media_debt / self.adjusted_bitrate).max(self.padding_debt / self.padding_bitrate);
338 if drain_debt_time.is_zero() {
339 // Give the main loop some time to do something else e.g. queue media.
340 drain_debt_time = Duration::from_micros(1);
341 }
342 
343 let padding_at = self
344 .last_handle_time
345 .map(|h| h + drain_debt_time)
346 .unwrap_or(now);
347 
348 // We explicitly don't return a queue to poll here. We need another call to
349 // handle_timeout to request the padding before we can poll the selected queue.
350 Some(((padding_at, PacerReason::Padding), None))
351 }
352 
353 fn maybe_update_adjusted_bitrate(&mut self, now: Instant) {
354 // Use probe's target bitrate if actively probing, otherwise use pacing bitrate
355 self.adjusted_bitrate = if let Some(probe) = self.probe_queue.front() {
356 probe.config().target_bitrate()
357 } else {
358 self.pacing_bitrate
359 };
360 
361 let (queue_time, queued_packets, queue_size) =
362 self.queue_states
363 .iter()
364 .fold((Duration::ZERO, 0, DataSize::ZERO), |acc, q| {
365 (
366 acc.0 + q.snapshot.total_queue_time(now),
367 acc.1 + q.snapshot.packet_count,
368 acc.2 + DataSize::from(q.snapshot.byte_size),
369 )
370 });
371 if queued_packets == 0 {
372 return;
373 }
374 
375 let avg_queue_time = queue_time / queued_packets;
376 
377 // The average time we want the packet in the queue to at most to wait to drain.
378 let target_queue_wait =
379 Duration::from_millis(1).max(self.queue_limit.saturating_sub(avg_queue_time));
380 // Min data rate to drain what's currently in the queue.
381 let min_rate = queue_size / target_queue_wait;
382 if min_rate > self.adjusted_bitrate {
383 // Min rate exceeds our pacing rate, increase the rate to force drain the queue.
384 self.adjusted_bitrate = min_rate.clamp(Bitrate::ZERO, MAX_BITRATE);
385 }
386 }
387 
388 fn add_padding_debt(&mut self, size: DataSize) {
389 self.padding_debt += size;
390 self.padding_debt = self
391 .padding_debt
392 .min(self.padding_bitrate * MAX_DEBT_IN_TIME);
393 log_pacer_padding_debt!(self.padding_debt.as_bytes_usize());
394 }
395 
396 /// Optimistically attempt to create a padding request.
397 ///
398 /// Returns `Some(PaddingRequest)` if padding is enabled and the current queue state
399 /// allows padding, otherwise returns `None`.
400 fn maybe_create_padding_request(&mut self, now: Instant) -> Option<PaddingRequest> {
401 // Queues must be empty.
402 let all_queues_empty = self
403 .queue_states
404 .iter()
405 .all(|q| q.snapshot.packet_count == 0);
406 if !all_queues_empty {
407 return None;
408 }
409 
410 // We must have a queue that supports padding.
411 let maybe_queue = self
412 .queue_states
413 .iter()
414 .filter(|q| q.use_for_padding)
415 .max_by_key(|q| q.snapshot.last_emitted);
416 
417 // Save whether we have a valid padding queue.
418 self.has_padding_queue = maybe_queue.is_some();
419 
420 if !self.has_padding_queue {
421 // No padding queue, no probes.
422 self.probe_queue.clear();
423 }
424 
425 let queue = maybe_queue?;
426 
427 // Check for PROBE padding FIRST (bypasses debt checks)
428 // Active probes need padding to hit their target bitrate when there's insufficient media.
429 if let Some(probe) = self.probe_queue.front_mut() {
430 // Delegate probe timing and padding calculation to ProbeClusterState
431 if !probe.should_send_now(now) {
432 // Not time yet - wait until next_probe_time
433 return None;
434 }
435 
436 // Get recommended padding amount from ProbeClusterState
437 // This handles the calculation of how much padding is needed based on probe timing
438 let padding_size = probe.next_packet(now);
439 let Some(padding_size) = padding_size else {
440 // Probe says no padding needed (already sent enough for this interval)
441 return None;
442 };
443 
444 return Some(PaddingRequest {
445 midrid: queue.midrid,
446 padding: padding_size.as_bytes_usize(),
447 });
448 }
449 
450 // Normal padding: requires zero debt
451 if self.media_debt != DataSize::ZERO || self.padding_debt != DataSize::ZERO {
452 return None;
453 }
454 
455 if self.padding_bitrate == Bitrate::ZERO {
456 return None;
457 }
458 
459 // We can generate padding
460 let padding = (self.padding_bitrate * PADDING_BURST_INTERVAL).as_bytes_usize();
461 
462 Some(PaddingRequest {
463 midrid: queue.midrid,
464 padding,
465 })
466 }
467 
468 fn request_immediate_timeout(&mut self) {
469 // Request timeout at the next microsecond to ensure time advances between packets.
470 // We can't use already_happened() because that would cause the test harness to
471 // set a very old timestamp, and while lib.rs prevents last_now from going backwards,
472 // it doesn't force it to advance, so all packets would get the same timestamp.
473 const MINIMAL_DELTA: Duration = Duration::from_micros(1);
474 
475 let Some(time) = self.last_handle_time.map(|t| t + MINIMAL_DELTA) else {
476 self.next_poll_time = None;
477 return;
478 };
479 
480 self.next_poll_time = Some((time, PacerReason::Immediate));
481 }
482}
483 
484#[cfg(test)]
485mod test {
486 use super::super::{QueuePriority, QueueSnapshot};
487 use super::*;
488 use crate::rtp_::{DataSize, Mid, RtpHeader};
489 use queue::{PacketKind, Queue, QueuedPacket};
490 use std::time::{Duration, Instant};
491 
492 #[test]
493 fn test_typical_behavior() {
494 let now = Instant::now();
495 let mut queue = Queue::default();
496 // 2,000 bits per second, 10 bytes per pacing interval(40ms)
497 let mut pacer = LeakyBucketPacer::new((10 * 200).into());
498 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(1));
499 
500 assert!(
501 pacer.poll_queue().is_none(),
502 "We initially attempt to poll any non-empty queue if we have never sent",
503 );
504 
505 enqueue_packet_noisy(
506 &mut pacer,
507 &mut queue,
508 1,
509 5,
510 PacketKind::Video,
511 now + duration_ms(21),
512 );
513 
514 assert_poll_success(
515 &mut pacer,
516 &mut queue,
517 now + duration_ms(21),
518 "First packet should be released because we have no debt",
519 |packet| {
520 assert_eq!(packet.header.sequence_number, 1);
521 },
522 );
523 
524 enqueue_packet_noisy(
525 &mut pacer,
526 &mut queue,
527 2,
528 8,
529 PacketKind::Video,
530 now + duration_ms(27),
531 );
532 enqueue_packet_noisy(
533 &mut pacer,
534 &mut queue,
535 3,
536 25,
537 PacketKind::Video,
538 now + duration_ms(28),
539 );
540 
541 assert_poll_success(
542 &mut pacer,
543 &mut queue,
544 now + duration_ms(28),
545 "Second packet should be released because the debt is within tolerance",
546 |packet| {
547 assert_eq!(packet.header.sequence_number, 2);
548 },
549 );
550 
551 // We have incurred too much media debt so polling will now fail until the debt can be
552 // reduced.
553 assert!(
554 pacer.poll_queue().is_none(),
555 "Third packet should not be released because we have too much debt"
556 );
557 
558 // Periodic timeout
559 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(41));
560 
561 assert_poll_success(
562 &mut pacer,
563 &mut queue,
564 now + duration_ms(41),
565 "Third packet should be released because we have cleared debt as time moved forward",
566 |packet| {
567 assert_eq!(packet.header.sequence_number, 3);
568 },
569 );
570 
571 enqueue_packet_noisy(
572 &mut pacer,
573 &mut queue,
574 4,
575 12,
576 PacketKind::Video,
577 now + duration_ms(45),
578 );
579 enqueue_packet_noisy(
580 &mut pacer,
581 &mut queue,
582 5,
583 25,
584 PacketKind::Video,
585 now + duration_ms(47),
586 );
587 
588 // We have incurred too much media debt so polling will now fail until the debt can be
589 // reduced.
590 assert!(
591 pacer.poll_queue().is_none(),
592 "Fourth packet should not be released because we have too much debt"
593 );
594 
595 enqueue_packet_noisy(
596 &mut pacer,
597 &mut queue,
598 6,
599 100,
600 PacketKind::Audio,
601 now + duration_ms(52),
602 );
603 
604 // Unpaced packets should be able to send even if we have too much media debt.
605 assert_poll_success(
606 &mut pacer,
607 &mut queue,
608 now + duration_ms(52),
609 "Sixth packet (audio) should be released despite too much media debt because \
610 audio packets are not paced",
611 |packet| {
612 assert_eq!(packet.kind, PacketKind::Audio);
613 assert_eq!(packet.header.sequence_number, 6);
614 },
615 );
616 
617 // A lot of time passes, now the bitrate should be adjusted to force drain the queues to
618 // avoid packets being queued for too long.
619 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(2053));
620 
621 assert_poll_success(
622 &mut pacer,
623 &mut queue,
624 now + duration_ms(2053),
625 "Fourth packet should be released after hitting the queue limit",
626 |packet| {
627 assert_eq!(packet.header.sequence_number, 4);
628 },
629 );
630 
631 assert_poll_success(
632 &mut pacer,
633 &mut queue,
634 now + duration_ms(2053),
635 "Fifth packet should be released after hitting the queue limit",
636 |packet| {
637 assert_eq!(packet.header.sequence_number, 5);
638 },
639 );
640 
641 assert!(queue.is_empty());
642 }
643 
644 #[test]
645 fn test_queue_drain() {
646 let now = Instant::now();
647 let mut queue = Queue::default();
648 // 2,000 bits per second, 10 bytes per pacing interval(40ms)
649 let mut pacer = LeakyBucketPacer::new((10 * 200).into());
650 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(1));
651 
652 enqueue_packet_noisy(
653 &mut pacer,
654 &mut queue,
655 1,
656 22,
657 PacketKind::Video,
658 now + duration_ms(21),
659 );
660 
661 assert_poll_success(
662 &mut pacer,
663 &mut queue,
664 now + duration_ms(21),
665 "First packet should be released because we have no debt",
666 |packet| {
667 assert_eq!(packet.header.sequence_number, 1);
668 },
669 );
670 
671 // Time moves forward
672 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(41));
673 
674 // Nothing happens for a while because there's nothing in the queues.
675 
676 enqueue_packet_noisy(
677 &mut pacer,
678 &mut queue,
679 2,
680 8,
681 PacketKind::Video,
682 // Debt will be just slightly above what can be drained in 40 ms
683 // after 66ms
684 now + duration_ms(66),
685 );
686 
687 assert!(
688 pacer.poll_queue().is_none(),
689 "Second packet should not be released because there's too much debt"
690 );
691 
692 enqueue_packet_noisy(
693 &mut pacer,
694 &mut queue,
695 3,
696 5,
697 PacketKind::Video,
698 now + duration_ms(70),
699 );
700 // Drain packet 2
701 assert_poll_success(
702 &mut pacer,
703 &mut queue,
704 now + duration_ms(70),
705 "Second packet should be released because of the adjusted bitrate to drain the queue",
706 |packet| {
707 assert_eq!(packet.header.sequence_number, 2);
708 },
709 );
710 
711 enqueue_packet_noisy(
712 &mut pacer,
713 &mut queue,
714 4,
715 1200,
716 PacketKind::Video,
717 now + duration_ms(71),
718 );
719 
720 // Drain packet 3
721 assert_poll_success(
722 &mut pacer,
723 &mut queue,
724 now + duration_ms(71),
725 "Third packet should be released because of the adjusted bitrate to drain the queue",
726 |packet| {
727 assert_eq!(packet.header.sequence_number, 3);
728 },
729 );
730 
731 // Drain packet 4
732 assert_poll_success(
733 &mut pacer,
734 &mut queue,
735 now + duration_ms(71),
736 "Fourth packet should be released because of the adjusted bitrate to drain the queue",
737 |packet| {
738 assert_eq!(packet.header.sequence_number, 4);
739 },
740 );
741 
742 // Time moves forward
743 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(81));
744 
745 enqueue_packet_noisy(
746 &mut pacer,
747 &mut queue,
748 5,
749 40,
750 PacketKind::Video,
751 now + duration_ms(81),
752 );
753 
754 assert!(
755 pacer.poll_queue().is_none(),
756 "Fifth packet shoud not be relaesed because there's too much debt"
757 );
758 }
759 
760 #[test]
761 fn test_padding_fill_in() {
762 let now = Instant::now();
763 let mut queue = Queue::default();
764 let mut pacer = LeakyBucketPacer::new((10 * 200).into());
765 // 2,000 bits per second, 10 bytes per pacing interval(40ms) with padding at 3,000 bits per
766 // second, 15 bytes per pacing interval(40ms)
767 pacer.set_pacing_rate((10 * 200).into());
768 pacer.set_padding_rate((15 * 200).into());
769 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(1));
770 
771 enqueue_packet_noisy(
772 &mut pacer,
773 &mut queue,
774 1,
775 22,
776 PacketKind::Video,
777 now + duration_ms(21),
778 );
779 
780 assert_poll_success(
781 &mut pacer,
782 &mut queue,
783 now + duration_ms(21),
784 "First packet should be released because we have no debt",
785 |packet| {
786 assert_eq!(packet.header.sequence_number, 1);
787 },
788 );
789 
790 // Time moves forward
791 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(41));
792 
793 // Nothing happens for a while because there's nothing in the queues.
794 
795 enqueue_packet_noisy(
796 &mut pacer,
797 &mut queue,
798 2,
799 8,
800 PacketKind::Video,
801 now + duration_ms(70),
802 );
803 
804 // Drain packet 2
805 assert_poll_success(
806 &mut pacer,
807 &mut queue,
808 now + duration_ms(70),
809 "Second packet should be released because of the adjusted bitrate to drain the queue",
810 |packet| {
811 assert_eq!(packet.header.sequence_number, 2);
812 },
813 );
814 
815 // Time moves forward, all debt is cleared out now
816 handle_timeout_noisy(&mut pacer, &mut queue, now + duration_ms(155));
817 
818 // Drain padding packet
819 assert_poll_success(
820 &mut pacer,
821 &mut queue,
822 now + duration_ms(165),
823 "The queued padding packet should be drained",
824 |packet| {
825 assert_eq!(packet.size(), 2);
826 assert_eq!(packet.header.sequence_number, 0);
827 },
828 );
829 
830 enqueue_packet_noisy(
831 &mut pacer,
832 &mut queue,
833 3,
834 15,
835 PacketKind::Video,
836 now + duration_ms(165),
837 );
838 
839 // Drain packet 3
840 assert_poll_success(
841 &mut pacer,
842 &mut queue,
843 now + duration_ms(165),
844 "Third packet should be released because the sent padding doesn't \
845 increase the media debt too much",
846 |packet| {
847 assert_eq!(packet.header.sequence_number, 3);
848 },
849 );
850 }
851 
852 #[test]
853 fn test_realistic() {
854 let config = RealisticTestConfig {
855 padding_rate: Bitrate::kbps(2500),
856 max_overshoot_factor: 0.05,
857 spike_probability: 3,
858 ..Default::default()
859 };
860 let (media_rate, padding_rate, total_rate) = run_realistic_test(config);
861 let expected_padding = config.padding_rate - config.media_rate;
862 // Expect result to be within 2 standard deviations.
863 let upper_bound =
864 config.media_rate + expected_padding * (1.0 + config.max_overshoot_factor * 2.0) as f64;
865 let lower_bound =
866 config.media_rate + expected_padding * (1.0 - config.max_overshoot_factor * 2.0) as f64;
867 
868 assert!(
869 total_rate >= lower_bound && total_rate <= upper_bound,
870 "Expected reuslting total rate to be within expected bounds. \
871 total_rate={total_rate}, media_rate={media_rate}, padding_rate={padding_rate}, \
872 config={config:?}, lower_bound={lower_bound}, upper_bound={upper_bound}"
873 );
874 }
875 
876 #[test]
877 fn test_queue_state_merge() {
878 let now = Instant::now();
879 
880 let mut state = QueueState {
881 midrid: MidRid(Mid::from("001"), None),
882 unpaced: false,
883 use_for_padding: true,
884 snapshot: QueueSnapshot {
885 created_at: now,
886 byte_size: 10_usize,
887 packet_count: 1332,
888 total_queue_time_origin: duration_ms(1_000),
889 last_emitted: Some(now + duration_ms(500)),
890 first_unsent: None,
891 priority: QueuePriority::Media,
892 },
893 };
894 
895 let other = QueueState {
896 midrid: MidRid(Mid::from("002"), None),
897 unpaced: false,
898 use_for_padding: false,
899 snapshot: QueueSnapshot {
900 created_at: now,
901 byte_size: 30_usize,
902 packet_count: 5,
903 total_queue_time_origin: duration_ms(337),
904 last_emitted: None,
905 first_unsent: Some(now + duration_ms(19)),
906 priority: QueuePriority::Padding,
907 },
908 };
909 
910 state.snapshot.merge(&other.snapshot);
911 
912 assert_eq!(state.midrid.mid(), Mid::from("001"));
913 assert_eq!(state.snapshot.byte_size, 40_usize);
914 assert_eq!(state.snapshot.packet_count, 1337);
915 assert_eq!(state.snapshot.total_queue_time_origin, duration_ms(1337));
916 
917 assert_eq!(state.snapshot.last_emitted, Some(now + duration_ms(500)));
918 assert_eq!(state.snapshot.first_unsent, Some(now + duration_ms(19)));
919 assert_eq!(state.snapshot.priority, QueuePriority::Media);
920 }
921 
922 #[test]
923 fn test_priority_ordering() {
924 assert!(QueuePriority::Media < QueuePriority::Padding);
925 assert!(QueuePriority::Media < QueuePriority::Empty);
926 assert!(QueuePriority::Padding < QueuePriority::Empty);
927 }
928 
929 fn assert_poll_success<F>(
930 pacer: &mut impl Pacer,
931 queue: &mut Queue,
932 now: Instant,
933 msg: &str,
934 do_asserts: F,
935 ) -> Instant
936 where
937 F: Fn(QueuedPacket),
938 {
939 let (qid, _cluster_id) = pacer.poll_queue().expect(msg);
940 let packet = queue.next_packet().unwrap();
941 let packet_size = packet.size();
942 do_asserts(packet);
943 pacer.register_send(now, DataSize::from(packet_size), qid);
944 queue.register_send(qid, now);
945 
946 let timeout = pacer.poll_timeout().0;
947 // After gating, the pacer requests a timeout at now + 1ยตs to ensure time advances
948 const MINIMAL_DELTA: Duration = Duration::from_micros(1);
949 assert!(
950 timeout <= Some(now + MINIMAL_DELTA) && timeout.is_some(),
951 "After a successful send the pacer should return an immediate timeout"
952 );
953 
954 // Simulate an immediate call to handle_timeout
955 handle_timeout_noisy(pacer, queue, now);
956 
957 timeout.unwrap()
958 }
959 
960 fn enqueue_packet_noisy(
961 pacer: &mut impl Pacer,
962 queue: &mut Queue,
963 seq_no: u16,
964 size: usize,
965 kind: PacketKind,
966 now: Instant,
967 ) {
968 let (header, payload_len, kind) = make_packet(seq_no, size, kind);
969 
970 let queued_packet = QueuedPacket {
971 queued_at: now,
972 header,
973 payload_len,
974 kind,
975 };
976 queue.enqueue_packet(queued_packet);
977 
978 // Matches the queueing behavior when the pacer is used in real code.
979 // Each packet being queued causes time to move forward in the pacer and the queue.
980 handle_timeout_noisy(pacer, queue, now);
981 }
982 
983 fn handle_timeout_noisy(pacer: &mut impl Pacer, queue: &mut Queue, now: Instant) {
984 queue.update_average_queue_time(now);
985 if let Some(padding_request) = pacer.handle_timeout(now, queue.queue_state(now)) {
986 queue.generate_padding(padding_request.padding, now);
987 
988 let timeout = pacer.poll_timeout().0;
989 if timeout.map(|t| t <= now).unwrap_or(false) {
990 // Refresh queue state
991 pacer.handle_timeout(now, queue.queue_state(now));
992 }
993 }
994 }
995 
996 fn duration_ms(ms: u64) -> Duration {
997 Duration::from_millis(ms)
998 }
999 
1000 fn make_packet(seq_no: u16, size: usize, kind: PacketKind) -> (RtpHeader, usize, PacketKind) {
1001 let header = RtpHeader {
1002 sequence_number: seq_no,
1003 ..Default::default()
1004 };
1005 
1006 (header, size, kind)
1007 }
1008 
1009 #[derive(Debug, Clone, Copy)]
1010 struct RealisticTestConfig {
1011 media_rate: Bitrate,
1012 padding_rate: Bitrate,
1013 duration: Duration,
1014 // Spike probability as a percentage
1015 spike_probability: u8,
1016 max_overshoot_factor: f32,
1017 frame_pacing: Duration,
1018 }
1019 
1020 impl Default for RealisticTestConfig {
1021 fn default() -> Self {
1022 RealisticTestConfig {
1023 media_rate: Bitrate::kbps(250),
1024 padding_rate: Bitrate::kbps(800),
1025 duration: Duration::from_secs(10),
1026 spike_probability: 0,
1027 max_overshoot_factor: 0.25,
1028 frame_pacing: Duration::from_millis(33), // ~30 FPS
1029 }
1030 }
1031 }
1032 
1033 /// Run a realistic test of the pacer with simulated media.
1034 ///
1035 /// Returns the media rate, padding, rate, and total rate achieved by the test.
1036 fn run_realistic_test(config: RealisticTestConfig) -> (Bitrate, Bitrate, Bitrate) {
1037 let RealisticTestConfig {
1038 media_rate,
1039 padding_rate,
1040 duration,
1041 spike_probability,
1042 max_overshoot_factor,
1043 frame_pacing,
1044 } = config;
1045 
1046 let base = Instant::now();
1047 let mut queue = Queue::default();
1048 let mut pacer = LeakyBucketPacer::new(media_rate);
1049 pacer.set_pacing_rate(padding_rate);
1050 pacer.set_padding_rate(padding_rate);
1051 
1052 let mut last_media_at = base - frame_pacing - Duration::from_millis(1);
1053 let mut media_sent = DataSize::ZERO;
1054 let mut padding_sent = DataSize::ZERO;
1055 let mut elapsed = Duration::ZERO;
1056 
1057 let generate_padding = |queue: &mut Queue, now: Instant, request: PaddingRequest| {
1058 let rand: f32 = fastrand::f32();
1059 let overshoot_factor: f32 = rand * max_overshoot_factor;
1060 let final_size = ((request.padding as f32) * (1.0 + overshoot_factor).round()) as usize;
1061 queue.generate_padding(final_size, now);
1062 };
1063 
1064 loop {
1065 if elapsed > duration {
1066 break;
1067 }
1068 
1069 let timeout = {
1070 if let Some((midrid, _cluster_id)) = pacer.poll_queue() {
1071 let packet = queue
1072 .next_packet()
1073 .unwrap_or_else(|| panic!("Should have a packet for {:?}", midrid));
1074 queue.register_send(midrid, base + elapsed);
1075 queue.update_average_queue_time(base + elapsed);
1076 pacer.register_send(
1077 base + elapsed,
1078 DataSize::bytes(packet.payload_len as i64),
1079 midrid,
1080 );
1081 if packet.kind == PacketKind::Padding {
1082 padding_sent += packet.payload_len.into();
1083 } else {
1084 media_sent += packet.payload_len.into();
1085 }
1086 continue;
1087 }
1088 
1089 pacer.poll_timeout()
1090 };
1091 
1092 let sleep_until_poll = timeout
1093 .0
1094 .map(|t| t.duration_since(base + elapsed))
1095 .unwrap_or(Duration::ZERO);
1096 
1097 let sleep_until_media =
1098 frame_pacing.saturating_sub((base + elapsed).duration_since(last_media_at));
1099 
1100 if sleep_until_poll < sleep_until_media {
1101 elapsed += sleep_until_poll;
1102 
1103 queue.update_average_queue_time(base + elapsed);
1104 if let Some(padding_request) =
1105 pacer.handle_timeout(base + elapsed, queue.queue_state(base + elapsed))
1106 {
1107 generate_padding(&mut queue, base + elapsed, padding_request);
1108 }
1109 continue;
1110 } else {
1111 elapsed += sleep_until_media;
1112 }
1113 
1114 let large_overshoot = (fastrand::u8(..) % 100) >= (100 - spike_probability);
1115 let mut to_add = if large_overshoot {
1116 (media_rate * 2.5) * frame_pacing
1117 } else {
1118 media_rate * frame_pacing
1119 };
1120 
1121 while to_add > DataSize::ZERO {
1122 let packet_size = to_add.min(DataSize::bytes(1100));
1123 let (header, size, kind) =
1124 make_packet(0, packet_size.as_bytes_usize(), PacketKind::Video);
1125 queue.enqueue_packet(QueuedPacket {
1126 queued_at: base + elapsed,
1127 header,
1128 payload_len: size,
1129 kind,
1130 });
1131 to_add -= packet_size;
1132 }
1133 last_media_at = base + elapsed;
1134 }
1135 
1136 let observed_media_rate = media_sent / duration;
1137 let observed_padding_rate = padding_sent / duration;
1138 let total_rate = (media_sent + padding_sent) / duration;
1139 
1140 (observed_media_rate, observed_padding_rate, total_rate)
1141 }
1142 
1143 /// A packet queue for use in tests of the pacer.
1144 mod queue {
1145 use std::collections::VecDeque;
1146 use std::time::{Duration, Instant};
1147 
1148 use crate::rtp_::{DataSize, RtpHeader};
1149 
1150 use super::*;
1151 
1152 // A packet queue
1153 pub(super) struct Queue {
1154 /// Queue for audio packets
1155 audio_queue: Inner,
1156 /// Queue for video packets
1157 video_queue: Inner,
1158 /// Queue for padding packets
1159 padding_queue: Inner,
1160 }
1161 
1162 pub(super) struct QueuedPacket {
1163 pub(super) queued_at: Instant,
1164 pub(super) header: RtpHeader,
1165 pub(super) payload_len: usize,
1166 pub(super) kind: PacketKind,
1167 }
1168 
1169 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1170 pub(super) enum PacketKind {
1171 Audio,
1172 Video,
1173 Padding,
1174 }
1175 
1176 impl Queue {
1177 pub(super) fn is_empty(&self) -> bool {
1178 self.audio_queue.is_empty() && self.video_queue.is_empty()
1179 }
1180 
1181 pub(super) fn update_average_queue_time(&mut self, now: Instant) {
1182 self.audio_queue.update_average_queue_time(now);
1183 self.video_queue.update_average_queue_time(now);
1184 }
1185 
1186 pub(super) fn enqueue_packet(&mut self, packet: QueuedPacket) {
1187 let queue = self.queue_for_kind_mut(packet.kind);
1188 queue.enqueue(packet);
1189 }
1190 
1191 pub(super) fn next_packet(&mut self) -> Option<QueuedPacket> {
1192 if !self.audio_queue.is_empty() {
1193 self.audio_queue.pop_packet()
1194 } else if !self.video_queue.is_empty() {
1195 self.video_queue.pop_packet()
1196 } else {
1197 self.padding_queue.pop_packet()
1198 }
1199 }
1200 
1201 pub(super) fn queue_state(&self, now: Instant) -> impl Iterator<Item = QueueState> {
1202 vec![
1203 self.audio_queue.queue_state(now),
1204 self.video_queue.queue_state(now),
1205 self.padding_queue.queue_state(now),
1206 ]
1207 .into_iter()
1208 }
1209 
1210 pub(super) fn register_send(&mut self, midrid: MidRid, now: Instant) {
1211 if self.video_queue.midrid == midrid {
1212 self.video_queue.last_emitted = Some(now);
1213 } else if self.audio_queue.midrid == midrid {
1214 self.audio_queue.last_emitted = Some(now);
1215 } else if self.padding_queue.midrid == midrid {
1216 self.padding_queue.last_emitted = Some(now);
1217 } else {
1218 panic!(
1219 "Attempted to register send on unknown queue with id {:?}",
1220 midrid
1221 );
1222 }
1223 }
1224 
1225 pub(super) fn generate_padding(&mut self, mut pad_size: usize, now: Instant) {
1226 while pad_size > 0 {
1227 let final_packet_size = pad_size.min(1200);
1228 let final_packet_size = DataSize::bytes(final_packet_size as i64);
1229 let (header, payload_len, kind) =
1230 make_packet(0, final_packet_size.as_bytes_usize(), PacketKind::Padding);
1231 self.enqueue_packet(QueuedPacket {
1232 queued_at: now,
1233 header,
1234 payload_len,
1235 kind,
1236 });
1237 self.update_average_queue_time(now);
1238 
1239 pad_size = pad_size.saturating_sub(final_packet_size.as_bytes_usize());
1240 }
1241 }
1242 
1243 fn queue_for_kind_mut(&mut self, kind: PacketKind) -> &mut Inner {
1244 match kind {
1245 PacketKind::Audio => &mut self.audio_queue,
1246 PacketKind::Video => &mut self.video_queue,
1247 PacketKind::Padding => &mut self.padding_queue,
1248 }
1249 }
1250 }
1251 
1252 impl Default for Queue {
1253 fn default() -> Self {
1254 Self {
1255 audio_queue: Inner::new(
1256 MidRid(Mid::from("001"), None),
1257 true,
1258 QueuePriority::Media,
1259 ),
1260 video_queue: Inner::new(
1261 MidRid(Mid::from("002"), None),
1262 false,
1263 QueuePriority::Media,
1264 ),
1265 padding_queue: Inner::new(
1266 MidRid(Mid::from("003"), None),
1267 false,
1268 QueuePriority::Padding,
1269 ),
1270 }
1271 }
1272 }
1273 
1274 impl QueuedPacket {
1275 pub(super) fn size(&self) -> usize {
1276 self.payload_len
1277 }
1278 }
1279 
1280 struct Inner {
1281 midrid: MidRid,
1282 last_emitted: Option<Instant>,
1283 queue: VecDeque<QueuedPacket>,
1284 packet_count: u32,
1285 total_time_spent_queued: Duration,
1286 last_update: Option<Instant>,
1287 is_audio: bool,
1288 priority: QueuePriority,
1289 }
1290 
1291 impl Inner {
1292 fn new(midrid: MidRid, is_audio: bool, priority: QueuePriority) -> Self {
1293 Self {
1294 midrid,
1295 last_emitted: None,
1296 queue: VecDeque::default(),
1297 packet_count: 0,
1298 total_time_spent_queued: Duration::ZERO,
1299 last_update: None,
1300 is_audio,
1301 priority,
1302 }
1303 }
1304 
1305 fn enqueue(&mut self, packet: QueuedPacket) {
1306 self.queue.push_back(packet);
1307 self.packet_count += 1;
1308 }
1309 
1310 fn pop_packet(&mut self) -> Option<QueuedPacket> {
1311 let packet = self.queue.pop_front()?;
1312 
1313 let time_spent_queued = self
1314 .last_update
1315 .map(|last_update| last_update - packet.queued_at)
1316 .unwrap_or(Duration::ZERO);
1317 self.total_time_spent_queued = self
1318 .total_time_spent_queued
1319 .saturating_sub(time_spent_queued);
1320 self.packet_count -= 1;
1321 
1322 Some(packet)
1323 }
1324 
1325 fn is_empty(&self) -> bool {
1326 self.queue.is_empty()
1327 }
1328 
1329 fn update_average_queue_time(&mut self, now: Instant) {
1330 let Some(last_update) = self.last_update else {
1331 self.last_update = Some(now);
1332 return;
1333 };
1334 
1335 let elapsed = now - last_update;
1336 self.total_time_spent_queued += elapsed * self.packet_count;
1337 self.last_update = Some(now);
1338 }
1339 
1340 fn queue_state(&self, now: Instant) -> QueueState {
1341 QueueState {
1342 midrid: self.midrid,
1343 unpaced: self.is_audio,
1344 use_for_padding: !self.is_audio && self.last_emitted.is_some(),
1345 snapshot: QueueSnapshot {
1346 created_at: now,
1347 byte_size: self.queue.iter().map(QueuedPacket::size).sum(),
1348 packet_count: self.packet_count,
1349 total_queue_time_origin: self.total_time_spent_queued,
1350 last_emitted: self.last_emitted,
1351 first_unsent: self.queue.iter().next().map(|p| p.queued_at),
1352 priority: self.priority,
1353 },
1354 }
1355 }
1356 }
1357 
1358 use std::fmt;
1359 
1360 impl fmt::Display for PacketKind {
1361 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1362 match self {
1363 PacketKind::Audio => write!(f, "audio"),
1364 PacketKind::Video => write!(f, "video"),
1365 PacketKind::Padding => write!(f, "padding"),
1366 }
1367 }
1368 }
1369 }
1370}