From 643f643ce9d37ecac9e1166ecaf89eea2deb5e40 Mon Sep 17 00:00:00 2001 From: dittops Date: Sun, 27 Sep 2026 18:22:57 +0530 Subject: [PATCH 1/2] test: a concurrent burst after idle gets the whole burst Real-Redis regression test for the refilled-burst path (token bucket 1/s, burst 2, local_allowance 0.8) across 0/2/4 ms store latency. Passes locally; written while chasing an intermittent under-admission of one request seen on a live gateway. Co-Authored-By: Claude Opus 5.5 --- tests/redis_cluster.rs | 64 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 64 insertions(+) diff --git a/tests/redis_cluster.rs b/tests/redis_cluster.rs index 81464e7..dcd0f58 100644 --- a/tests/redis_cluster.rs +++ b/tests/redis_cluster.rs @@ -255,3 +255,67 @@ 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 + }) + .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, 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 +} From 0d17796e012b5c4c69282bb2f6501fe143b06d8a Mon Sep 17 00:00:00 2001 From: dittops Date: Sun, 27 Sep 2026 22:48:53 +0530 Subject: [PATCH 2/2] test: a staggered burst after idle gets the whole burst with a second window Token bucket 1/s burst 2 plus a non-binding per-minute window, with the requests of each burst arriving 0-8 ms apart: 2.6 s after a drain, the refilled burst admits exactly 2 of 3, every round. Co-Authored-By: Claude Opus 5.5 --- tests/redis_cluster.rs | 60 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 60 insertions(+) diff --git a/tests/redis_cluster.rs b/tests/redis_cluster.rs index dcd0f58..f5dc00f 100644 --- a/tests/redis_cluster.rs +++ b/tests/redis_cluster.rs @@ -319,3 +319,63 @@ async fn burst(l: &Arc, n: usize) -> usize { } 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, 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 +}