mod token_bucket;
pub use token_bucket::TokenBucketRateLimiter as RateLimiter;
mod rt;
#[cfg(test)]
mod tests {
use std::time::{Duration, Instant};
use tokio::spawn;
use super::*;
#[tokio::test]
#[cfg(any(feature = "rt-tokio", feature = "rt-async-std"))]
async fn test_try_acquire() {
use rt::delay;
let rl = RateLimiter::new(3);
rl.burst(5);
rl.try_acquire().unwrap();
let duration = rl.try_acquire().unwrap_err();
assert!(duration > Duration::from_millis(330));
assert!(duration < Duration::from_millis(340));
delay(duration).await;
rl.try_acquire().unwrap();
}
#[tokio::test]
#[cfg(any(feature = "rt-tokio", feature = "rt-async-std"))]
async fn test_acquire() {
let rl = RateLimiter::new(3);
rl.burst(5);
let start = Instant::now();
rl.acquire().await;
assert!(start.elapsed() < Duration::from_millis(10));
rl.acquire().await;
assert!(start.elapsed() > Duration::from_millis(330));
assert!(start.elapsed() < Duration::from_millis(340));
rl.acquire().await;
assert!(start.elapsed() > Duration::from_millis(660));
assert!(start.elapsed() < Duration::from_millis(680));
let res = rl.acquire_with_timeout(Duration::from_millis(5000)).await;
assert!(res);
assert!(
start.elapsed() >= Duration::from_secs(1),
"got: {:?}",
start.elapsed()
);
assert!(start.elapsed() < Duration::from_millis(1030));
let res = rl.acquire_with_timeout(Duration::from_millis(10)).await;
assert!(!res);
assert!(start.elapsed() < Duration::from_millis(1050));
}
#[tokio::test]
#[cfg(any(feature = "rt-tokio", feature = "rt-async-std"))]
async fn test_clone() {
let rl = RateLimiter::new(3);
rl.burst(5);
let start = Instant::now();
rl.acquire().await;
assert!(start.elapsed() < Duration::from_millis(10));
rl.acquire().await;
assert!(start.elapsed() > Duration::from_millis(330));
assert!(start.elapsed() < Duration::from_millis(340));
rl.acquire().await;
assert!(start.elapsed() > Duration::from_millis(660));
assert!(start.elapsed() < Duration::from_millis(680));
let rl2 = rl.clone();
let jh = spawn(async move {
let rl = rl2;
let start = Instant::now();
let res = rl.acquire_with_timeout(Duration::from_millis(700)).await;
assert!(res);
assert!(
start.elapsed() <= Duration::from_millis(700),
"got: {:?}",
start.elapsed()
);
});
let res = rl.acquire_with_timeout(Duration::from_millis(700)).await;
assert!(res);
assert!(jh.await.is_ok());
}
#[tokio::test]
#[cfg(any(feature = "rt-tokio", feature = "rt-async-std"))]
async fn test_cancel_task() {
use rt::delay;
let rl = RateLimiter::new(3);
rl.burst(5);
let start = Instant::now();
rl.acquire().await;
assert!(start.elapsed() < Duration::from_millis(10));
rl.acquire().await;
assert!(start.elapsed() > Duration::from_millis(330));
assert!(start.elapsed() < Duration::from_millis(340));
rl.acquire().await;
assert!(start.elapsed() > Duration::from_millis(660));
assert!(start.elapsed() < Duration::from_millis(680));
let rl2 = rl.clone();
let jh = spawn(async move {
let rl = rl2;
let start = Instant::now();
let res = rl.acquire_with_timeout(Duration::from_millis(700)).await;
assert!(res);
assert!(start.elapsed() <= Duration::from_millis(700));
});
delay(Duration::from_millis(100)).await;
jh.abort();
let start = Instant::now();
let res = rl.acquire_with_timeout(Duration::from_millis(5000)).await;
assert!(res);
assert!(
start.elapsed() <= Duration::from_millis(700),
"got: {:?}",
start.elapsed()
);
let start = Instant::now();
let res = rl.acquire_with_timeout(Duration::from_millis(5000)).await;
assert!(res);
assert!(
start.elapsed() <= Duration::from_millis(10),
"got: {:?}",
start.elapsed()
);
}
}