The account's rate limit, shared by every backend of the cluster
(contract Execution, "The account's rate limit is shared
cluster-wide"). postjevsql_core::gcra decides; this module gives it
its three words in shared memory and a clock every backend reads
alike.
The words live in a named DSM segment (GetNamedDSMSegment, PG17+),
so the extension need not be in shared_preload_libraries: the first
backend to ask creates and zeroes it, and the rest attach. Zeroed
words are a limiter at rest: both TATs in the past, and the adaptive
scale at its ceiling.
Attaches this backend to the shared words, creating them if no
backend has yet. Called from begin, never from a future.
34fn words() -> &'static Words { 35 if let Some(words) = WORDS.get() { 36 return words; 37 } 38 let mut found = false; 39 // SAFETY: called on the backend thread, outside any future, where an 40 // ERROR from dsm_registry.c longjmps into the caller's pg_guard. The 41 // segment is pinned and its mapping kept for the backend's life, so 42 // the pointer stays valid as `'static`. `init` runs under the 43 // registry's lock, before any other backend can attach. 44 let ptr = unsafe { 45 pg_sys::GetNamedDSMSegment(c"postjevsql_ratelimit".as_ptr(), size_of::<Words>(), Some(init), &mut found) 46 }; 47 // SAFETY: the segment is `size_of::<Words>()` bytes, suitably aligned 48 // (DSM segments are MAXALIGNed), zeroed by `init`, and only ever 49 // accessed through atomics. 50 let words = unsafe { &*(ptr as *const Words) }; 51 WORDS.set(Some(words)); 52 words 53}
Nanoseconds on CLOCK_MONOTONIC, which every process on the host
shares (an Instant cannot be compared across processes).
The cluster's limiter under this backend's ceilings. Built in begin,
so the segment is attached before any future runs.
80impl RateLimit {
tokens_per_second and requests_per_minute of 0 are unlimited.
ratio prices a request's body until its answer reports the real
count.
Tokens a request is charged before its answer reports the real count: the body's length at the learned characters per token.
Waits until both limits admit a request with this body, and
charges it. The wait is a timer on the executor, so a cancel or
statement_timeout ends it like any other wait.
102 pub async fn admit(&self, body: &[u8]) { 103 let tokens = self.estimate(body); 104 loop { 105 match self.limiter().admit(now(), tokens, self.ceilings) { 106 Admission::Admitted => return, 107 // At least a millisecond, so a rounding sliver cannot spin. 108 Admission::Wait(wait) => executor::sleep(wait.max(Duration::from_millis(1))).await, 109 } 110 } 111 }
The server answered status to a request with this body.
114 pub fn settle(&self, body: &[u8], status: u16, response: &[u8]) { 115 let limiter = self.limiter(); 116 match status { 117 200 => { 118 limiter.adaptive.succeeded(); 119 if let Some(reported) = reported_input_tokens(response) { 120 limiter.correct(self.estimate(body), reported, self.ceilings); 121 } 122 } 123 429 => limiter.throttled(self.ceilings), 124 _ => {} 125 } 126 }
The request never reached the server: nothing was processed, so its tokens are returned. (The request slot stays spent; a redial is rare and the slot is one of 1,200 a minute.)