From ca2ad9f1173b537d0fdb47bdc75294d674fa852b Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Sat, 19 Sep 2026 02:00:01 +0000 Subject: [PATCH] test: drive the setup limiter's clock in the refill test The drained-bucket refill test raced the wall clock: its draining loop competed with the refill it was draining against, and it then slept for real time before the legitimate setup. On a loaded runner either half could land on the wrong side of a token. SessionSetupRateLimiter now reads the time through a clock function, Instant::now in production, which a test can replace. TokenBucket gains with_params_at and try_acquire_at so the limiter passes its own reading through; every other bucket keeps calling Instant::now exactly as before. The test pauses tokio's clock, points the limiter at it, delivers exactly one setup past the burst and expects exactly one refusal, then advances the clock by the time a full burst takes to refill before establishing the legitimate session. --- src/node/rate_limit.rs | 56 ++++++++++++++++++++++++++++++--------- src/node/tests/session.rs | 53 ++++++++++++++++++------------------ 2 files changed, 69 insertions(+), 40 deletions(-) diff --git a/src/node/rate_limit.rs b/src/node/rate_limit.rs index b456c002..726a1d84 100644 --- a/src/node/rate_limit.rs +++ b/src/node/rate_limit.rs @@ -93,11 +93,19 @@ impl TokenBucket { /// * `capacity` - Maximum number of tokens (burst capacity) /// * `refill_rate` - Tokens added per second pub fn with_params(capacity: u32, refill_rate: f64) -> Self { + Self::with_params_at(capacity, refill_rate, Instant::now()) + } + + /// Create a token bucket with custom parameters, full as of `now`. + /// + /// For a caller that keeps its own clock and passes the same clock's + /// readings to [`Self::try_acquire_at`]. + pub fn with_params_at(capacity: u32, refill_rate: f64, now: Instant) -> Self { Self { capacity, tokens: capacity as f64, refill_rate, - last_refill: Instant::now(), + last_refill: now, } } @@ -114,7 +122,20 @@ impl TokenBucket { /// Returns `true` if n tokens were available and consumed, `false` if /// rate limited (insufficient tokens). pub fn try_acquire_n(&mut self, n: u32) -> bool { - self.refill(); + self.try_acquire_n_at(n, Instant::now()) + } + + /// Try to consume one token, refilling as of `now`. + /// + /// `now` must come from the same clock as every earlier reading this + /// bucket was given. + pub fn try_acquire_at(&mut self, now: Instant) -> bool { + self.try_acquire_n_at(1, now) + } + + /// Try to consume n tokens, refilling as of `now`. + fn try_acquire_n_at(&mut self, n: u32, now: Instant) -> bool { + self.refill_at(now); if self.tokens >= n as f64 { self.tokens -= n as f64; @@ -127,14 +148,14 @@ impl TokenBucket { /// Check if tokens are available without consuming them. #[cfg(test)] pub fn available(&mut self) -> bool { - self.refill(); + self.refill_at(Instant::now()); self.tokens >= 1.0 } /// Get the current number of available tokens. #[cfg(test)] pub fn tokens(&mut self) -> f64 { - self.refill(); + self.refill_at(Instant::now()); self.tokens } @@ -144,9 +165,8 @@ impl TokenBucket { self.capacity } - /// Refill tokens based on elapsed time. - fn refill(&mut self) { - let now = Instant::now(); + /// Refill tokens based on the time elapsed up to `now`. + fn refill_at(&mut self, now: Instant) { let elapsed = now.duration_since(self.last_refill); let elapsed_secs = elapsed.as_secs_f64(); @@ -174,7 +194,7 @@ impl TokenBucket { /// estimated time until one token will be available. #[cfg(test)] pub fn time_until_available(&mut self) -> std::time::Duration { - self.refill(); + self.refill_at(Instant::now()); if self.tokens >= 1.0 { std::time::Duration::ZERO @@ -458,6 +478,8 @@ pub struct SessionSetupRateLimiter { stranger: (u32, f64), /// Burst and refill rate for a new link's established bucket. established: (u32, f64), + /// Where the limiter reads the time. `Instant::now` outside tests. + clock: fn() -> Instant, } impl SessionSetupRateLimiter { @@ -469,6 +491,7 @@ impl SessionSetupRateLimiter { buckets: HashMap::new(), stranger, established, + clock: Instant::now, } } @@ -477,22 +500,22 @@ impl SessionSetupRateLimiter { /// Returns `false` when the class's bucket for that link is empty, in /// which case the caller must drop the message before doing any work. pub fn try_admit(&mut self, link_peer: &NodeAddr, class: Msg1Class) -> bool { - let now = Instant::now(); + let now = (self.clock)(); let stranger = self.stranger; let established = self.established; let link = self .buckets .entry(*link_peer) .or_insert_with(|| LinkBuckets { - stranger: TokenBucket::with_params(stranger.0, stranger.1), - established: TokenBucket::with_params(established.0, established.1), + stranger: TokenBucket::with_params_at(stranger.0, stranger.1, now), + established: TokenBucket::with_params_at(established.0, established.1, now), seen: now, }); link.seen = now; let admitted = match class { - Msg1Class::Stranger => link.stranger.try_acquire(), - Msg1Class::EstablishedLink => link.established.try_acquire(), + Msg1Class::Stranger => link.stranger.try_acquire_at(now), + Msg1Class::EstablishedLink => link.established.try_acquire_at(now), }; if admitted { @@ -502,6 +525,13 @@ impl SessionSetupRateLimiter { admitted } + /// Read the time from `clock` instead of `Instant::now`, so a test can + /// drive the refill. + #[cfg(test)] + pub fn set_clock(&mut self, clock: fn() -> Instant) { + self.clock = clock; + } + /// Number of link peers currently holding buckets. #[cfg(test)] pub fn len(&self) -> usize { diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index c01ebac0..78524156 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -4660,40 +4660,39 @@ async fn test_forged_setups_from_one_link_peer_stop_creating_session_entries_onc cleanup_nodes(&mut nodes).await; } +/// The limiter's clock for the refill test: tokio's paused clock, which the +/// test moves with `tokio::time::advance` and nothing else moves. +fn paused_now() -> std::time::Instant { + tokio::time::Instant::now().into_std() +} + #[tokio::test] async fn test_a_drained_setup_bucket_refills_and_admits_the_next_legitimate_setup() { - // The refill has to be slow enough that the draining loop below cannot be - // outrun by the refill it is draining against. At the 50/s this test used - // to run at, a token returned every 20 ms, so on a loaded runner the loop - // outlived its own window, the third setup was admitted, and the - // precondition failed on arrangement rather than on behaviour. At 2/s a - // delivery would have to take 500 ms to lose that race. - let mut nodes = make_setup_limited_pair(2, 2.0).await; + const BURST: u32 = 2; + const RATE: f64 = 2.0; + let mut nodes = make_setup_limited_pair(BURST, RATE).await; + + // From here the limiter reads a clock only the test moves, so the drain + // cannot race a refill however slowly each delivery runs, and the refill + // below is exactly the one the test grants. + tokio::time::pause(); + nodes[1].node.setup_rate_limiter.set_clock(paused_now); - // Deliver until one is actually refused, rather than assuming three is - // enough: a delivery the refill absorbs costs one more iteration and - // nothing else. The cap is what a runner slow enough to lose even this - // race trips, and it says so rather than reporting a drained bucket that - // was never drained. let before = nodes[1].node.stats().session.setup_rate_limited; - let mut delivered = 0; - while nodes[1].node.stats().session.setup_rate_limited == before { - assert!( - delivered < 50, - "the bucket must actually be drained before the refill is tested; \ - 50 forged setups drew no refusal, so each delivery is outlasting \ - the 500 ms refill interval" - ); + for _ in 0..=BURST { deliver_forged_setup_over_link(&mut nodes).await; - delivered += 1; } + assert_eq!( + nodes[1].node.stats().session.setup_rate_limited, + before + 1, + "the burst must be admitted and the one setup past it refused" + ); - // A full burst back from empty at 2/s, so the legitimate setup below meets - // the same bucket however many tokens the drain left behind. The point - // being made is that the denial is transient and clears on its own; the - // length of the window is a function of the configured rate, not of the - // claim. - tokio::time::sleep(Duration::from_millis(1200)).await; + // A full burst back from empty at the configured rate, so the legitimate + // setup below meets a bucket the refill alone has restored. The denial + // is transient and clears on its own; how long it lasts is a function of + // the configured rate. + tokio::time::advance(Duration::from_secs_f64(f64::from(BURST) / RATE)).await; establish_pair_session(&mut nodes).await; cleanup_nodes(&mut nodes).await;