diff --git a/knot2/crates/knot-resource/src/admission.rs b/knot2/crates/knot-resource/src/admission.rs index e90ae1d8d..a246b8396 100644 --- a/knot2/crates/knot-resource/src/admission.rs +++ b/knot2/crates/knot-resource/src/admission.rs @@ -1,10 +1,11 @@ use std::collections::HashMap; use std::hash::Hash; use std::net::IpAddr; +use std::num::{NonZeroU32, NonZeroU64}; use std::sync::{Arc, Mutex}; use std::time::Duration; -use knot_types::UnixMicros; +use knot_types::{AccountDid, UnixMicros}; const MAX_TRACKED_PEERS: usize = 100_000; @@ -12,28 +13,59 @@ const MAX_PACED_KEYS: usize = 4_096; const SWEEP_INTERVAL_MICROS: u64 = 1_000_000; +const MAX_METERED_SUBJECTS: usize = 16_384; + knot_types::scalar_newtype! { - pub struct Burst(u32); - pub struct RefillMicros(u64); pub struct PerPeerInflight(usize); pub struct GlobalInflight(usize); } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct RateLimit { - pub burst: Burst, - pub refill: RefillMicros, +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub struct Burst(NonZeroU32); + +impl Burst { + pub fn of(value: u32) -> Option { + NonZeroU32::new(value).map(Self) + } + + pub const fn get(self) -> u32 { + self.0.get() + } + + pub const fn literal(value: u32) -> Self { + match NonZeroU32::new(value) { + Some(value) => Self(value), + None => panic!("Burst spelled out in source allows at least one call"), + } + } } -impl RateLimit { - pub const fn interval(self) -> RefillMicros { - match self.refill.get() { - 0 => RefillMicros::new(1), - micros => RefillMicros::new(micros), +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub struct RefillMicros(NonZeroU64); + +impl RefillMicros { + pub fn of(value: u64) -> Option { + NonZeroU64::new(value).map(Self) + } + + pub const fn get(self) -> u64 { + self.0.get() + } + + pub const fn literal(value: u64) -> Self { + match NonZeroU64::new(value) { + Some(value) => Self(value), + None => panic!("Refill interval spelled out in source is at least tick"), } } } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct RateLimit { + pub burst: Burst, + pub refill: RefillMicros, +} + #[derive(Debug, Clone, Copy)] pub struct LimitConfig { pub rate: Option, @@ -45,8 +77,8 @@ impl Default for LimitConfig { fn default() -> Self { Self { rate: Some(RateLimit { - burst: Burst::new(20), - refill: RefillMicros::new(100_000), + burst: Burst::literal(20), + refill: RefillMicros::literal(100_000), }), per_peer_inflight: Some(PerPeerInflight::new(8)), global_inflight: Some(GlobalInflight::new(64)), @@ -72,12 +104,12 @@ impl LimitConfig { } } -struct Bucket { +struct Tokens { tokens: u32, last_refill: UnixMicros, } -impl Bucket { +impl Tokens { fn new(rate: RateLimit, now: UnixMicros) -> Self { Self { tokens: rate.burst.get(), @@ -85,16 +117,28 @@ impl Bucket { } } + fn intervals(&self, rate: RateLimit, now: UnixMicros) -> u64 { + now.get().saturating_sub(self.last_refill.get()) / rate.refill.get() + } + fn tokens_at(&self, rate: RateLimit, now: UnixMicros) -> u32 { - let elapsed = now.get().saturating_sub(self.last_refill.get()); - let gained = (elapsed / rate.interval().get()).min(u64::from(rate.burst.get())) as u32; + let gained = self.intervals(rate, now).min(u64::from(rate.burst.get())) as u32; self.tokens.saturating_add(gained).min(rate.burst.get()) } fn replenish(&mut self, rate: RateLimit, now: UnixMicros) -> bool { - if now.get().saturating_sub(self.last_refill.get()) >= rate.interval().get() { - self.tokens = self.tokens_at(rate, now); + if now.get() < self.last_refill.get() { self.last_refill = now; + return self.tokens > 0; + } + let intervals = self.intervals(rate, now); + if intervals > 0 { + self.tokens = self.tokens_at(rate, now); + self.last_refill = UnixMicros::new( + self.last_refill + .get() + .saturating_add(intervals.saturating_mul(rate.refill.get())), + ); } self.tokens > 0 } @@ -105,7 +149,7 @@ impl Bucket { } struct PeerState { - bucket: Option, + bucket: Option, inflight: usize, } @@ -126,13 +170,13 @@ impl PeerState { const CONCENTRATED_REFUSALS: u32 = 1_024; #[derive(Default)] -struct RefusalMajority { +struct RejectionMajority { peer: Option, votes: u32, reported: bool, } -impl RefusalMajority { +impl RejectionMajority { fn observe(&mut self, peer: IpAddr) -> Option { match (self.peer == Some(peer), self.votes) { (true, _) => self.votes = self.votes.saturating_add(1), @@ -152,7 +196,7 @@ struct Inner { peers: HashMap, PeerState>, global_inflight: usize, last_sweep: UnixMicros, - refusal_majority: RefusalMajority, + rejection_majority: RejectionMajority, } impl Inner { @@ -181,7 +225,7 @@ pub struct PreAuthLimiter { } #[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum Refusal { +pub enum Rejection { RateLimited, Saturated, } @@ -200,7 +244,7 @@ impl PreAuthLimiter { peers: HashMap::new(), global_inflight: 0, last_sweep: UnixMicros::new(0), - refusal_majority: RefusalMajority::default(), + rejection_majority: RejectionMajority::default(), }), } } @@ -215,7 +259,7 @@ impl PreAuthLimiter { self: &Arc, peer: Option, now: UnixMicros, - ) -> Result { + ) -> Result { let mut inner = self.lock(); let rate = self.config.rate; let over_global = self @@ -225,10 +269,10 @@ impl PreAuthLimiter { let per_peer = self.config.per_peer_inflight; let decision = match inner.has_room_for(peer, rate, now) { - false => Err(Refusal::Saturated), + false => Err(Rejection::Saturated), true => { let state = inner.peers.entry(peer).or_insert_with(|| PeerState { - bucket: rate.map(|rate| Bucket::new(rate, now)), + bucket: rate.map(|rate| Tokens::new(rate, now)), inflight: 0, }); let ready = match (rate, state.bucket.as_mut()) { @@ -237,8 +281,8 @@ impl PreAuthLimiter { }; let over_peer = per_peer.is_some_and(|limit| state.inflight >= limit.get()); match (ready, over_peer || over_global) { - (false, _) => Err(Refusal::RateLimited), - (_, true) => Err(Refusal::Saturated), + (false, _) => Err(Rejection::RateLimited), + (_, true) => Err(Rejection::Saturated), (true, false) => { if let Some(bucket) = state.bucket.as_mut() { bucket.tokens -= 1; @@ -254,7 +298,7 @@ impl PreAuthLimiter { if inner.peers.get(&peer).is_some_and(PeerState::forgettable) { inner.peers.remove(&peer); } - let concentrated = peer.and_then(|peer| inner.refusal_majority.observe(peer)); + let concentrated = peer.and_then(|peer| inner.rejection_majority.observe(peer)); drop(inner); if let Some(peer) = concentrated { tracing::warn!( @@ -328,11 +372,11 @@ impl HostKey { } #[derive(Debug, Clone, PartialEq, Eq, Hash)] -pub struct SubjectKey(String); +pub struct SubjectKey(AccountDid); impl SubjectKey { - pub fn new(subject: &str) -> Self { - Self(subject.to_ascii_lowercase()) + pub fn new(subject: &AccountDid) -> Self { + Self(subject.clone()) } } @@ -377,7 +421,7 @@ impl Pacer { pub fn reserve(&self, key: &K, now: UnixMicros) -> Duration { let mut booked = self.lock(); match booked.has_room_for(key, now) { - false => Duration::from_micros(self.rate.interval().get()), + false => Duration::from_micros(self.rate.refill.get()), true => { let wait = self.wait_for(&booked, key, now); self.claim_turn(&mut booked, key, now); @@ -403,7 +447,7 @@ impl Pacer { fn wait_for(&self, booked: &Booked, key: &K, now: UnixMicros) -> Duration { let tolerance = self .rate - .interval() + .refill .get() .saturating_mul(u64::from(self.rate.burst.get().saturating_sub(1))); Duration::from_micros( @@ -418,7 +462,7 @@ impl Pacer { let until = self .turn(booked, key, now) .get() - .saturating_add(self.rate.interval().get()); + .saturating_add(self.rate.refill.get()); booked.turns.insert(key.clone(), UnixMicros::new(until)); } @@ -436,6 +480,147 @@ impl Pacer { } } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Bucket { + Peer, + ActorWrites, + ActorReservations, + KnotReservations, + RepoStorage, +} + +impl Bucket { + pub const fn name(self) -> &'static str { + match self { + Self::Peer => "peer", + Self::ActorWrites => "actor-writes", + Self::ActorReservations => "actor-reservations", + Self::KnotReservations => "knot-reservations", + Self::RepoStorage => "repo-storage", + } + } +} + +struct Metering { + tokens: Tokens, + last_seen: UnixMicros, +} + +impl Metering { + fn new(rate: RateLimit, now: UnixMicros) -> Self { + Self { + tokens: Tokens::new(rate, now), + last_seen: now, + } + } + + fn spend(&mut self, rate: RateLimit, now: UnixMicros) -> Result<(), Rejection> { + self.last_seen = now; + if self.tokens.replenish(rate, now) { + self.tokens.tokens -= 1; + Ok(()) + } else { + Err(Rejection::RateLimited) + } + } +} + +struct Metered { + buckets: HashMap, + last_sweep: UnixMicros, +} + +impl Metered { + fn make_room(&mut self, rate: RateLimit, now: UnixMicros) -> bool { + if self.buckets.len() < MAX_METERED_SUBJECTS { + return true; + } + if now.get().saturating_sub(self.last_sweep.get()) >= SWEEP_INTERVAL_MICROS { + self.last_sweep = now; + self.buckets + .retain(|_, metering| !metering.tokens.full(rate, now)); + } + if self.buckets.len() < MAX_METERED_SUBJECTS { + return true; + } + let coldest = self + .buckets + .iter() + .min_by_key(|(_, metering)| metering.last_seen.get()) + .map(|(subject, metering)| (subject.clone(), metering.last_seen)); + + // Freezing cold buckets are provably evictable without consequence. + match coldest { + Some((subject, seen)) if now.get().saturating_sub(seen.get()) >= rate.refill.get() => { + self.buckets.remove(&subject); + true + } + _ => false, + } + } +} + +pub struct SubjectLimiter { + rate: RateLimit, + inner: Mutex, +} + +impl SubjectLimiter { + pub fn new(rate: RateLimit) -> Self { + Self { + rate, + inner: Mutex::new(Metered { + buckets: HashMap::new(), + last_sweep: UnixMicros::new(0), + }), + } + } + + pub fn admit(&self, subject: &SubjectKey, now: UnixMicros) -> Result<(), Rejection> { + let rate = self.rate; + let mut metered = self + .inner + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + if let Some(metering) = metered.buckets.get_mut(subject) { + return metering.spend(rate, now); + } + if !metered.make_room(rate, now) { + return Err(Rejection::Saturated); + } + metered + .buckets + .entry(subject.clone()) + .or_insert_with(|| Metering::new(rate, now)) + .spend(rate, now) + } + + pub fn tracked(&self) -> usize { + self.inner + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .buckets + .len() + } +} + +knot_types::scalar_newtype! { + pub struct SocialBudgetBytes(u64); + pub struct StructureBytes(u64) => ordered; +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct SocialCharge { + pub current: StructureBytes, + pub adding: StructureBytes, +} + +impl SocialBudgetBytes { + pub const fn allows(self, charge: SocialCharge) -> bool { + charge.current.get().saturating_add(charge.adding.get()) <= self.get() + } +} + #[cfg(test)] mod tests { use super::*; @@ -445,11 +630,15 @@ mod tests { Arc::new(PreAuthLimiter::with_config(config)) } + fn limit(burst: u32, refill_micros: u64) -> RateLimit { + RateLimit { + burst: Burst::of(burst).expect("Test burst allows at least one call"), + refill: RefillMicros::of(refill_micros).expect("Test refill is real interval"), + } + } + fn rate(burst: u32, refill_micros: u64) -> Option { - Some(RateLimit { - burst: Burst::new(burst), - refill: RefillMicros::new(refill_micros), - }) + Some(limit(burst, refill_micros)) } fn peer(last: u8) -> Option { @@ -478,7 +667,7 @@ mod tests { #[test] fn only_an_address_taking_most_of_the_refusals_is_reported_and_only_once() { let reported = |rounds, peer: fn(u32) -> u8| { - let mut majority = RefusalMajority::default(); + let mut majority = RejectionMajority::default(); (0..rounds) .filter_map(|round| majority.observe(ip(peer(round)))) .collect::>() @@ -510,7 +699,7 @@ mod tests { assert!(limiter.admit(peer(1), at(0)).is_ok()); assert_eq!( limiter.admit(peer(1), at(0)).err(), - Some(Refusal::RateLimited), + Some(Rejection::RateLimited), "the limiter refuses a third request inside the same instant even when nothing is in flight" ); assert!( @@ -519,7 +708,7 @@ mod tests { ); assert_eq!( limiter.admit(peer(1), at(600)).err(), - Some(Refusal::RateLimited) + Some(Rejection::RateLimited) ); assert!( limiter.admit(peer(1), at(1_000)).is_ok(), @@ -534,29 +723,29 @@ mod tests { per_peer_inflight: Some(PerPeerInflight::new(1)), global_inflight: Some(GlobalInflight::new(2)), }); - let held = limiter + let guard = limiter .admit(peer(1), at(0)) - .expect("the limiter admits the first request"); + .expect("Limiter allows first request"); let other = limiter .admit(peer(2), at(0)) .expect("a second peer fills the global budget"); assert_eq!( limiter.admit(peer(1), at(0)).err(), - Some(Refusal::Saturated), + Some(Rejection::Saturated), "the limiter sheds a second concurrent request from one peer" ); assert_eq!( limiter.admit(peer(3), at(0)).err(), - Some(Refusal::Saturated), + Some(Rejection::Saturated), "the limiter sheds a third peer once the global in-flight limit is reached" ); drop(other); drop( limiter .admit(peer(3), at(0)) - .expect("freeing a global slot admits the peer that was shed"), + .expect("freeing a global slot readmits the peer that was shed"), ); - drop(held); + drop(guard); (0..50).for_each(|_| { limiter .admit(peer(1), at(0)) @@ -570,7 +759,7 @@ mod tests { }); assert_eq!( limiter.admit(peer(1), at(0)).err(), - Some(Refusal::RateLimited), + Some(Rejection::RateLimited), "a dropped guard without a refund keeps its token spent" ); } @@ -589,29 +778,29 @@ mod tests { .map(|_| { per_peer .admit(peer(1), at(0)) - .expect("the limiter admits both concurrent operations from one peer") + .expect("Limiter allows both concurrent operations from one peer") }) .collect(); assert_eq!( per_peer.admit(peer(1), at(0)).err(), - Some(Refusal::Saturated), - "only the per-peer count refuses a request in this budget" + Some(Rejection::Saturated), + "only per-peer count turns requests away in this budget" ); drop(concurrent); - let held: Vec = (0..overflowing) + let guards: Vec = (0..overflowing) .map(|index| { per_peer .admit(rotating(index), at(0)) .expect("a budget that keeps no idle state has no table to overflow") }) .collect(); - assert_eq!(tracked(&per_peer), held.len()); - drop(held); + assert_eq!(tracked(&per_peer), guards.len()); + drop(guards); assert_eq!( tracked(&per_peer), 0, "with no tokens to remember, an idle peer leaves no entry, \ - so address rotation mustn't fill the table and start shedding newcomers" + so address rotation mustn't fill the table and start shedding newcomers" ); let global = limiter(LimitConfig { @@ -621,18 +810,17 @@ mod tests { }); let _saturating = global .admit(peer(1), at(0)) - .expect("the limiter admits the first peer"); + .expect("Limiter allows first peer"); (0..overflowing).for_each(|index| { assert_eq!( global.admit(rotating(index), at(0)).err(), - Some(Refusal::Saturated) + Some(Rejection::Saturated) ); }); assert_eq!( tracked(&global), 1, - "a refusal returns no guard, so an entry it left behind would never be freed, \ - and the sweep that bounds the table only reclaims idle rate state" + "rejection doesn't release its guard, entry it left behind would never be freed, and sweep that bounds table only reclaims idle rate state" ); let unmetered = limiter(LimitConfig::unmetered()); @@ -640,7 +828,7 @@ mod tests { .map(|_| { unmetered .admit(peer(1), at(0)) - .expect("an unmetered budget admits every request from every peer") + .expect("Unmetered budget allows every request from every peer") }) .collect(); drop(guards); @@ -661,12 +849,12 @@ mod tests { limiter .admit(peer(201), at(SWEEP_INTERVAL_MICROS - 1)) .err(), - Some(Refusal::Saturated), + Some(Rejection::Saturated), "inside the sweep interval a full map sheds unseen peers without rescanning" ); assert!( limiter.admit(peer(202), at(SWEEP_INTERVAL_MICROS)).is_ok(), - "once the interval elapses the sweep evicts replenished entries and admits the newcomer" + "once the interval elapses, sweep evicts replenished entries and readmits the newcomer" ); (0..(MAX_TRACKED_PEERS as u64 + 50_000)).for_each(|index| { let _ = limiter.admit( @@ -683,10 +871,7 @@ mod tests { #[test] fn a_pacer_spends_a_hosts_burst_at_once_then_spaces_it_and_leaves_every_other_host_alone() { - let pacer = HostPacer::new(RateLimit { - burst: Burst::new(3), - refill: RefillMicros::new(100), - }); + let pacer = HostPacer::new(limit(3, 100)); let plc = HostKey::new("plc.directory"); let pds = HostKey::new("PDS.Nel.Pet"); let waits: Vec = (0..5) @@ -696,13 +881,13 @@ mod tests { waits, vec![0, 0, 0, 100, 200], "a cold host takes its whole burst without waiting. Every visit after that is one \ - refill interval further out" + refill interval further out" ); assert_eq!( pacer.reserve(&pds, at(0)), Duration::ZERO, "a knot whose accounts spread over many PDSes must fill at the sum of their rates, \ - so one busy host mustn't delay another" + so one busy host mustn't delay another" ); assert_eq!( HostKey::new("PDS.Nel.Pet"), @@ -718,38 +903,27 @@ mod tests { #[test] fn a_turn_that_isnt_due_is_refused_without_pushing_the_schedule_further_out() { - let pacer = SubjectPacer::new(RateLimit { - burst: Burst::new(1), - refill: RefillMicros::new(1_000), - }); - let nel = SubjectKey::new("did:plc:nel"); + let pacer = SubjectPacer::new(limit(1, 1_000)); + let nel = did("nel"); assert!( pacer.reserve_now(&nel, at(0)), "a subject the caller hasn't read takes its turn straight away" ); assert!( !pacer.reserve_now(&nel, at(500)), - "a caller that mustn't wait is refused inside the interval" + "Caller that mustn't wait is turned away inside the interval" ); assert!( !pacer.reserve_now(&nel, at(999)), - "refusals must leave the booking alone, or whoever keeps trying pushes the turn \ - further out every time and the knot never reads the subject again" + "rejections must leave booking alone, or whoever keeps trying pushes their turn \ + further out every time and the knot never reads the subject again" ); assert!(pacer.reserve_now(&nel, at(1_000))); - assert_eq!( - SubjectKey::new("DID:PLC:NEL"), - SubjectKey::new("did:plc:nel"), - "a DID a client typed in mixed case is the same account and shares its schedule" - ); } #[test] fn a_flood_of_distinct_hosts_doesnt_grow_the_schedule_past_its_limit_or_delay_a_newcomer() { - let pacer = HostPacer::new(RateLimit { - burst: Burst::new(1), - refill: RefillMicros::new(1_000_000), - }); + let pacer = HostPacer::new(limit(1, 1_000_000)); (0..MAX_PACED_KEYS as u64 + 5_000).for_each(|index| { let _ = pacer.reserve(&HostKey::new(&format!("{index}.nel.pet")), at(0)); }); @@ -757,13 +931,13 @@ mod tests { assert!( tracked <= MAX_PACED_KEYS, "a grant set spread over more did:web hosts than the schedule can track mustn't \ - grow it past its limit, saw {tracked}" + grow it past its limit, saw {tracked}" ); assert_eq!( pacer.reserve(&HostKey::new("plc.directory"), at(0)), Duration::from_micros(1_000_000), "a host the full schedule can't track waits one refill interval, so a caller that \ - outgrows the schedule slows itself down" + outgrows the schedule slows itself down" ); let settled = at(SWEEP_INTERVAL_MICROS + 2_000_000); assert_eq!( @@ -774,12 +948,12 @@ mod tests { pacer.lock().turns.len(), 1, "the sweep reclaims the schedule and tracks the newcomer once every booked turn \ - has passed" + has passed" ); } #[test] - fn the_majority_counter_sees_a_refusal_from_a_full_peer_table() { + fn the_majority_counter_sees_a_rejection_from_a_full_peer_table() { let limiter = limiter(LimitConfig { rate: rate(1, 1_000), per_peer_inflight: Some(PerPeerInflight::new(8)), @@ -792,13 +966,154 @@ mod tests { limiter .admit(peer(201), at(SWEEP_INTERVAL_MICROS - 1)) .err(), - Some(Refusal::Saturated) + Some(Rejection::Saturated) ); let inner = limiter.lock(); assert_eq!( - (inner.refusal_majority.peer, inner.refusal_majority.votes), + ( + inner.rejection_majority.peer, + inner.rejection_majority.votes + ), (peer(201), 1), - "that refusal has to reach the counter like every other refusal, since a peer shed by a full table is what the warning most needs to report" + "Peer shed by full table is what this warning exists for, and its rejection counts like any other" + ); + } + + fn writes(burst: u32, refill_micros: u64) -> SubjectLimiter { + SubjectLimiter::new(limit(burst, refill_micros)) + } + + fn did(suffix: &str) -> SubjectKey { + SubjectKey::new(&AccountDid::new(format!("did:plc:{suffix}")).unwrap()) + } + + #[test] + fn two_accounts_behind_one_address_spend_their_own_write_budgets() { + let limiter = writes(2, 1_000); + let (nel, olaren) = (did("nel"), did("olaren")); + assert!(limiter.admit(&nel, at(0)).is_ok()); + assert!(limiter.admit(&nel, at(0)).is_ok()); + assert_eq!( + limiter.admit(&nel, at(0)), + Err(Rejection::RateLimited), + "Third write inside one refill window is over burst" + ); + assert!( + limiter.admit(&olaren, at(0)).is_ok(), + "Shared NAT is one address and two writers, and both still write" + ); + assert!( + limiter.admit(&nel, at(1_000)).is_ok(), + "one refill interval later, first writer has its token again" + ); + } + + #[test] + fn one_account_keeps_one_bucket_and_two_dids_never_share_one() { + let limiter = writes(1, 10_000); + let nel = did("nel"); + assert!(limiter.admit(&nel, at(0)).is_ok()); + assert_eq!( + limiter.admit(&nel, at(0)), + Err(Rejection::RateLimited), + "Key is account, not connection account arrived on" + ); + let upper = SubjectKey::new(&AccountDid::new("did:web:knot.test:user:Nel").unwrap()); + let lower = SubjectKey::new(&AccountDid::new("did:web:knot.test:user:nel").unwrap()); + assert!( + limiter.admit(&upper, at(0)).is_ok(), + "two spellings are two accounts, did:web path is case-sensitive" + ); + assert!( + limiter.admit(&lower, at(0)).is_ok(), + "and second spelling has its own budget to spend" + ); + } + + #[test] + fn a_clock_stepped_backwards_doesnt_mint_a_token_or_strand_a_bucket() { + let limiter = writes(2, 1_000); + let nel = did("nel"); + assert!(limiter.admit(&nel, at(5_000)).is_ok()); + assert!( + limiter.admit(&nel, at(4_000)).is_ok(), + "second token of burst is still there, when clock steps back" + ); + assert_eq!( + limiter.admit(&nel, at(4_000)), + Err(Rejection::RateLimited), + "and step back doesn't mint one" + ); + assert!( + limiter.admit(&nel, at(5_000)).is_ok(), + "measuring refill from stepped-back clock keeps bucket from waiting for old future to come around again" + ); + } + + #[test] + fn a_flood_of_accounts_never_grows_the_write_map_without_bound() { + let limiter = writes(1, 10_000); + (0..MAX_METERED_SUBJECTS + 64).for_each(|index| { + let _ = limiter.admit(&did(&format!("flood{index}")), at(0)); + }); + assert!( + limiter.tracked() <= MAX_METERED_SUBJECTS, + "stranger, with a fresh DID per request, would otherwise leak memory" ); } + + #[test] + fn a_full_table_sheds_a_newcomer_until_its_coldest_writer_has_been_quiet_for_a_refill() { + let limiter = writes(1, 1_000_000); + let squatters: Vec = (0..MAX_METERED_SUBJECTS) + .map(|index| did(&format!("squatter{index}"))) + .collect(); + squatters.iter().enumerate().for_each(|(index, squatter)| { + let _ = limiter.admit(squatter, at(index as u64)); + }); + assert_eq!(limiter.tracked(), MAX_METERED_SUBJECTS); + + assert_eq!( + limiter.admit(&did("newcomer"), at(500_000)), + Err(Rejection::Saturated), + "every bucket in full table was spent inside its last refill. forgetting one would hand its account fresh burst, and knot sheds instead" + ); + assert_eq!( + limiter.admit(&squatters[MAX_METERED_SUBJECTS - 1], at(500_000)), + Err(Rejection::RateLimited), + "An account the table already meters is answered by its own bucket, not shed" + ); + assert!( + limiter.admit(&did("newcomer"), at(1_000_000)).is_ok(), + "once coldest writer has been quiet for refill, its bucket makes room" + ); + assert_eq!( + limiter.tracked(), + MAX_METERED_SUBJECTS, + "Newcomer took evicted slot rather than growing table" + ); + } + + #[test] + fn a_social_budget_allows_a_write_that_fits_and_turns_away_the_one_past_it() { + let budget = SocialBudgetBytes::new(1_000); + let cases: &[(u64, u64, bool)] = &[ + (0, 1_000, true), + (900, 100, true), + (900, 101, false), + (1_000, 0, true), + (1_000, 1, false), + (u64::MAX, 1, false), + ]; + cases.iter().for_each(|(current, adding, allowed)| { + assert_eq!( + budget.allows(SocialCharge { + current: StructureBytes::new(*current), + adding: StructureBytes::new(*adding), + }), + *allowed, + "Repository with {current} bytes and taking {adding} more" + ); + }); + } } diff --git a/knot2/crates/knot-resource/src/lib.rs b/knot2/crates/knot-resource/src/lib.rs index 1e401f8c5..fb32d1f72 100644 --- a/knot2/crates/knot-resource/src/lib.rs +++ b/knot2/crates/knot-resource/src/lib.rs @@ -6,8 +6,9 @@ mod mem; mod slots; pub use admission::{ - AdmitGuard, Burst, GlobalInflight, HostKey, HostPacer, LimitConfig, PeerPacer, PerPeerInflight, - PreAuthLimiter, RateLimit, RefillMicros, Refusal, SubjectKey, SubjectPacer, + AdmitGuard, Bucket, Burst, GlobalInflight, HostKey, HostPacer, LimitConfig, PeerPacer, + PerPeerInflight, PreAuthLimiter, RateLimit, RefillMicros, Rejection, SocialBudgetBytes, + SocialCharge, StructureBytes, SubjectKey, SubjectLimiter, SubjectPacer, }; pub use cpu::{Saturate, ThreadCount, gix_thread_limit, map_chunks, map_spans, saturate, threads}; pub use disk::{