gcra.rsannotatedgcra.rssource326 lines · 11.7 KB · raw

The account's rate limits as GCRA, one atomic word per limit (contract Execution, "The account's rate limit is shared cluster-wide").

Each [Limit] word holds a theoretical arrival time (TAT) in nanoseconds on a clock every backend shares. A request of cost c is admitted at now when TAT - now <= burst, and moves the TAT to max(TAT, now) + c * interval. A token bucket would need two words and a lock; this needs one compare-and-swap loop per limit.

[Adaptive] is the shared effective rate: it halves on a 429 and climbs back after successes, under the hard ceilings the caller passes in.

No Postgres and no I/O: the caller supplies the words (in shared memory, in postjevsql-pg), the clock and the wait.

18use std::sync::atomic::{AtomicU64, Ordering};
19use std::time::Duration;

A rate: amount units per period, with room for one period of burst. amount == 0 means unlimited.

23#[derive(Clone, Copy, Debug, PartialEq)]
24pub struct Rate {
25    pub amount: f64,
26    pub period: Duration,
27}
29impl Rate {

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}

One limit's word. It lives wherever every sharer can reach it.

38pub struct Limit<'a>(pub &'a AtomicU64);

What admission decided.

41#[derive(Clone, Copy, Debug, PartialEq, Eq)]
42pub enum Admission {

Charged; send now.

44    Admitted,

Nothing was charged; ask again after this long.

46    Wait(Duration),
47}
49impl Limit<'_> {

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    }

Charges cost if it fits at now; otherwise charges nothing and 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    }

Moves the TAT by delta units of cost: a correction once the real cost is known (negative returns an overestimate), or the undo of a 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}

The shared effective rate, as a fraction of the ceilings in 1/[Adaptive::FULL] steps.

94pub struct Adaptive<'a>(pub &'a AtomicU64);
96impl Adaptive<'_> {

The ceiling itself. A fresh (zeroed) word reads as this.

98    pub const FULL: u64 = 1024;

Halving stops here: 1/64 of the ceiling.

100    pub const FLOOR: u64 = Self::FULL / 64;

Each success climbs back this much (additive increase).

102    pub const STEP: u64 = 8;

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    }
109    fn level(word: u64) -> u64 {
110        if word == 0 { Self::FULL } else { word.clamp(Self::FLOOR, Self::FULL) }
111    }

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    }

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}

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}

The ceilings, from jev.max_tokens_per_second and 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}
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<'_> {

Admits one request estimated at tokens at now (nanoseconds on the shared clock) only if both limits pass. A refusal charges 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    }

A 429 on a request this limiter admitted: halve the rate, and re-price that request's slot at the halved rate, so the next request (often its retry) waits the new interval rather than the old one. Without the re-price the halving would first show one 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    }

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

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