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;