From 6bcffe4f48f54dcafde72505922096dc976f19df Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Philipp=20Kr=C3=BCger?= Date: Wed, 15 Apr 2026 14:08:54 +0200 Subject: [PATCH] Separate `Pacer::delay` into mutating `update` and non-mutating `delay` And write a proptest verifying it against the original implementation (`update_and_delay` after this commit). --- noq-proto/src/connection/pacing.rs | 133 ++++++++++++++++++++++++++--- noq-proto/src/connection/paths.rs | 2 +- 2 files changed, 120 insertions(+), 15 deletions(-) diff --git a/noq-proto/src/connection/pacing.rs b/noq-proto/src/connection/pacing.rs index 794b8070c..d9916cac6 100644 --- a/noq-proto/src/connection/pacing.rs +++ b/noq-proto/src/connection/pacing.rs @@ -12,7 +12,7 @@ use tracing::warn; /// The bucket refills at a rate slightly faster /// than one congestion window per RTT, as recommended in /// -#[derive(Debug)] +#[derive(Debug, PartialEq)] pub(super) struct Pacer { capacity: u64, last_window: u64, @@ -39,6 +39,70 @@ impl Pacer { self.tokens = self.tokens.saturating_sub(packet_length.into()) } + pub(super) fn delay( + &self, + smoothed_rtt: Duration, + bytes_to_send: u64, + window: u64, + ) -> Option { + // if we can already send a packet, there is no need for delay + if self.tokens >= bytes_to_send { + return None; + } + + // we disable pacing for extremely large windows + let window = u32::try_from(window).ok()?; + + let unscaled_delay = smoothed_rtt + .checked_mul((bytes_to_send.max(self.capacity) - self.tokens) as _) + .unwrap_or(Duration::MAX) + / window; + + // divisions come before multiplications to prevent overflow + // this is the time at which the pacing window becomes empty + Some((unscaled_delay / 5) * 4) + } + + pub(super) fn update(&mut self, smoothed_rtt: Duration, mtu: u16, window: u64, now: Instant) { + debug_assert_ne!( + window, 0, + "zero-sized congestion control window is nonsense" + ); + + if window != self.last_window || mtu != self.last_mtu { + self.capacity = optimal_capacity(smoothed_rtt, window, mtu); + + // Clamp the tokens + self.tokens = self.capacity.min(self.tokens); + self.last_window = window; + self.last_mtu = mtu; + } + + // we disable pacing for extremely large windows + let Ok(window) = u32::try_from(window) else { + return; + }; + + let time_elapsed = now.checked_duration_since(self.prev).unwrap_or_else(|| { + warn!("received a timestamp early than a previous recorded time, ignoring"); + Default::default() + }); + + if smoothed_rtt.as_nanos() == 0 { + return; + } + + let elapsed_rtts = time_elapsed.as_secs_f64() / smoothed_rtt.as_secs_f64(); + let new_tokens = (window as f64 * 1.25 * elapsed_rtts).round() as u64; + self.tokens = self.tokens.saturating_add(new_tokens).min(self.capacity); + + // In the unlikely event that we're getting polled faster than tokens are generated, ensure + // that `elapsed_rtts` can grow until we make progress. + if new_tokens > 0 { + self.prev = now; + } + } + /// Return how long we need to wait before sending `bytes_to_send`. /// /// If we can send a packet right away, this returns `None`. Otherwise, returns @@ -47,7 +111,7 @@ impl Pacer { /// /// The 5/4 ratio used here comes from the suggestion that N = 1.25 in the draft IETF /// RFC for QUIC. - pub(super) fn delay( + pub(super) fn update_and_delay( &mut self, smoothed_rtt: Duration, bytes_to_send: u64, @@ -166,6 +230,9 @@ const MAX_BURST_SIZE: u64 = 256; #[cfg(test)] mod tests { + use proptest::{prelude::Strategy, prop_assert_eq}; + use test_strategy::proptest; + use super::*; #[test] @@ -176,17 +243,17 @@ mod tests { assert!( Pacer::new(rtt, 30000, 1500, new_instant) - .delay(Duration::from_micros(0), 0, 1500, 1, old_instant) + .update_and_delay(Duration::from_micros(0), 0, 1500, 1, old_instant) .is_none() ); assert!( Pacer::new(rtt, 30000, 1500, new_instant) - .delay(Duration::from_micros(0), 1600, 1500, 1, old_instant) + .update_and_delay(Duration::from_micros(0), 1600, 1500, 1, old_instant) .is_none() ); assert!( Pacer::new(rtt, 30000, 1500, new_instant) - .delay(Duration::from_micros(0), 1500, 1500, 3000, old_instant) + .update_and_delay(Duration::from_micros(0), 1500, 1500, 3000, old_instant) .is_none() ); } @@ -229,27 +296,27 @@ mod tests { assert_eq!(pacer.tokens, pacer.capacity); let initial_tokens = pacer.tokens; - pacer.delay(rtt, mtu as u64, mtu, window * 2, now); + pacer.update_and_delay(rtt, mtu as u64, mtu, window * 2, now); assert_eq!( pacer.capacity, (2 * window as u128 * TARGET_BURST_INTERVAL.as_nanos() / rtt.as_nanos()) as u64 ); assert_eq!(pacer.tokens, initial_tokens); - pacer.delay(rtt, mtu as u64, mtu, window / 2, now); + pacer.update_and_delay(rtt, mtu as u64, mtu, window / 2, now); assert_eq!( pacer.capacity, (window as u128 / 2 * TARGET_BURST_INTERVAL.as_nanos() / rtt.as_nanos()) as u64 ); assert_eq!(pacer.tokens, initial_tokens / 2); - pacer.delay(rtt, mtu as u64, mtu * 2, window, now); + pacer.update_and_delay(rtt, mtu as u64, mtu * 2, window, now); assert_eq!( pacer.capacity, (window as u128 * TARGET_BURST_INTERVAL.as_nanos() / rtt.as_nanos()) as u64 ); - pacer.delay(rtt, mtu as u64, 20_000, window, now); + pacer.update_and_delay(rtt, mtu as u64, 20_000, window, now); assert_eq!(pacer.capacity, 20_000_u64 * MIN_BURST_SIZE); } @@ -265,7 +332,7 @@ mod tests { for _ in 0..packet_capacity { assert_eq!( - pacer.delay(rtt, mtu as u64, mtu, window, old_instant), + pacer.update_and_delay(rtt, mtu as u64, mtu, window, old_instant), None, "When capacity is available packets should be sent immediately" ); @@ -276,7 +343,7 @@ mod tests { let pace_duration = Duration::from_nanos((TARGET_BURST_INTERVAL.as_nanos() * 4 / 5) as u64); let actual_delay = pacer - .delay(rtt, mtu as u64, mtu, window, old_instant) + .update_and_delay(rtt, mtu as u64, mtu, window, old_instant) .expect("Send must be delayed"); let diff = actual_delay.abs_diff(pace_duration); @@ -288,7 +355,7 @@ mod tests { ); // Refill half of the tokens assert_eq!( - pacer.delay( + pacer.update_and_delay( rtt, mtu as u64, mtu, @@ -301,7 +368,7 @@ mod tests { for _ in 0..packet_capacity / 2 { assert_eq!( - pacer.delay(rtt, mtu as u64, mtu, window, old_instant), + pacer.update_and_delay(rtt, mtu as u64, mtu, window, old_instant), None, "When capacity is available packets should be sent immediately" ); @@ -311,7 +378,7 @@ mod tests { // Refill all capacity by waiting more than the expected duration assert_eq!( - pacer.delay( + pacer.update_and_delay( rtt, mtu as u64, mtu, @@ -322,4 +389,42 @@ mod tests { ); assert_eq!(pacer.tokens, pacer.capacity); } + + #[derive(Debug, Clone, Copy, test_strategy::Arbitrary)] + struct PacerParams { + smoothed_rtt: Duration, + window: u64, + mtu: u16, + } + + impl PacerParams { + fn into_pacer(self, now: Instant) -> Pacer { + Pacer::new(self.smoothed_rtt, self.window, self.mtu, now) + } + } + + #[proptest(cases = 256000)] + fn pacer_separate_and_combined_check_equal( + params: PacerParams, + transmitted: u16, + smoothed_rtt: Duration, + bytes_to_send: u64, + mtu: u16, + #[strategy(1u64..1_000_000_000)] window: u64, + #[strategy((0u64..1_000_000_000).prop_map(Duration::from_micros))] after: Duration, + ) { + let start = Instant::now(); + let mut pacer1 = params.into_pacer(start); + let mut pacer2 = params.into_pacer(start); + + pacer1.on_transmit(transmitted); + pacer2.on_transmit(transmitted); + + let now = start + after; + pacer1.update(smoothed_rtt, mtu, window, now); + let separate = pacer1.delay(smoothed_rtt, bytes_to_send, window); + let combined = pacer2.update_and_delay(smoothed_rtt, bytes_to_send, mtu, window, now); + prop_assert_eq!(separate, combined); + prop_assert_eq!(pacer1, pacer2); + } } diff --git a/noq-proto/src/connection/paths.rs b/noq-proto/src/connection/paths.rs index ceeefc006..ae99161b5 100644 --- a/noq-proto/src/connection/paths.rs +++ b/noq-proto/src/connection/paths.rs @@ -623,7 +623,7 @@ impl PathData { /// See [`Pacer::delay`]. pub(super) fn pacing_delay(&mut self, bytes_to_send: u64, now: Instant) -> Option { let smoothed_rtt = self.rtt.get(); - self.pacing.delay( + self.pacing.update_and_delay( smoothed_rtt, bytes_to_send, self.current_mtu(),