1//! The account's rate limits as GCRA, one atomic word per limit 2//! (contract *Execution*, "The account's rate limit is shared 3//! cluster-wide"). 4//! 5//! Each [`Limit`] word holds a theoretical arrival time (TAT) in 6//! nanoseconds on a clock every backend shares. A request of cost `c` 7//! is admitted at `now` when `TAT - now <= burst`, and moves the TAT to 8//! `max(TAT, now) + c * interval`. A token bucket would need two words 9//! and a lock; this needs one compare-and-swap loop per limit. 10//! 11//! [`Adaptive`] is the shared effective rate: it halves on a 429 and 12//! climbs back after successes, under the hard ceilings the caller 13//! passes in. 14//! 15//! No Postgres and no I/O: the caller supplies the words (in shared 16//! memory, in `postjevsql-pg`), the clock and the wait. 17 18use std::sync::atomic::{AtomicU64, Ordering}; 19use std::time::Duration; 20 21/// A rate: `amount` units per `period`, with room for one `period` of 22/// burst. `amount == 0` means unlimited. 23#[derive(Clone, Copy, Debug, PartialEq)] 24pub struct Rate { 25 pub amount: f64, 26 pub period: Duration, 27} 28 29impl Rate { 30 /// Nanoseconds a unit of cost advances the TAT; `None` when unlimited. 31 fn interval_nanos(self, scale: f64) -> Option<f64> { 32 let per = self.amount * scale; 33 (per > 0.0).then(|| self.period.as_nanos() as f64 / per) 34 } 35} 36 37/// One limit's word. It lives wherever every sharer can reach it. 38pub struct Limit<'a>(pub &'a AtomicU64); 39 40/// What admission decided. 41#[derive(Clone, Copy, Debug, PartialEq, Eq)] 42pub enum Admission { 43 /// Charged; send now. 44 Admitted, 45 /// Nothing was charged; ask again after this long. 46 Wait(Duration), 47} 48 49impl Limit<'_> { 50 /// How long a request of `cost` would wait at `now`, without charging. 51 fn wait(&self, now: u64, rate: Rate, scale: f64) -> Duration { 52 if rate.interval_nanos(scale).is_none() { 53 return Duration::ZERO; 54 } 55 let tat = self.0.load(Ordering::Acquire); 56 let burst = rate.period.as_nanos() as u64; 57 Duration::from_nanos(tat.saturating_sub(now).saturating_sub(burst)) 58 } 59 60 /// Charges `cost` if it fits at `now`; otherwise charges nothing and 61 /// says how long to wait. 62 fn charge(&self, now: u64, cost: f64, rate: Rate, scale: f64) -> Admission { 63 let Some(interval) = rate.interval_nanos(scale) else { return Admission::Admitted }; 64 let burst = rate.period.as_nanos() as u64; 65 let step = (cost.max(0.0) * interval) as u64; 66 let mut tat = self.0.load(Ordering::Acquire); 67 loop { 68 let ahead = tat.saturating_sub(now); 69 if ahead > burst { 70 return Admission::Wait(Duration::from_nanos(ahead - burst)); 71 } 72 let next = tat.max(now).saturating_add(step); 73 match self.0.compare_exchange_weak(tat, next, Ordering::AcqRel, Ordering::Acquire) { 74 Ok(_) => return Admission::Admitted, 75 Err(seen) => tat = seen, 76 } 77 } 78 } 79 80 /// Moves the TAT by `delta` units of cost: a correction once the real 81 /// cost is known (negative returns an overestimate), or the undo of a 82 /// charge the other limit refused. 83 pub fn adjust(&self, delta: f64, rate: Rate, scale: f64) { 84 let Some(interval) = rate.interval_nanos(scale) else { return }; 85 let step = (delta.abs() * interval) as u64; 86 let _ = self.0.fetch_update(Ordering::AcqRel, Ordering::Acquire, |tat| { 87 Some(if delta >= 0.0 { tat.saturating_add(step) } else { tat.saturating_sub(step) }) 88 }); 89 } 90} 91 92/// The shared effective rate, as a fraction of the ceilings in 93/// 1/[`Adaptive::FULL`] steps. 94pub struct Adaptive<'a>(pub &'a AtomicU64); 95 96impl Adaptive<'_> { 97 /// The ceiling itself. A fresh (zeroed) word reads as this. 98 pub const FULL: u64 = 1024; 99 /// Halving stops here: 1/64 of the ceiling. 100 pub const FLOOR: u64 = Self::FULL / 64; 101 /// Each success climbs back this much (additive increase). 102 pub const STEP: u64 = 8; 103 104 /// The fraction of the ceilings in force now, in (0, 1]. 105 pub fn scale(&self) -> f64 { 106 Self::level(self.0.load(Ordering::Acquire)) as f64 / Self::FULL as f64 107 } 108 109 fn level(word: u64) -> u64 { 110 if word == 0 { Self::FULL } else { word.clamp(Self::FLOOR, Self::FULL) } 111 } 112 113 /// A 429: halve, down to the floor. 114 pub fn throttled(&self) { 115 let _ = self.0.fetch_update(Ordering::AcqRel, Ordering::Acquire, |w| { 116 Some((Self::level(w) / 2).max(Self::FLOOR)) 117 }); 118 } 119 120 /// An answer: climb back toward the ceiling. 121 pub fn succeeded(&self) { 122 let _ = self.0.fetch_update(Ordering::AcqRel, Ordering::Acquire, |w| { 123 let level = Self::level(w); 124 (level < Self::FULL).then(|| (level + Self::STEP).min(Self::FULL)) 125 }); 126 } 127} 128 129/// Both of the account's limits and the adaptive scale. 130pub struct Limiter<'a> { 131 pub tokens: Limit<'a>, 132 pub requests: Limit<'a>, 133 pub adaptive: Adaptive<'a>, 134} 135 136/// The ceilings, from `jev.max_tokens_per_second` and 137/// `jev.max_requests_per_minute`. 138#[derive(Clone, Copy, Debug, PartialEq)] 139pub struct Ceilings { 140 pub tokens_per_second: f64, 141 pub requests_per_minute: f64, 142} 143 144impl Ceilings { 145 fn tokens(self) -> Rate { 146 Rate { amount: self.tokens_per_second, period: Duration::from_secs(1) } 147 } 148 149 fn requests(self) -> Rate { 150 Rate { amount: self.requests_per_minute, period: Duration::from_secs(60) } 151 } 152} 153 154impl Limiter<'_> { 155 /// Admits one request estimated at `tokens` at `now` (nanoseconds on 156 /// the shared clock) only if both limits pass. A refusal charges 157 /// neither. 158 pub fn admit(&self, now: u64, tokens: f64, ceilings: Ceilings) -> Admission { 159 let scale = self.adaptive.scale(); 160 // A request bigger than a whole second's tokens is still admitted 161 // once the TAT is within the burst; later requests wait it out. 162 let wait = self 163 .tokens 164 .wait(now, ceilings.tokens(), scale) 165 .max(self.requests.wait(now, ceilings.requests(), scale)); 166 if !wait.is_zero() { 167 return Admission::Wait(wait); 168 } 169 match self.requests.charge(now, 1.0, ceilings.requests(), scale) { 170 Admission::Admitted => {} 171 wait => return wait, 172 } 173 match self.tokens.charge(now, tokens, ceilings.tokens(), scale) { 174 Admission::Admitted => Admission::Admitted, 175 wait => { 176 self.requests.adjust(-1.0, ceilings.requests(), scale); 177 wait 178 } 179 } 180 } 181 182 /// A 429 on a request this limiter admitted: halve the rate, and 183 /// re-price that request's slot at the halved rate, so the next 184 /// request (often its retry) waits the new interval rather than the 185 /// old one. Without the re-price the halving would first show one 186 /// request late. 187 pub fn throttled(&self, ceilings: Ceilings) { 188 let before = self.adaptive.scale(); 189 self.adaptive.throttled(); 190 let after = self.adaptive.scale(); 191 if after < before { 192 self.requests.adjust(1.0, ceilings.requests(), after); 193 self.requests.adjust(-1.0, ceilings.requests(), before); 194 } 195 } 196 197 /// Replaces the estimate with the tokens the answer reported. 198 pub fn correct(&self, estimated: f64, reported: f64, ceilings: Ceilings) { 199 self.tokens.adjust(reported - estimated, ceilings.tokens(), self.adaptive.scale()); 200 } 201} 202 203#[cfg(test)] 204mod tests { 205 use super::*; 206 207 const SEC: u64 = 1_000_000_000; 208 209 struct Words([AtomicU64; 3]); 210 211 impl Words { 212 fn new() -> Self { 213 Words([AtomicU64::new(0), AtomicU64::new(0), AtomicU64::new(0)]) 214 } 215 fn limiter(&self) -> Limiter<'_> { 216 Limiter { tokens: Limit(&self.0[0]), requests: Limit(&self.0[1]), adaptive: Adaptive(&self.0[2]) } 217 } 218 } 219 220 fn ceilings(tps: f64, rpm: f64) -> Ceilings { 221 Ceilings { tokens_per_second: tps, requests_per_minute: rpm } 222 } 223 224 /// Admitted requests over `seconds`, asking as fast as allowed. 225 fn admitted_in(limiter: &Limiter<'_>, start: u64, seconds: u64, tokens: f64, c: Ceilings) -> u64 { 226 let mut now = start; 227 let mut n = 0; 228 while now < start + seconds * SEC { 229 match limiter.admit(now, tokens, c) { 230 Admission::Admitted => n += 1, 231 Admission::Wait(d) => now += d.as_nanos() as u64, 232 } 233 } 234 n 235 } 236 237 #[test] 238 fn unlimited_admits_everything() { 239 let w = Words::new(); 240 for _ in 0..10_000 { 241 assert_eq!(w.limiter().admit(SEC, 1e6, ceilings(0.0, 0.0)), Admission::Admitted); 242 } 243 } 244 245 #[test] 246 fn requests_per_minute_holds_after_the_burst() { 247 let w = Words::new(); 248 let c = ceilings(0.0, 120.0); 249 // A minute's burst of 121 at once, then one more at +0.5 s. 250 let first = admitted_in(&w.limiter(), 100 * SEC, 1, 1.0, c); 251 assert_eq!(first, 122); 252 let next = admitted_in(&w.limiter(), 101 * SEC, 10, 1.0, c); 253 assert!((19..=23).contains(&next), "{next}"); 254 } 255 256 #[test] 257 fn tokens_per_second_holds() { 258 let w = Words::new(); 259 let c = ceilings(1000.0, 0.0); 260 let n = admitted_in(&w.limiter(), 100 * SEC, 10, 100.0, c); 261 // 10 s at 10 requests/s, plus one second of burst. 262 assert!((109..=112).contains(&n), "{n}"); 263 } 264 265 #[test] 266 fn a_refusal_charges_neither_limit() { 267 let w = Words::new(); 268 let c = ceilings(100.0, 6000.0); 269 let l = w.limiter(); 270 assert_eq!(l.admit(100 * SEC, 150.0, c), Admission::Admitted); 271 let requests = w.0[1].load(Ordering::Relaxed); 272 let tokens = w.0[0].load(Ordering::Relaxed); 273 assert!(matches!(l.admit(100 * SEC, 150.0, c), Admission::Wait(_))); 274 assert_eq!(w.0[1].load(Ordering::Relaxed), requests); 275 assert_eq!(w.0[0].load(Ordering::Relaxed), tokens); 276 } 277 278 #[test] 279 fn a_429_halves_and_successes_restore() { 280 let w = Words::new(); 281 let l = w.limiter(); 282 assert_eq!(l.adaptive.scale(), 1.0); 283 l.adaptive.throttled(); 284 assert_eq!(l.adaptive.scale(), 0.5); 285 let c = ceilings(0.0, 120.0); 286 let halved = admitted_in(&l, 101 * SEC, 60, 1.0, c); 287 // A minute's burst at the halved rate, then one a second. 288 assert!((119..=122).contains(&halved), "{halved}"); 289 for _ in 0..1000 { 290 l.adaptive.throttled(); 291 } 292 assert_eq!(l.adaptive.scale(), 1.0 / 64.0); 293 for _ in 0..1000 { 294 l.adaptive.succeeded(); 295 } 296 assert_eq!(l.adaptive.scale(), 1.0); 297 } 298 299 #[test] 300 fn a_429_spaces_the_next_request_at_the_halved_rate() { 301 let w = Words::new(); 302 let c = ceilings(0.0, 60.0); 303 let l = w.limiter(); 304 // Spend the burst, so each request waits its interval. 305 let mut now = 100 * SEC; 306 while l.admit(now, 1.0, c) == Admission::Admitted {} 307 let Admission::Wait(d) = l.admit(now, 1.0, c) else { unreachable!() }; 308 now += d.as_nanos() as u64; 309 assert_eq!(l.admit(now, 1.0, c), Admission::Admitted); 310 l.throttled(c); 311 // At 60/min halved, the next goes 2 s later, not 1 s. 312 let Admission::Wait(d) = l.admit(now, 1.0, c) else { panic!("admitted at once") }; 313 assert!((1_999..=2_001).contains(&d.as_millis()), "{d:?}"); 314 } 315 316 #[test] 317 fn correction_returns_an_overestimate() { 318 let w = Words::new(); 319 let c = ceilings(100.0, 0.0); 320 let l = w.limiter(); 321 assert_eq!(l.admit(100 * SEC, 200.0, c), Admission::Admitted); 322 assert!(matches!(l.admit(100 * SEC, 100.0, c), Admission::Wait(_))); 323 l.correct(200.0, 50.0, c); 324 assert_eq!(l.admit(100 * SEC, 100.0, c), Admission::Admitted); 325 } 326}