Skip to main content

rust_zero_core/
limit.rs

1use std::{
2    collections::HashMap,
3    hash::Hash,
4    sync::{Arc, Mutex},
5    time::{Duration, Instant, SystemTime, UNIX_EPOCH},
6};
7
8#[cfg(feature = "stores-redis")]
9use std::sync::atomic::{AtomicU64, AtomicU8, Ordering};
10
11#[cfg(feature = "stores-redis")]
12use sha2::{Digest, Sha256};
13
14#[cfg(feature = "stores-redis")]
15use crate::{RedisStore, RedisStoreError};
16
17#[cfg(feature = "stores-redis")]
18const TOKEN_BUCKET_SCRIPT: &str = r#"
19local clock = redis.call('TIME')
20local now = (clock[1] * 1000) + math.floor(clock[2] / 1000)
21local rate = tonumber(ARGV[1])
22local burst = tonumber(ARGV[2])
23local permits = tonumber(ARGV[3])
24local values = redis.call('HMGET', KEYS[1], 'tokens', 'updated')
25local tokens = tonumber(values[1]) or burst
26local updated = tonumber(values[2]) or now
27if now > updated then
28  tokens = math.min(burst, tokens + ((now - updated) * rate / 1000))
29end
30local retry = 0
31local outcome = 0
32if tokens >= permits then
33  tokens = tokens - permits
34  if tokens < 1 then outcome = 1 end
35else
36  outcome = 2
37  retry = math.ceil((permits - tokens) * 1000 / rate)
38end
39redis.call('HSET', KEYS[1], 'tokens', tokens, 'updated', now)
40redis.call('PEXPIRE', KEYS[1], math.max(1000, math.ceil(burst * 2000 / rate)))
41return {outcome, retry}
42"#;
43
44#[cfg(feature = "stores-redis")]
45const PERIOD_SCRIPT: &str = r#"
46local clock = redis.call('TIME')
47local now = (clock[1] * 1000) + math.floor(clock[2] / 1000)
48local period = tonumber(ARGV[1])
49local quota = tonumber(ARGV[2])
50local boundary = now - (now % period) + period
51local used = tonumber(redis.call('GET', KEYS[1])) or 0
52if used >= quota then return {2, boundary - now} end
53used = redis.call('INCR', KEYS[1])
54if used == 1 then redis.call('PEXPIREAT', KEYS[1], boundary) end
55if used == quota then return {1, 0} end
56return {0, 0}
57"#;
58
59/// Result of requesting capacity from a rate limiter.
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub enum LimitDecision {
62    Allowed,
63    HitQuota,
64    OverQuota { retry_after: Duration },
65}
66
67impl LimitDecision {
68    pub fn is_allowed(self) -> bool {
69        !matches!(self, Self::OverQuota { .. })
70    }
71}
72
73struct Bucket {
74    available: f64,
75    last_refill: Instant,
76}
77
78/// A process-local token-bucket limiter with burst support.
79#[derive(Clone)]
80pub struct TokenLimiter {
81    rate_per_second: f64,
82    burst: u32,
83    bucket: Arc<Mutex<Bucket>>,
84}
85
86impl TokenLimiter {
87    pub fn new(rate_per_second: u32, burst: u32) -> Self {
88        assert!(rate_per_second > 0, "rate must be greater than zero");
89        assert!(burst > 0, "burst must be greater than zero");
90        Self {
91            rate_per_second: f64::from(rate_per_second),
92            burst,
93            bucket: Arc::new(Mutex::new(Bucket {
94                available: f64::from(burst),
95                last_refill: Instant::now(),
96            })),
97        }
98    }
99
100    pub fn allow(&self) -> bool {
101        self.take(1).is_allowed()
102    }
103
104    pub fn take(&self, permits: u32) -> LimitDecision {
105        assert!(permits > 0, "permit count must be greater than zero");
106        if permits > self.burst {
107            return LimitDecision::OverQuota {
108                retry_after: Duration::MAX,
109            };
110        }
111
112        let mut bucket = self.bucket.lock().expect("token limiter lock poisoned");
113        let now = Instant::now();
114        let elapsed = now.duration_since(bucket.last_refill).as_secs_f64();
115        bucket.available =
116            (bucket.available + elapsed * self.rate_per_second).min(f64::from(self.burst));
117        bucket.last_refill = now;
118
119        if bucket.available >= f64::from(permits) {
120            bucket.available -= f64::from(permits);
121            if bucket.available < 1.0 {
122                LimitDecision::HitQuota
123            } else {
124                LimitDecision::Allowed
125            }
126        } else {
127            LimitDecision::OverQuota {
128                retry_after: Duration::from_secs_f64(
129                    (f64::from(permits) - bucket.available) / self.rate_per_second,
130                ),
131            }
132        }
133    }
134}
135
136struct PeriodWindow {
137    used: u32,
138    window: u128,
139}
140
141/// A keyed fixed-window quota limiter.
142///
143/// A distributed backend can wrap this API, while this implementation provides go-zero's
144/// in-process rescue behavior when no Redis deployment is available.
145pub struct PeriodLimiter<K> {
146    period: Duration,
147    quota: u32,
148    windows: Mutex<HashMap<K, PeriodWindow>>,
149}
150
151impl<K> PeriodLimiter<K>
152where
153    K: Eq + Hash,
154{
155    pub fn new(period: Duration, quota: u32) -> Self {
156        assert!(!period.is_zero(), "period must be greater than zero");
157        assert!(quota > 0, "quota must be greater than zero");
158        Self {
159            period,
160            quota,
161            windows: Mutex::new(HashMap::new()),
162        }
163    }
164
165    pub fn take(&self, key: K) -> LimitDecision {
166        let now = unix_millis();
167        let period = self.period.as_millis();
168        let current_window = now / period;
169        let retry_after = duration_from_millis_saturating(
170            (current_window + 1)
171                .saturating_mul(period)
172                .saturating_sub(now),
173        );
174        let mut windows = self.windows.lock().expect("period limiter lock poisoned");
175        let window = windows.entry(key).or_insert(PeriodWindow {
176            used: 0,
177            window: current_window,
178        });
179        if window.window != current_window {
180            window.used = 0;
181            window.window = current_window;
182        }
183
184        if window.used >= self.quota {
185            return LimitDecision::OverQuota { retry_after };
186        }
187
188        window.used += 1;
189        if window.used == self.quota {
190            LimitDecision::HitQuota
191        } else {
192            LimitDecision::Allowed
193        }
194    }
195
196    pub fn clear(&self, key: &K) {
197        self.windows
198            .lock()
199            .expect("period limiter lock poisoned")
200            .remove(key);
201    }
202}
203
204fn unix_millis() -> u128 {
205    SystemTime::now()
206        .duration_since(UNIX_EPOCH)
207        .unwrap_or_default()
208        .as_millis()
209}
210
211fn duration_from_millis_saturating(millis: u128) -> Duration {
212    Duration::from_millis(u64::try_from(millis).unwrap_or(u64::MAX))
213}
214
215/// Observable state for a Redis-backed limiter and its process-local rescue path.
216#[cfg(feature = "stores-redis")]
217#[derive(Debug, Clone, Copy, PartialEq, Eq)]
218pub struct RedisLimiterSnapshot {
219    pub backend_available: bool,
220    pub backend_failures: u64,
221    pub backend_recoveries: u64,
222    pub redis_allowed: u64,
223    pub redis_hit_quota: u64,
224    pub redis_over_quota: u64,
225    pub rescue_allowed: u64,
226    pub rescue_hit_quota: u64,
227    pub rescue_over_quota: u64,
228}
229
230#[cfg(feature = "stores-redis")]
231struct RedisLimiterMonitor {
232    // 0 = not observed, 1 = available, 2 = unavailable, 3 = recovery probe in flight.
233    backend_state: AtomicU8,
234    next_probe_at_ms: AtomicU64,
235    backend_failures: AtomicU64,
236    backend_recoveries: AtomicU64,
237    redis_allowed: AtomicU64,
238    redis_hit_quota: AtomicU64,
239    redis_over_quota: AtomicU64,
240    rescue_allowed: AtomicU64,
241    rescue_hit_quota: AtomicU64,
242    rescue_over_quota: AtomicU64,
243}
244
245#[cfg(feature = "stores-redis")]
246impl Default for RedisLimiterMonitor {
247    fn default() -> Self {
248        Self {
249            backend_state: AtomicU8::new(0),
250            next_probe_at_ms: AtomicU64::new(0),
251            backend_failures: AtomicU64::new(0),
252            backend_recoveries: AtomicU64::new(0),
253            redis_allowed: AtomicU64::new(0),
254            redis_hit_quota: AtomicU64::new(0),
255            redis_over_quota: AtomicU64::new(0),
256            rescue_allowed: AtomicU64::new(0),
257            rescue_hit_quota: AtomicU64::new(0),
258            rescue_over_quota: AtomicU64::new(0),
259        }
260    }
261}
262
263#[cfg(feature = "stores-redis")]
264impl RedisLimiterMonitor {
265    fn should_attempt_redis(&self) -> bool {
266        match self.backend_state.load(Ordering::Acquire) {
267            0 | 1 => true,
268            2 => {
269                let now = unix_millis_u64();
270                now >= self.next_probe_at_ms.load(Ordering::Acquire)
271                    && self
272                        .backend_state
273                        .compare_exchange(2, 3, Ordering::AcqRel, Ordering::Acquire)
274                        .is_ok()
275            }
276            _ => false,
277        }
278    }
279
280    fn redis_success(&self, decision: LimitDecision) {
281        if matches!(self.backend_state.swap(1, Ordering::AcqRel), 2 | 3) {
282            self.backend_recoveries.fetch_add(1, Ordering::Relaxed);
283        }
284        decision_counter(
285            decision,
286            &self.redis_allowed,
287            &self.redis_hit_quota,
288            &self.redis_over_quota,
289        );
290    }
291
292    fn redis_failure(&self, probe_interval_ms: u64) {
293        self.next_probe_at_ms.store(
294            unix_millis_u64().saturating_add(probe_interval_ms),
295            Ordering::Release,
296        );
297        self.backend_state.store(2, Ordering::Release);
298        self.backend_failures.fetch_add(1, Ordering::Relaxed);
299    }
300
301    fn rescue(&self, decision: LimitDecision) {
302        decision_counter(
303            decision,
304            &self.rescue_allowed,
305            &self.rescue_hit_quota,
306            &self.rescue_over_quota,
307        );
308    }
309
310    fn snapshot(&self) -> RedisLimiterSnapshot {
311        RedisLimiterSnapshot {
312            backend_available: self.backend_state.load(Ordering::Acquire) == 1,
313            backend_failures: self.backend_failures.load(Ordering::Relaxed),
314            backend_recoveries: self.backend_recoveries.load(Ordering::Relaxed),
315            redis_allowed: self.redis_allowed.load(Ordering::Relaxed),
316            redis_hit_quota: self.redis_hit_quota.load(Ordering::Relaxed),
317            redis_over_quota: self.redis_over_quota.load(Ordering::Relaxed),
318            rescue_allowed: self.rescue_allowed.load(Ordering::Relaxed),
319            rescue_hit_quota: self.rescue_hit_quota.load(Ordering::Relaxed),
320            rescue_over_quota: self.rescue_over_quota.load(Ordering::Relaxed),
321        }
322    }
323}
324
325#[cfg(feature = "stores-redis")]
326fn decision_counter(
327    decision: LimitDecision,
328    allowed: &AtomicU64,
329    hit: &AtomicU64,
330    over: &AtomicU64,
331) {
332    match decision {
333        LimitDecision::Allowed => allowed,
334        LimitDecision::HitQuota => hit,
335        LimitDecision::OverQuota { .. } => over,
336    }
337    .fetch_add(1, Ordering::Relaxed);
338}
339
340/// Atomic Redis token bucket with a process-local rescue bucket.
341///
342/// If Redis fails or times out, callers use the local bucket and one caller periodically probes
343/// for recovery. Dropping the returned future cancels the caller's wait and never grants local
344/// capacity speculatively.
345#[cfg(feature = "stores-redis")]
346#[derive(Clone)]
347pub struct RedisTokenLimiter {
348    store: RedisStore,
349    key: String,
350    rate_per_second: u32,
351    burst: u32,
352    rescue: TokenLimiter,
353    recovery_probe_interval_ms: u64,
354    monitor: Arc<RedisLimiterMonitor>,
355}
356
357#[cfg(feature = "stores-redis")]
358impl RedisTokenLimiter {
359    pub fn new(
360        store: RedisStore,
361        key: impl Into<String>,
362        rate_per_second: u32,
363        burst: u32,
364    ) -> Self {
365        assert!(rate_per_second > 0, "rate must be greater than zero");
366        assert!(burst > 0, "burst must be greater than zero");
367        let key = key.into();
368        assert!(!key.is_empty(), "Redis limiter key must not be empty");
369        Self {
370            store,
371            key,
372            rate_per_second,
373            burst,
374            rescue: TokenLimiter::new(rate_per_second, burst),
375            recovery_probe_interval_ms: 1_000,
376            monitor: Arc::new(RedisLimiterMonitor::default()),
377        }
378    }
379
380    /// Changes how often one caller probes Redis while the backend is unavailable.
381    pub fn with_recovery_probe_interval(mut self, interval: Duration) -> Self {
382        assert!(
383            !interval.is_zero(),
384            "recovery probe interval must be positive"
385        );
386        self.recovery_probe_interval_ms =
387            u64::try_from(interval.as_millis()).expect("recovery probe interval is too large");
388        self
389    }
390
391    pub async fn take(&self, permits: u32) -> LimitDecision {
392        assert!(permits > 0, "permit count must be greater than zero");
393        if permits > self.burst {
394            return LimitDecision::OverQuota {
395                retry_after: Duration::MAX,
396            };
397        }
398        if !self.monitor.should_attempt_redis() {
399            let decision = self.rescue.take(permits);
400            self.monitor.rescue(decision);
401            return decision;
402        }
403        let arguments = [
404            self.rate_per_second.to_string(),
405            self.burst.to_string(),
406            permits.to_string(),
407        ];
408        match self
409            .store
410            .eval::<Vec<i64>, _, _>(TOKEN_BUCKET_SCRIPT, &[&self.key], &arguments)
411            .await
412            .and_then(parse_redis_decision)
413        {
414            Ok(decision) => {
415                self.monitor.redis_success(decision);
416                decision
417            }
418            Err(_) => {
419                self.monitor.redis_failure(self.recovery_probe_interval_ms);
420                let decision = self.rescue.take(permits);
421                self.monitor.rescue(decision);
422                decision
423            }
424        }
425    }
426
427    pub async fn allow(&self) -> bool {
428        self.take(1).await.is_allowed()
429    }
430
431    pub fn snapshot(&self) -> RedisLimiterSnapshot {
432        self.monitor.snapshot()
433    }
434}
435
436#[cfg(feature = "stores-redis")]
437struct BoundedPeriodLimiter {
438    period: Duration,
439    quota: u32,
440    max_keys: usize,
441    windows: Mutex<HashMap<String, PeriodWindow>>,
442}
443
444#[cfg(feature = "stores-redis")]
445impl BoundedPeriodLimiter {
446    fn take(&self, key: &str) -> LimitDecision {
447        let now = unix_millis();
448        let period = self.period.as_millis();
449        let current_window = now / period;
450        let retry_after = duration_from_millis_saturating(
451            (current_window + 1)
452                .saturating_mul(period)
453                .saturating_sub(now),
454        );
455        let mut windows = self.windows.lock().expect("period limiter lock poisoned");
456        if !windows.contains_key(key) && windows.len() >= self.max_keys {
457            windows.retain(|_, window| window.window == current_window);
458            if windows.len() >= self.max_keys {
459                return LimitDecision::OverQuota { retry_after };
460            }
461        }
462        let window = windows.entry(key.to_owned()).or_insert(PeriodWindow {
463            used: 0,
464            window: current_window,
465        });
466        if window.window != current_window {
467            window.used = 0;
468            window.window = current_window;
469        }
470        if window.used >= self.quota {
471            return LimitDecision::OverQuota { retry_after };
472        }
473        window.used += 1;
474        if window.used == self.quota {
475            LimitDecision::HitQuota
476        } else {
477            LimitDecision::Allowed
478        }
479    }
480}
481
482/// Atomic Redis fixed-window limiter keyed by an application identity.
483///
484/// Windows are aligned to Redis server-time boundaries. Application keys are SHA-256 hashed before
485/// use, avoiding accidental disclosure and keeping Redis keys bounded. During Redis failures the
486/// rescue map rejects new identities after `max_rescue_keys` rather than evicting active quotas.
487#[cfg(feature = "stores-redis")]
488#[derive(Clone)]
489pub struct RedisPeriodLimiter {
490    store: RedisStore,
491    namespace: String,
492    period: Duration,
493    quota: u32,
494    rescue: Arc<BoundedPeriodLimiter>,
495    recovery_probe_interval_ms: u64,
496    monitor: Arc<RedisLimiterMonitor>,
497}
498
499#[cfg(feature = "stores-redis")]
500impl RedisPeriodLimiter {
501    pub fn new(
502        store: RedisStore,
503        namespace: impl Into<String>,
504        period: Duration,
505        quota: u32,
506        max_rescue_keys: usize,
507    ) -> Self {
508        assert!(!period.is_zero(), "period must be greater than zero");
509        assert!(period.as_millis() <= u64::MAX.into(), "period is too large");
510        assert!(quota > 0, "quota must be greater than zero");
511        assert!(max_rescue_keys > 0, "rescue key capacity must be positive");
512        let namespace = namespace.into();
513        assert!(
514            !namespace.is_empty(),
515            "Redis limiter namespace must not be empty"
516        );
517        Self {
518            store,
519            namespace,
520            period,
521            quota,
522            rescue: Arc::new(BoundedPeriodLimiter {
523                period,
524                quota,
525                max_keys: max_rescue_keys,
526                windows: Mutex::new(HashMap::new()),
527            }),
528            recovery_probe_interval_ms: 1_000,
529            monitor: Arc::new(RedisLimiterMonitor::default()),
530        }
531    }
532
533    /// Changes how often one caller probes Redis while the backend is unavailable.
534    pub fn with_recovery_probe_interval(mut self, interval: Duration) -> Self {
535        assert!(
536            !interval.is_zero(),
537            "recovery probe interval must be positive"
538        );
539        self.recovery_probe_interval_ms =
540            u64::try_from(interval.as_millis()).expect("recovery probe interval is too large");
541        self
542    }
543
544    pub async fn take(&self, key: &str) -> LimitDecision {
545        if !self.monitor.should_attempt_redis() {
546            let rescue_key = hex_digest(key);
547            let decision = self.rescue.take(&rescue_key);
548            self.monitor.rescue(decision);
549            return decision;
550        }
551        let redis_key = format!("{}:{}", self.namespace, hex_digest(key));
552        let arguments = [self.period.as_millis().to_string(), self.quota.to_string()];
553        match self
554            .store
555            .eval::<Vec<i64>, _, _>(PERIOD_SCRIPT, &[redis_key], &arguments)
556            .await
557            .and_then(parse_redis_decision)
558        {
559            Ok(decision) => {
560                self.monitor.redis_success(decision);
561                decision
562            }
563            Err(_) => {
564                self.monitor.redis_failure(self.recovery_probe_interval_ms);
565                let rescue_key = hex_digest(key);
566                let decision = self.rescue.take(&rescue_key);
567                self.monitor.rescue(decision);
568                decision
569            }
570        }
571    }
572
573    pub fn snapshot(&self) -> RedisLimiterSnapshot {
574        self.monitor.snapshot()
575    }
576}
577
578#[cfg(feature = "stores-redis")]
579fn parse_redis_decision(response: Vec<i64>) -> Result<LimitDecision, RedisStoreError> {
580    let [outcome, retry_after] = response.as_slice() else {
581        return Err(RedisStoreError::UnexpectedResponse(format!(
582            "limiter script returned {response:?}"
583        )));
584    };
585    match outcome {
586        0 => Ok(LimitDecision::Allowed),
587        1 => Ok(LimitDecision::HitQuota),
588        2 if *retry_after >= 0 => Ok(LimitDecision::OverQuota {
589            retry_after: Duration::from_millis(*retry_after as u64),
590        }),
591        _ => Err(RedisStoreError::UnexpectedResponse(format!(
592            "limiter script returned {response:?}"
593        ))),
594    }
595}
596
597#[cfg(feature = "stores-redis")]
598fn hex_digest(value: &str) -> String {
599    let digest = Sha256::digest(value.as_bytes());
600    let mut encoded = String::with_capacity(digest.len() * 2);
601    for byte in digest {
602        use std::fmt::Write as _;
603        write!(encoded, "{byte:02x}").expect("writing to a String cannot fail");
604    }
605    encoded
606}
607
608#[cfg(feature = "stores-redis")]
609fn unix_millis_u64() -> u64 {
610    u64::try_from(unix_millis()).unwrap_or(u64::MAX)
611}
612
613#[cfg(test)]
614mod tests {
615    use super::*;
616
617    #[test]
618    fn token_limiter_honors_burst_capacity() {
619        let limiter = TokenLimiter::new(1, 2);
620        assert!(limiter.allow());
621        assert!(limiter.allow());
622        assert!(matches!(limiter.take(1), LimitDecision::OverQuota { .. }));
623    }
624
625    #[test]
626    fn period_limiter_reports_the_quota_boundary() {
627        let limiter = PeriodLimiter::new(Duration::from_secs(1), 2);
628        assert_eq!(limiter.take("user"), LimitDecision::Allowed);
629        assert_eq!(limiter.take("user"), LimitDecision::HitQuota);
630        assert!(matches!(
631            limiter.take("user"),
632            LimitDecision::OverQuota { .. }
633        ));
634        assert_eq!(limiter.take("other"), LimitDecision::Allowed);
635    }
636
637    #[cfg(feature = "stores-redis")]
638    #[tokio::test]
639    async fn redis_limiter_rescue_is_bounded_and_observable() {
640        let store = RedisStore::new(
641            crate::RedisStoreConfig::new("redis://127.0.0.1:1/")
642                .with_operation_timeout(Duration::from_millis(20)),
643        )
644        .unwrap();
645        let token = RedisTokenLimiter::new(store.clone(), "unavailable-token", 1, 2);
646        assert_eq!(token.take(1).await, LimitDecision::Allowed);
647        assert_eq!(token.take(1).await, LimitDecision::HitQuota);
648        assert!(matches!(
649            token.take(1).await,
650            LimitDecision::OverQuota { .. }
651        ));
652        assert_eq!(token.snapshot().backend_failures, 1);
653        assert_eq!(token.snapshot().rescue_over_quota, 1);
654
655        let period =
656            RedisPeriodLimiter::new(store, "unavailable-period", Duration::from_secs(60), 2, 1);
657        assert_eq!(period.take("known").await, LimitDecision::Allowed);
658        assert_eq!(period.take("known").await, LimitDecision::HitQuota);
659        assert!(matches!(
660            period.take("new-identity").await,
661            LimitDecision::OverQuota { .. }
662        ));
663
664        let monitor = RedisLimiterMonitor::default();
665        monitor.redis_failure(1_000);
666        assert!(!monitor.should_attempt_redis());
667        monitor.next_probe_at_ms.store(0, Ordering::Release);
668        assert!(monitor.should_attempt_redis());
669        monitor.redis_success(LimitDecision::Allowed);
670        assert!(monitor.snapshot().backend_available);
671        assert_eq!(monitor.snapshot().backend_recoveries, 1);
672    }
673
674    #[cfg(feature = "stores-redis")]
675    async fn exercise_real_redis_limiters(store: RedisStore, namespace: &str) {
676        let token_key = format!("{namespace}:token");
677        let period_namespace = format!("{namespace}:period");
678        let period_key = format!("{}:{}", period_namespace, hex_digest("shared-user"));
679        let concurrent_namespace = format!("{namespace}:concurrent");
680        let concurrent_key = format!(
681            "{}:{}",
682            concurrent_namespace,
683            hex_digest("shared-concurrent-user")
684        );
685        store
686            .delete(&[&token_key, &period_key, &concurrent_key])
687            .await
688            .unwrap();
689
690        let token_a = RedisTokenLimiter::new(store.clone(), &token_key, 1, 2);
691        let token_b = RedisTokenLimiter::new(store.clone(), &token_key, 1, 2);
692        assert_eq!(token_a.take(1).await, LimitDecision::Allowed);
693        assert_eq!(token_b.take(1).await, LimitDecision::HitQuota);
694        assert!(matches!(
695            token_a.take(1).await,
696            LimitDecision::OverQuota { .. }
697        ));
698        assert!(token_a.snapshot().backend_available);
699
700        let period_a = RedisPeriodLimiter::new(
701            store.clone(),
702            &period_namespace,
703            Duration::from_secs(60),
704            2,
705            16,
706        );
707        let period_b =
708            RedisPeriodLimiter::new(store, &period_namespace, Duration::from_secs(60), 2, 16);
709        assert_eq!(period_a.take("shared-user").await, LimitDecision::Allowed);
710        assert_eq!(period_b.take("shared-user").await, LimitDecision::HitQuota);
711        let rejected = period_a.take("shared-user").await;
712        assert!(matches!(
713            rejected,
714            LimitDecision::OverQuota { retry_after } if !retry_after.is_zero()
715        ));
716
717        let concurrent = RedisPeriodLimiter::new(
718            period_b.store.clone(),
719            concurrent_namespace,
720            Duration::from_secs(60),
721            8,
722            16,
723        );
724        let decisions = futures::future::join_all((0..32).map(|_| {
725            let limiter = concurrent.clone();
726            async move { limiter.take("shared-concurrent-user").await }
727        }))
728        .await;
729        assert_eq!(
730            decisions
731                .iter()
732                .filter(|decision| decision.is_allowed())
733                .count(),
734            8
735        );
736    }
737
738    #[cfg(feature = "stores-redis")]
739    #[tokio::test]
740    async fn redis_limiter_integration_shares_standalone_quotas() {
741        let Ok(url) = std::env::var("RUST_ZERO_REDIS_URL") else {
742            return;
743        };
744        let namespace = format!("rust-zero-limit:{}", std::process::id());
745        let store = RedisStore::new(
746            crate::RedisStoreConfig::new(url).with_key_prefix(format!("{namespace}:")),
747        )
748        .unwrap();
749        exercise_real_redis_limiters(store, "standalone").await;
750    }
751
752    #[cfg(feature = "stores-redis")]
753    #[tokio::test]
754    async fn redis_cluster_integration_shares_limiter_quotas() {
755        let Ok(nodes) = std::env::var("RUST_ZERO_REDIS_CLUSTER_URLS") else {
756            return;
757        };
758        let store = RedisStore::new(
759            crate::RedisStoreConfig::cluster(
760                nodes
761                    .split(',')
762                    .map(str::trim)
763                    .filter(|node| !node.is_empty()),
764            )
765            .with_key_prefix(format!("rust-zero-limit-cluster:{}:", std::process::id())),
766        )
767        .unwrap();
768        exercise_real_redis_limiters(store, "cluster").await;
769    }
770}