Files
dtourolleandClaude Opus 5 545c7d92a2
CI / fmt, clippy, test (push) Failing after 1m20s
CI / static musl binary (push) Has been skipped
CI / advisories and licences (push) Successful in 25s
Federation: replicate content, re-derive judgement (UR-008)
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
2026-07-31 09:28:32 +02:00

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);
}
}