Implements §9a. The replication surface is four reads and no writes: a change feed, fetch by content_id, a batch have, and a human-facing peer directory — plus a capabilities endpoint carrying the accepted envelope versions, which lets a client discover a schema mismatch in one request instead of a 400 per manifest across a library sweep. Pull, never push: a pulling server chooses what it ingests and when. Push would let any peer inject work into the validation queue — the same abuse surface as anonymous upload, at higher volume. Nothing inherits a peer's judgement. A pulled manifest runs the full §6 stage 1 and 2 validation and this server's own cast check, and the fetched body must hash to the content_id that was asked for — the check that stops an intermediary or a misbehaving peer substituting content under a trusted id. A peer's retraction flags for review rather than delisting, because auto-delisting would hand every peer a remote delete primitive; only the opt-in per-peer abuse channel delists, because a takedown propagating at the speed of manual review is the wrong failure mode for that one case. A test caught a real bug in the first cut: the feed cursor was a ULID, and ULIDs are only monotonic *between* milliseconds — two generated in the same millisecond carry independent random components and can sort opposite to write order. A peer resuming from `seq > cursor` would then silently skip an entry: replication losing manifests with no error anywhere. The cursor is now an AUTOINCREMENT integer, and the test asserts strict monotonicity rather than merely sortedness. Peer administration is deliberately not an API. §9a requires that a peering exist only because an operator typed a URL, so nothing a remote server returns can establish or widen one; there_is_no_endpoint_that_creates_a_peering asserts that absence rather than trusting it. 212 tests. Coverage 25/32 (78%). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> TRACES: UR-008 | PR-006
240 lines
8.3 KiB
Rust
240 lines
8.3 KiB
Rust
//! §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<HashMap<(String, Surface), Window>>,
|
|
}
|
|
|
|
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<Quota, Quota> {
|
|
self.check_at(key, surface, Instant::now())
|
|
}
|
|
|
|
fn check_at(&self, key: &str, surface: Surface, now: Instant) -> Result<Quota, Quota> {
|
|
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);
|
|
}
|
|
}
|