use nmbrs_rate::{RateLimiter, RateSpec};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
#[tokio::test]
async fn rate_accuracy_100_ops_per_sec() {
let limiter = RateLimiter::start(RateSpec::new(100.0));
tokio::time::sleep(Duration::from_millis(100)).await;
let start = Instant::now();
for _ in 0..50 {
limiter.acquire().await;
}
let elapsed = start.elapsed();
limiter.stop().await;
let ms = elapsed.as_millis();
assert!(
(200..=1500).contains(&ms),
"50 ops at 100/s took {ms}ms, expected ~500ms"
);
}
#[tokio::test]
async fn rate_accuracy_1000_ops_per_sec() {
let limiter = RateLimiter::start(RateSpec::new(1000.0));
tokio::time::sleep(Duration::from_millis(50)).await;
let start = Instant::now();
for _ in 0..200 {
limiter.acquire().await;
}
let elapsed = start.elapsed();
limiter.stop().await;
let ms = elapsed.as_millis();
assert!(
(80..=800).contains(&ms),
"200 ops at 1000/s took {ms}ms, expected ~200ms"
);
}
#[tokio::test]
async fn rate_accuracy_10000_ops_per_sec() {
let limiter = RateLimiter::start(RateSpec::new(10000.0));
tokio::time::sleep(Duration::from_millis(50)).await;
let start = Instant::now();
for _ in 0..1000 {
limiter.acquire().await;
}
let elapsed = start.elapsed();
limiter.stop().await;
let ms = elapsed.as_millis();
assert!(
(30..=500).contains(&ms),
"1000 ops at 10000/s took {ms}ms, expected ~100ms"
);
}
#[tokio::test]
async fn concurrent_acquires() {
let limiter = Arc::new(RateLimiter::start(RateSpec::new(5000.0)));
tokio::time::sleep(Duration::from_millis(50)).await;
let counter = Arc::new(AtomicU64::new(0));
let mut handles = Vec::new();
for _ in 0..4 {
let limiter = limiter.clone();
let counter = counter.clone();
handles.push(tokio::spawn(async move {
for _ in 0..50 {
limiter.acquire().await;
counter.fetch_add(1, Ordering::Relaxed);
}
}));
}
for handle in handles {
handle.await.unwrap();
}
assert_eq!(counter.load(Ordering::Relaxed), 200);
}
#[tokio::test]
async fn many_tasks_contending() {
let limiter = Arc::new(RateLimiter::start(RateSpec::new(10000.0)));
tokio::time::sleep(Duration::from_millis(50)).await;
let counter = Arc::new(AtomicU64::new(0));
let mut handles = Vec::new();
for _ in 0..20 {
let limiter = limiter.clone();
let counter = counter.clone();
handles.push(tokio::spawn(async move {
for _ in 0..10 {
limiter.acquire().await;
counter.fetch_add(1, Ordering::Relaxed);
}
}));
}
for handle in handles {
handle.await.unwrap();
}
assert_eq!(counter.load(Ordering::Relaxed), 200);
}
#[tokio::test]
async fn very_low_rate() {
let spec = RateSpec::new(0.5);
assert_eq!(spec.unit, nmbrs_rate::TimeUnit::Micros);
let limiter = RateLimiter::start(spec);
tokio::time::sleep(Duration::from_millis(100)).await;
let start = Instant::now();
limiter.acquire().await;
let elapsed = start.elapsed();
limiter.stop().await;
assert!(
elapsed.as_millis() < 5000,
"took too long: {}ms",
elapsed.as_millis()
);
}
#[tokio::test]
async fn very_high_rate() {
let limiter = RateLimiter::start(RateSpec::new(1_000_000.0));
tokio::time::sleep(Duration::from_millis(100)).await;
let start = Instant::now();
for _ in 0..1000 {
limiter.acquire().await;
}
let elapsed = start.elapsed();
limiter.stop().await;
assert!(
elapsed.as_millis() < 2000,
"1000 ops at 1M/s took {}ms",
elapsed.as_millis()
);
}
#[tokio::test]
async fn burst_ratio_affects_recovery() {
let spec = RateSpec::with_burst(1000.0, 2.0); let limiter = RateLimiter::start(spec);
tokio::time::sleep(Duration::from_millis(200)).await;
let start = Instant::now();
for _ in 0..100 {
limiter.acquire().await;
}
let elapsed = start.elapsed();
limiter.stop().await;
assert!(
elapsed.as_millis() < 200,
"burst recovery should speed up, took {}ms",
elapsed.as_millis()
);
}
#[test]
fn spec_parse_edge_cases() {
assert!(RateSpec::parse("0.000001").is_ok());
assert!(RateSpec::parse("10000000").is_ok());
assert!(RateSpec::parse(" 1000 , 1.1 , start ").is_ok());
let s = RateSpec::parse("42").unwrap();
assert_eq!(s.ops_per_sec, 42.0);
assert_eq!(s.burst_ratio, 1.1);
assert_eq!(s.verb, nmbrs_rate::Verb::Start);
}
#[test]
fn spec_parse_rejects_bad_input() {
assert!(RateSpec::parse("").is_err());
assert!(RateSpec::parse("0").is_err());
assert!(RateSpec::parse("-100").is_err());
assert!(RateSpec::parse("abc").is_err());
assert!(RateSpec::parse("100,1.1,bogus").is_err());
}
#[tokio::test]
async fn acquire_order_is_fair() {
let limiter = Arc::new(RateLimiter::start(RateSpec::new(1000.0)));
tokio::time::sleep(Duration::from_millis(50)).await;
let counters: Vec<Arc<AtomicU64>> = (0..4).map(|_| Arc::new(AtomicU64::new(0))).collect();
let mut handles = Vec::new();
for counter in &counters {
let limiter = limiter.clone();
let counter = counter.clone();
handles.push(tokio::spawn(async move {
for _ in 0..25 {
limiter.acquire().await;
counter.fetch_add(1, Ordering::Relaxed);
}
}));
}
for handle in handles {
handle.await.unwrap();
}
for (i, c) in counters.iter().enumerate() {
assert_eq!(
c.load(Ordering::Relaxed),
25,
"task {i} didn't complete all acquires"
);
}
}
#[tokio::test]
async fn stop_is_clean() {
let limiter = RateLimiter::start(RateSpec::new(1000.0));
limiter.acquire().await;
limiter.stop().await;
}
#[tokio::test]
async fn drop_is_clean() {
{
let limiter = RateLimiter::start(RateSpec::new(1000.0));
limiter.acquire().await;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}