gcra.rsannotatedgcra.rssource326 lines · 11.7 KB · raw
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}