Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
124 changes: 124 additions & 0 deletions tests/redis_cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -255,3 +255,127 @@ async fn concurrency_cap_is_cluster_wide() {
let g = reps[1].acquire("session").await;
assert!(matches!(g, Ok(Some(_))), "released slots come back");
}

// A burst after idle gets the whole burst, however the concurrent requests fall across syncs and
// whatever the Redis round-trip within `redis_timeout_ms`. Found live on a gateway: 3 concurrent
// requests to a token bucket (1/s, burst 2) after idle were admitted 1 or 2 at random.
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
async fn a_concurrent_burst_after_idle_gets_the_whole_burst() {
for delay_ms in [0u64, 2, 4] {
let svc = format!("burst{}", fastrand::u32(..));
let Some(reps) = cluster(1, &svc, |s| {
Arc::new(Slow {
inner: s,
delay: Duration::from_millis(delay_ms),
}) as Arc<dyn Store>
})
.await
else {
return;
};
let l = Arc::new(reps.into_iter().next().unwrap());
let cfg = RateLimitConfig {
algorithm: RateLimitAlgorithm::TokenBucket,
requests_per_second: Some(1),
burst_size: Some(2),
local_allowance: 0.8,
redis_timeout_ms: 10,
cache_ttl_ms: 500,
..Default::default()
};
l.set_policy("dep", Some(&cfg), None);
let mut got = Vec::new();
for _ in 0..6 {
// Drain the bucket, then leave it idle long enough to refill completely.
let drained = burst(&l, 5).await;
tokio::time::sleep(Duration::from_millis(2_600)).await;
got.push((drained, burst(&l, 3).await));
}
eprintln!("store delay {delay_ms} ms: (drain, after idle) admits {got:?}");
for (drained, after) in &got {
assert!(*drained <= 2, "delay {delay_ms}: {drained} > burst 2");
assert_eq!(
*after, 2,
"delay {delay_ms}: the refilled burst of 2 was not honoured: {got:?}"
);
}
l.shutdown().await;
}
}

/// `n` concurrent checks; how many were admitted.
async fn burst(l: &Arc<Limiter>, n: usize) -> usize {
let tasks: Vec<_> = (0..n)
.map(|_| {
let l = l.clone();
tokio::spawn(async move { l.check("dep").await })
})
.collect();
let mut admitted = 0;
for t in tasks {
if matches!(t.await.unwrap(), resil::Decision::Allow(_)) {
admitted += 1;
}
}
admitted
}

// A token bucket 1/s burst 2 PLUS a non-binding per-minute window (which, being a token bucket, gets
// burst 2 too), with requests arriving a few ms apart as they do through a load balancer: the
// second window must not cost the refilled burst a token.
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
async fn a_staggered_burst_after_idle_gets_the_whole_burst_with_a_second_window() {
let mut failures = Vec::new();
for rpm in [1_000u32, 3_000, 6_000] {
for gap_ms in [0u64, 1, 3, 8] {
let svc = format!("stag{}", fastrand::u32(..));
let Some(reps) = cluster(1, &svc, |s| s).await else {
return;
};
let l = Arc::new(reps.into_iter().next().unwrap());
let cfg = RateLimitConfig {
algorithm: RateLimitAlgorithm::TokenBucket,
requests_per_second: Some(1),
requests_per_minute: Some(rpm),
burst_size: Some(2),
local_allowance: 0.8,
redis_timeout_ms: 10,
cache_ttl_ms: 500,
..Default::default()
};
l.set_policy("dep", Some(&cfg), None);
for round in 0..4 {
let _ = staggered(&l, 5, gap_ms).await;
tokio::time::sleep(Duration::from_millis(2_600)).await;
let after = staggered(&l, 3, gap_ms).await;
if after != 2 {
failures.push((rpm, gap_ms, round, after));
}
}
l.shutdown().await;
}
}
assert!(
failures.is_empty(),
"(rpm, gap_ms, round, admitted of a refilled burst of 2): {failures:?}"
);
}

/// `n` checks started `gap_ms` apart, running concurrently; how many were admitted.
async fn staggered(l: &Arc<Limiter>, n: usize, gap_ms: u64) -> usize {
let mut tasks = Vec::new();
for _ in 0..n {
let l = l.clone();
tasks.push(tokio::spawn(async move { l.check("dep").await }));
if gap_ms > 0 {
tokio::time::sleep(Duration::from_millis(gap_ms)).await;
}
}
let mut admitted = 0;
for t in tasks {
if matches!(t.await.unwrap(), resil::Decision::Allow(_)) {
admitted += 1;
}
}
admitted
}
Loading