Skip to main content

moirai_async/timer/
limiter.rs

1use crate::timer::interval::Interval;
2use std::time::Duration;
3
4/// Rate limiter using token bucket algorithm
5pub struct RateLimiter {
6    permits: u32,
7    current_permits: u32,
8    interval: Interval,
9}
10
11/// A permit from the rate limiter
12pub struct RatePermit;
13
14impl RateLimiter {
15    /// Create a new rate limiter with specified permits per second
16    pub fn new(permits_per_second: u32) -> Self {
17        let interval_duration = if permits_per_second > 0 {
18            Duration::from_nanos(1_000_000_000 / permits_per_second as u64)
19        } else {
20            Duration::from_secs(1)
21        };
22
23        Self {
24            permits: permits_per_second,
25            current_permits: permits_per_second,
26            interval: Interval::new(interval_duration),
27        }
28    }
29
30    /// Acquire a permit to perform an operation
31    pub async fn acquire(&mut self) -> RatePermit {
32        if self.current_permits > 0 {
33            self.current_permits -= 1;
34            return RatePermit;
35        }
36
37        // Wait for next interval, then refill and consume one. `saturating_sub`
38        // guards the degenerate `permits_per_second == 0` construction: without
39        // it `self.permits - 1` underflows to `u32::MAX` (a zero-rate limiter
40        // that grants ~4.3 billion permits per interval in release builds, or
41        // panics under overflow-checks). A zero rate therefore refills to 0 and
42        // grants exactly the one permit this call is returning.
43        self.interval.next().await;
44        self.current_permits = self.permits.saturating_sub(1);
45        RatePermit
46    }
47
48    /// Try to acquire a permit without waiting
49    pub fn try_acquire(&mut self) -> Option<RatePermit> {
50        if self.current_permits > 0 {
51            self.current_permits -= 1;
52            Some(RatePermit)
53        } else {
54            None
55        }
56    }
57}
58
59#[cfg(test)]
60mod tests {
61    use super::*;
62
63    #[test]
64    fn try_acquire_exhausts_then_denies() {
65        let mut limiter = RateLimiter::new(3);
66        assert!(limiter.try_acquire().is_some());
67        assert!(limiter.try_acquire().is_some());
68        assert!(limiter.try_acquire().is_some());
69        assert!(
70            limiter.try_acquire().is_none(),
71            "a 3-permit limiter must deny the fourth immediate acquire"
72        );
73    }
74
75    #[test]
76    fn zero_rate_limiter_does_not_underflow_on_refill() {
77        // Regression: `new(0)` refilling via `permits - 1` underflowed to
78        // u32::MAX. The refill must saturate at 0 so the bucket never grants a
79        // spurious ~4.3 billion permits.
80        let mut limiter = RateLimiter::new(0);
81        assert!(
82            limiter.try_acquire().is_none(),
83            "a zero-rate limiter starts with no immediate permits"
84        );
85        // Exercise the refill arithmetic directly (the async path performs the
86        // same `permits.saturating_sub(1)`): it must not panic or wrap.
87        limiter.current_permits = limiter.permits.saturating_sub(1);
88        assert_eq!(limiter.current_permits, 0);
89    }
90}