use std::time::{Duration, Instant};
pub const MIN_BURST_BITS: f64 = 8.0 * 1500.0;
#[derive(Debug, Clone)]
pub struct Pacer {
bitrate: f64,
configured_burst_bits: Option<f64>,
budget_bits: f64,
last_refill: Option<Instant>,
}
impl Pacer {
pub fn new(bits_per_second: f64) -> Self {
Self {
bitrate: bits_per_second.max(0.0),
configured_burst_bits: None,
budget_bits: Self::burst_for(bits_per_second),
last_refill: None,
}
}
pub fn with_burst_bits(mut self, burst_bits: f64) -> Self {
self.configured_burst_bits = Some(burst_bits.max(MIN_BURST_BITS));
self.budget_bits = self.budget_bits.min(self.burst_bits());
self
}
fn burst_for(bits_per_second: f64) -> f64 {
(bits_per_second / 10.0).max(MIN_BURST_BITS)
}
pub fn target_bitrate(&self) -> f64 {
self.bitrate
}
pub fn set_target_bitrate(&mut self, bits_per_second: f64) {
self.bitrate = bits_per_second.max(0.0);
self.budget_bits = self.budget_bits.min(self.burst_bits());
}
pub fn refill(&mut self, now: Instant) {
if let Some(last) = self.last_refill {
let elapsed = now.saturating_duration_since(last).as_secs_f64();
self.budget_bits = (self.budget_bits + elapsed * self.bitrate).min(self.burst_bits());
}
self.last_refill = Some(now);
}
pub fn can_afford(&self, bits: f64) -> bool {
self.budget_bits >= bits
}
pub fn consume(&mut self, bits: f64) {
self.budget_bits -= bits;
}
pub fn time_until_affordable(&self, bits: f64) -> Option<Duration> {
if self.budget_bits >= bits {
return Some(Duration::ZERO);
}
if self.bitrate <= 0.0 {
return None;
}
Some(Duration::from_secs_f64(
(bits - self.budget_bits) / self.bitrate,
))
}
pub fn affordable_at(&self, bits: f64) -> Option<Instant> {
let last = self.last_refill?;
self.time_until_affordable(bits).map(|wait| last + wait)
}
pub fn budget_bits(&self) -> f64 {
self.budget_bits
}
pub fn can_release(&self, bits: f64) -> bool {
let burst_bits = self.burst_bits();
if bits > burst_bits {
return self.budget_bits >= burst_bits;
}
self.budget_bits >= bits
}
pub fn releasable_at(&self, bits: f64) -> Option<Instant> {
self.affordable_at(bits.min(self.burst_bits()))
}
pub fn burst_bits(&self) -> f64 {
self.configured_burst_bits
.unwrap_or_else(|| Self::burst_for(self.bitrate))
}
}
#[cfg(test)]
mod tests {
use super::*;
fn pacer() -> Pacer {
Pacer::new(1_000_000.0).with_burst_bits(MIN_BURST_BITS)
}
#[test]
fn a_new_bucket_starts_full() {
let pacer = pacer();
assert_eq!(MIN_BURST_BITS, pacer.budget_bits());
assert!(pacer.can_afford(MIN_BURST_BITS));
}
#[test]
fn spending_reduces_the_budget() {
let mut pacer = pacer();
pacer.consume(1000.0);
assert_eq!(MIN_BURST_BITS - 1000.0, pacer.budget_bits());
}
#[test]
fn the_budget_refills_from_elapsed_time() {
let now = Instant::now();
let mut pacer = pacer();
pacer.refill(now);
pacer.consume(pacer.budget_bits());
assert_eq!(0.0, pacer.budget_bits());
pacer.refill(now + Duration::from_millis(10));
assert!(
(pacer.budget_bits() - 10_000.0).abs() < 1.0,
"{}",
pacer.budget_bits()
);
}
#[test]
fn refilling_in_steps_matches_refilling_at_once() {
let now = Instant::now();
let mut stepped = pacer();
stepped.refill(now);
stepped.consume(stepped.budget_bits());
for step in 1..=10 {
stepped.refill(now + Duration::from_millis(step));
}
let mut at_once = pacer();
at_once.refill(now);
at_once.consume(at_once.budget_bits());
at_once.refill(now + Duration::from_millis(10));
assert!((stepped.budget_bits() - at_once.budget_bits()).abs() < 1.0);
}
#[test]
fn the_budget_is_capped_at_the_burst() {
let now = Instant::now();
let mut pacer = pacer();
pacer.refill(now);
pacer.refill(now + Duration::from_secs(60));
assert_eq!(MIN_BURST_BITS, pacer.budget_bits());
}
#[test]
fn the_time_until_affordable_is_zero_when_it_already_is() {
let pacer = pacer();
assert_eq!(Some(Duration::ZERO), pacer.time_until_affordable(100.0));
}
#[test]
fn the_time_until_affordable_scales_with_the_shortfall() {
let now = Instant::now();
let mut pacer = pacer();
pacer.refill(now);
pacer.consume(pacer.budget_bits());
let wait = pacer
.time_until_affordable(10_000.0)
.expect("a finite wait");
assert!((wait.as_secs_f64() - 0.010).abs() < 0.0005, "got {wait:?}");
assert_eq!(Some(now + wait), pacer.affordable_at(10_000.0));
}
#[test]
fn nothing_becomes_affordable_at_a_zero_rate() {
let now = Instant::now();
let mut pacer = Pacer::new(0.0);
pacer.refill(now);
pacer.consume(pacer.budget_bits());
assert_eq!(None, pacer.time_until_affordable(1000.0));
assert_eq!(None, pacer.affordable_at(1000.0));
}
#[test]
fn changing_the_rate_changes_how_fast_the_budget_refills() {
let now = Instant::now();
let mut pacer = Pacer::new(1_000_000.0);
pacer.refill(now);
pacer.consume(pacer.budget_bits());
pacer.set_target_bitrate(2_000_000.0);
pacer.refill(now + Duration::from_millis(10));
assert!(
(pacer.budget_bits() - 20_000.0).abs() < 1.0,
"{}",
pacer.budget_bits()
);
assert_eq!(2_000_000.0, pacer.target_bitrate());
}
#[test]
fn lowering_the_rate_clamps_the_budget_to_the_new_burst() {
let now = Instant::now();
let mut pacer = Pacer::new(100_000_000.0);
pacer.refill(now);
let before = pacer.budget_bits();
pacer.set_target_bitrate(1000.0);
assert!(pacer.budget_bits() < before);
assert_eq!(
MIN_BURST_BITS,
pacer.budget_bits(),
"clamped to the floor burst"
);
}
#[test]
fn a_negative_rate_is_treated_as_zero() {
let mut pacer = Pacer::new(-5.0);
assert_eq!(0.0, pacer.target_bitrate());
pacer.set_target_bitrate(-1.0);
assert_eq!(0.0, pacer.target_bitrate());
}
#[test]
fn a_configured_burst_survives_a_rate_change() {
let mut pacer = Pacer::new(1_000_000.0).with_burst_bits(MIN_BURST_BITS);
assert_eq!(MIN_BURST_BITS, pacer.burst_bits());
pacer.set_target_bitrate(100_000_000.0);
assert_eq!(
MIN_BURST_BITS,
pacer.burst_bits(),
"a rate change must not widen a burst the caller set"
);
assert_eq!(100_000_000.0, pacer.target_bitrate());
}
#[test]
fn a_derived_burst_follows_the_rate() {
let mut pacer = Pacer::new(1_000_000.0);
assert_eq!(100_000.0, pacer.burst_bits());
pacer.set_target_bitrate(2_000_000.0);
assert_eq!(200_000.0, pacer.burst_bits());
}
#[test]
fn an_oversized_packet_waits_for_a_full_budget() {
let now = Instant::now();
let mut pacer = pacer();
pacer.refill(now);
let oversized = MIN_BURST_BITS * 2.0;
assert!(pacer.can_release(oversized), "a full budget releases it");
pacer.consume(oversized);
assert!(
!pacer.can_release(oversized),
"the next one waits for the debt to be repaid"
);
let at = pacer.releasable_at(oversized).expect("a finite wait");
assert_eq!(
Some(Duration::from_secs_f64(oversized / 1_000_000.0)),
pacer.time_until_affordable(MIN_BURST_BITS)
);
pacer.refill(at);
assert!(pacer.can_release(oversized), "and then it goes");
}
#[test]
fn a_packet_larger_than_the_burst_can_still_be_sent() {
let now = Instant::now();
let mut pacer = pacer();
pacer.refill(now);
let oversized = MIN_BURST_BITS * 2.0;
assert!(!pacer.can_afford(oversized));
pacer.refill(now + Duration::from_secs(10));
assert!(!pacer.can_afford(oversized));
pacer.consume(oversized);
assert!(
pacer.budget_bits() < 0.0,
"the overshoot is paid back over time"
);
pacer.refill(now + Duration::from_secs(20));
assert_eq!(
MIN_BURST_BITS,
pacer.budget_bits(),
"and recovers to the burst"
);
}
}