use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use dashmap::mapref::entry::Entry;
use super::gcra::{Gcra, Verdict};
const SWEEP_INTERVAL_MS: u64 = 60_000;
const SWEEP_HIGH_WATER: usize = 200_000;
const HW_SWEEP_MIN_MS: u64 = 1_000;
#[derive(Debug)]
pub struct GcraStore {
tats: dashmap::DashMap<String, u64>,
base: Instant,
last_sweep_ms: AtomicU64,
}
impl Default for GcraStore {
fn default() -> Self {
Self::new()
}
}
impl GcraStore {
pub fn new() -> Self {
Self {
tats: dashmap::DashMap::new(),
base: Instant::now(),
last_sweep_ms: AtomicU64::new(0),
}
}
pub fn check(&self, key: &str, gcra: &Gcra) -> Verdict {
let now = self.now();
self.maybe_sweep(now);
match self.tats.entry(key.to_string()) {
Entry::Occupied(mut o) => {
let stored = Some(Duration::from_nanos(*o.get()));
let verdict = gcra.check(stored, now);
if verdict.allowed {
*o.get_mut() = dur_nanos(verdict.new_tat);
}
verdict
}
Entry::Vacant(v) => {
let verdict = gcra.check(None, now);
if verdict.allowed {
v.insert(dur_nanos(verdict.new_tat));
}
verdict
}
}
}
fn now(&self) -> Duration {
self.base.elapsed()
}
fn evict_drained(&self, now: Duration) {
let now_nanos = dur_nanos(now);
self.tats.retain(|_, tat| *tat > now_nanos);
}
fn maybe_sweep(&self, now: Duration) {
let now_ms = u64::try_from(now.as_millis()).unwrap_or(u64::MAX);
let last = self.last_sweep_ms.load(Ordering::Relaxed);
let elapsed = now_ms.saturating_sub(last);
let due = elapsed >= SWEEP_INTERVAL_MS
|| (elapsed >= HW_SWEEP_MIN_MS && self.tats.len() > SWEEP_HIGH_WATER);
if !due {
return;
}
if self
.last_sweep_ms
.compare_exchange(last, now_ms, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
self.evict_drained(now);
}
}
#[cfg(test)]
fn len(&self) -> usize {
self.tats.len()
}
}
fn dur_nanos(d: Duration) -> u64 {
u64::try_from(d.as_nanos()).unwrap_or(u64::MAX)
}
#[cfg(test)]
mod tests {
use super::*;
fn gcra(rate: u64, burst: u64) -> Gcra {
Gcra::from_profile(super::super::gcra::Profile {
rate,
window: Duration::from_secs(60),
burst,
})
}
#[test]
fn admits_burst_then_blocks() {
let store = GcraStore::new();
let g = gcra(60, 2); assert!(store.check("k", &g).allowed);
assert!(store.check("k", &g).allowed);
assert!(!store.check("k", &g).allowed);
}
#[test]
fn keys_are_independent() {
let store = GcraStore::new();
let g = gcra(60, 1); assert!(store.check("a", &g).allowed);
assert!(store.check("b", &g).allowed);
assert!(!store.check("a", &g).allowed);
}
#[test]
fn eviction_reclaims_drained_keys() {
let store = GcraStore::new();
let g = gcra(600, 1); store.check("a", &g);
store.check("b", &g);
assert_eq!(store.len(), 2);
std::thread::sleep(Duration::from_millis(120));
store.evict_drained(store.now());
assert_eq!(store.len(), 0);
}
}