//! §5 rate limiting. //! //! A fixed-window counter keyed on `(token_or_ip, surface)`, held in process //! memory — no external counter store. §5 is explicit that a sliding window is //! not worth the complexity at this volume, and that counters resetting on //! restart is acceptable for abuse throttling. //! //! Read limits are applied *behind* the CDN cache, so a cache hit costs a client //! nothing against its budget — that is a deployment property (§8), not //! something this module can enforce. use std::collections::HashMap; use std::sync::Mutex; use std::time::{Duration, Instant}; /// The rate-limited surfaces of §5. Distinct from routes: the batch and single /// forms of `exists` are separate surfaces with separate budgets. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub enum Surface { ExistsSingle, ExistsBatch, ManifestFetch, SeriesFetch, ManifestUpload, BundleUpload, Report, Search, FederationChanges, FederationFetch, FederationHave, FederationPeers, } impl Surface { /// Requests per hour, per §5's table. pub fn limit(self) -> u32 { match self { Surface::ExistsSingle => 600, Surface::ExistsBatch => 60, Surface::ManifestFetch => 300, Surface::SeriesFetch => 120, Surface::ManifestUpload => 100, Surface::BundleUpload => 20, Surface::Report => 20, Surface::Search => 60, // §5: hourly polling is the default, so this allows generous // catch-up without letting one peer walk the feed continuously. Surface::FederationChanges => 120, // Bootstrap pulls are bulk by nature; capped so one peer cannot // saturate egress. Surface::FederationFetch => 5000, Surface::FederationHave => 120, // A public directory read by humans — no reason for volume. Surface::FederationPeers => 60, } } pub fn as_str(self) -> &'static str { match self { Surface::ExistsSingle => "exists", Surface::ExistsBatch => "exists_batch", Surface::ManifestFetch => "manifest_fetch", Surface::SeriesFetch => "series_fetch", Surface::ManifestUpload => "manifest_upload", Surface::BundleUpload => "bundle_upload", Surface::Report => "report", Surface::Search => "search", Surface::FederationChanges => "federation_changes", Surface::FederationFetch => "federation_fetch", Surface::FederationHave => "federation_have", Surface::FederationPeers => "federation_peers", } } } const WINDOW: Duration = Duration::from_secs(3600); /// Headers §5 requires on every rate-limited response. #[derive(Debug, Clone, Copy)] pub struct Quota { pub limit: u32, pub remaining: u32, /// Seconds until the window resets. pub reset: u64, } #[derive(Debug, Clone, Copy)] struct Window { started: Instant, count: u32, } /// TRACES: UR-004 | DR-006 | SR-004 pub struct RateLimiter { windows: Mutex>, } impl Default for RateLimiter { fn default() -> Self { Self::new() } } impl RateLimiter { pub fn new() -> Self { Self { windows: Mutex::new(HashMap::new()) } } /// Records one request against `(key, surface)`. /// /// `Ok(quota)` when within budget, `Err(quota)` when the limit is exceeded — /// in which case the caller returns `429` with `Retry-After` set from /// `quota.reset`. A rejected request does **not** increment the counter, so a /// client hammering a closed window cannot extend its own lockout. pub fn check(&self, key: &str, surface: Surface) -> Result { self.check_at(key, surface, Instant::now()) } fn check_at(&self, key: &str, surface: Surface, now: Instant) -> Result { let limit = surface.limit(); let mut windows = self.windows.lock().expect("rate limiter poisoned"); // Opportunistic eviction of stale windows, so an IP-keyed map cannot // grow without bound behind CGNAT. if windows.len() > 10_000 { windows.retain(|_, w| now.duration_since(w.started) < WINDOW); } let entry = windows.entry((key.to_string(), surface)).or_insert(Window { started: now, count: 0 }); let elapsed = now.duration_since(entry.started); if elapsed >= WINDOW { *entry = Window { started: now, count: 0 }; } let reset = WINDOW.saturating_sub(now.duration_since(entry.started)).as_secs(); if entry.count >= limit { return Err(Quota { limit, remaining: 0, reset }); } entry.count += 1; Ok(Quota { limit, remaining: limit - entry.count, reset }) } } #[cfg(test)] mod tests { use super::*; #[test] fn allows_up_to_the_limit_then_rejects() { let rl = RateLimiter::new(); let limit = Surface::Report.limit(); for i in 0..limit { let q = rl.check("ip", Surface::Report).expect("within budget"); assert_eq!(q.remaining, limit - i - 1); } let q = rl.check("ip", Surface::Report).expect_err("over budget"); assert_eq!(q.remaining, 0); } #[test] fn surfaces_have_independent_budgets() { let rl = RateLimiter::new(); for _ in 0..Surface::BundleUpload.limit() { rl.check("t", Surface::BundleUpload).unwrap(); } assert!(rl.check("t", Surface::BundleUpload).is_err()); // §5: a bundle counts as a single write against its own limit, and must // not consume the single-manifest budget. assert!(rl.check("t", Surface::ManifestUpload).is_ok()); } #[test] fn keys_are_independent() { let rl = RateLimiter::new(); for _ in 0..Surface::Report.limit() { rl.check("a", Surface::Report).unwrap(); } assert!(rl.check("a", Surface::Report).is_err()); assert!(rl.check("b", Surface::Report).is_ok()); } #[test] fn window_resets_after_an_hour() { let rl = RateLimiter::new(); let t0 = Instant::now(); for _ in 0..Surface::Report.limit() { rl.check_at("ip", Surface::Report, t0).unwrap(); } assert!(rl.check_at("ip", Surface::Report, t0).is_err()); // Still closed just inside the window. assert!(rl.check_at("ip", Surface::Report, t0 + Duration::from_secs(3599)).is_err()); // Open again once it rolls over. assert!(rl.check_at("ip", Surface::Report, t0 + Duration::from_secs(3600)).is_ok()); } #[test] fn rejected_requests_do_not_extend_the_lockout() { let rl = RateLimiter::new(); let t0 = Instant::now(); for _ in 0..Surface::Report.limit() { rl.check_at("ip", Surface::Report, t0).unwrap(); } // Hammer the closed window; the counter must not keep climbing, so the // window still expires on schedule. for _ in 0..50 { assert!(rl.check_at("ip", Surface::Report, t0 + Duration::from_secs(10)).is_err()); } assert!(rl.check_at("ip", Surface::Report, t0 + Duration::from_secs(3600)).is_ok()); } #[test] fn reset_counts_down_within_the_window() { let rl = RateLimiter::new(); let t0 = Instant::now(); let q = rl.check_at("ip", Surface::ExistsSingle, t0).unwrap(); assert_eq!(q.reset, 3600); let q = rl.check_at("ip", Surface::ExistsSingle, t0 + Duration::from_secs(600)).unwrap(); assert_eq!(q.reset, 3000); } #[test] fn limits_match_the_spec_table() { assert_eq!(Surface::ExistsSingle.limit(), 600); assert_eq!(Surface::ExistsBatch.limit(), 60); assert_eq!(Surface::ManifestFetch.limit(), 300); assert_eq!(Surface::SeriesFetch.limit(), 120); assert_eq!(Surface::ManifestUpload.limit(), 100); assert_eq!(Surface::BundleUpload.limit(), 20); assert_eq!(Surface::Report.limit(), 20); assert_eq!(Surface::Search.limit(), 60); assert_eq!(Surface::FederationChanges.limit(), 120); assert_eq!(Surface::FederationFetch.limit(), 5000); assert_eq!(Surface::FederationHave.limit(), 120); assert_eq!(Surface::FederationPeers.limit(), 60); } }