g1t/services/repos/src/resilience.rs

327 lines13,057 bytesCodeBlame
1//! What g1t does when the git store fails for a moment.
2//!
3//! - Errors are sorted ([`Failure`]): a rate limit, a failure that may pass
4//! (the store's `INTERNAL_ERROR`, `UPSTREAM_UNAVAILABLE`, a repository
5//! still being made, HTTP 5xx), and an answer that will not change
6//! (`NOT_FOUND`, `ALREADY_EXISTS`, bad input).
7//! - Calls that only read, and minting credentials, are tried again with
8//! exponential backoff and jitter ([`backoff_ms`]). Pushes never are.
9//! - Each namespace has a circuit breaker ([`Breaker`]), per isolate: after
10//! [`TRIP_AFTER`] failures in a row that may pass, calls are refused at
11//! once for [`COOL_MS`], then one probe is let through; it closes the
12//! breaker or opens it again.
13//! - What reaches the caller says so ([`busy`]): git gets 429 or 503 with
14//! `Retry-After`; the site says the git store is busy.
15
16use std::cell::RefCell;
17use std::collections::HashMap;
18
19/// How a failure is treated.
20#[derive(Clone, Copy, Debug, PartialEq, Eq)]
21pub enum Failure {
22 /// The store refused for its rate limit.
23 RateLimited,
24 /// May pass if asked again.
25 Transient,
26 /// The same question gets the same answer.
27 Permanent,
28}
29
30/// Sorts an error the binding threw, by its `code` and message.
31pub fn classify(code: Option<&str>, message: &str) -> Failure {
32 let lower = message.to_ascii_lowercase();
33 if code == Some("RATE_LIMITED")
34 || lower.contains("rate limit")
35 || lower.contains("ratelimit")
36 || lower.contains("too many requests")
37 || lower.contains("status 429")
38 {
39 return Failure::RateLimited;
40 }
41 match code {
42 Some("INTERNAL_ERROR" | "UPSTREAM_UNAVAILABLE" | "CREATE_IN_PROGRESS" | "IMPORT_IN_PROGRESS" | "FORK_IN_PROGRESS") => {
43 Failure::Transient
44 }
45 Some(_) => Failure::Permanent,
46 // No code: the call did not get an answer (a dropped connection, an
47 // overloaded runtime), which may pass.
48 None => Failure::Transient,
49 }
50}
51
52/// Sorts an HTTP answer from the store's git server; `None` when it worked
53/// or failed for good (401, 404, ...).
54pub fn classify_status(status: u16) -> Option<Failure> {
55 match status {
56 429 => Some(Failure::RateLimited),
57 500 | 502 | 503 | 504 => Some(Failure::Transient),
58 _ => None,
59 }
60}
61
62/// Attempts in all for a call that may be tried again.
63pub const ATTEMPTS: u32 = 3;
64const BASE_MS: u64 = 80;
65const RATE_LIMITED_BASE_MS: u64 = 400;
66const MAX_DELAY_MS: u64 = 2_000;
67
68/// How long to wait before attempt `attempt + 1` (`attempt` from 0), with
69/// `jitter` in [0, 1): half the exponential step, plus up to the other half
70/// at random, so callers that failed together do not retry together.
71pub fn backoff_ms(failure: Failure, attempt: u32, jitter: f64) -> u64 {
72 let base = if failure == Failure::RateLimited { RATE_LIMITED_BASE_MS } else { BASE_MS };
73 let step = base.saturating_mul(1 << attempt.min(10)).min(MAX_DELAY_MS);
74 step / 2 + ((step / 2) as f64 * jitter.clamp(0.0, 1.0)) as u64
75}
76
77/// Whether to try again after `failure` on attempt `attempt` (from 0).
78pub fn retry(failure: Failure, attempt: u32) -> bool {
79 failure != Failure::Permanent && attempt + 1 < ATTEMPTS
80}
81
82/// Failures in a row that may pass before a namespace's breaker opens.
83pub const TRIP_AFTER: u32 = 5;
84/// How long an open breaker refuses calls before letting a probe through.
85pub const COOL_MS: u64 = 10_000;
86
87#[derive(Clone, Copy, Debug, PartialEq, Eq)]
88enum State {
89 Closed { failures: u32 },
90 Open { until: u64 },
91 /// One probe is out; everything else waits for it.
92 HalfOpen,
93}
94
95/// What a breaker says about a call.
96#[derive(Clone, Copy, Debug, PartialEq, Eq)]
97pub enum Admit {
98 Go,
99 /// The call is the probe: its outcome decides.
100 Probe,
101 /// Refused; try again in this many milliseconds.
102 Wait(u64),
103}
104
105#[derive(Clone, Copy, Debug, PartialEq, Eq)]
106pub struct Breaker {
107 state: State,
108 /// When the probe went out, so a probe that never reports does not hold
109 /// the breaker half open for good.
110 probe_at: u64,
111}
112
113impl Default for Breaker {
114 fn default() -> Self {
115 Breaker { state: State::Closed { failures: 0 }, probe_at: 0 }
116 }
117}
118
119impl Breaker {
120 pub fn admit(&mut self, now: u64) -> Admit {
121 match self.state {
122 State::Closed { .. } => Admit::Go,
123 State::Open { until } if now < until => Admit::Wait(until - now),
124 State::Open { .. } => {
125 self.state = State::HalfOpen;
126 self.probe_at = now;
127 Admit::Probe
128 }
129 State::HalfOpen if now.saturating_sub(self.probe_at) > COOL_MS => {
130 self.probe_at = now;
131 Admit::Probe
132 }
133 State::HalfOpen => Admit::Wait(1_000),
134 }
135 }
136
137 pub fn succeeded(&mut self) {
138 self.state = State::Closed { failures: 0 };
139 }
140
141 /// Records a failure; only those that may pass count.
142 pub fn failed(&mut self, failure: Failure, now: u64) {
143 if failure == Failure::Permanent {
144 // The store answered: it is up.
145 self.succeeded();
146 return;
147 }
148 self.state = match self.state {
149 State::Closed { failures } if failures + 1 < TRIP_AFTER => State::Closed { failures: failures + 1 },
150 _ => State::Open { until: now + COOL_MS },
151 };
152 }
153
154 pub fn is_open(&self, now: u64) -> bool {
155 matches!(self.state, State::Open { until } if now < until)
156 }
157}
158
159thread_local! {
160 static BREAKERS: RefCell<HashMap<String, Breaker>> = RefCell::new(HashMap::new());
161}
162
163/// The breaker of namespace `store`, in this isolate.
164pub fn with_breaker<T>(store: &str, f: impl FnOnce(&mut Breaker) -> T) -> T {
165 BREAKERS.with(|breakers| f(breakers.borrow_mut().entry(store.to_owned()).or_default()))
166}
167
168/// Whether namespace `store`'s breaker refuses calls in this isolate now.
169pub fn open_now(store: &str, now: u64) -> bool {
170 BREAKERS.with(|breakers| breakers.borrow().get(store).is_some_and(|breaker| breaker.is_open(now)))
171}
172
173/// What marks an error as the git store being busy, through every `?` and
174/// `format!` it passes.
175const BUSY: &str = "git-store-busy:";
176
177/// The git store is busy: rate limited, unavailable after retries, or its
178/// breaker open; or read-only for now, served from the fallback store
179/// (fallback.rs). Seconds to wait.
180#[derive(Clone, Copy, Debug, PartialEq, Eq)]
181pub struct Busy {
182 pub rate_limited: bool,
183 pub retry_after: u64,
184 /// Reads work; writes wait until the store is back.
185 pub read_only: bool,
186}
187
188/// How long a writer is told to wait while the store is read-only.
189pub const READ_ONLY_RETRY_AFTER: u64 = 300;
190
191impl Busy {
192 /// Writes refused while the store is read-only.
193 pub fn read_only() -> Self {
194 Busy { rate_limited: false, retry_after: READ_ONLY_RETRY_AFTER, read_only: true }
195 }
196
197 pub fn error(self, detail: &str) -> worker::Error {
198 let kind = if self.read_only {
199 "read-only"
200 } else if self.rate_limited {
201 "rate-limited"
202 } else {
203 "unavailable"
204 };
205 worker::Error::RustError(format!("{BUSY}{kind}:{}: {detail}", self.retry_after))
206 }
207
208 /// The HTTP status for git: 429 for a rate limit, else 503.
209 pub fn status(self) -> u16 {
210 if self.rate_limited { 429 } else { 503 }
211 }
212
213 /// What people are told.
214 pub fn message(self) -> String {
215 if self.read_only {
216 "g1t's git storage is read-only while it recovers: clones, fetches and pages work, and pushes, merges and new repositories wait until it is back. https://status.g1t.sh has the latest.\n".to_owned()
217 } else if self.rate_limited {
218 format!("g1t's git storage is handling more requests than it allows right now. Try again in {} seconds.\n", self.retry_after)
219 } else {
220 format!("g1t's git storage is not answering right now. Try again in {} seconds; https://status.g1t.sh has the latest.\n", self.retry_after)
221 }
222 }
223}
224
225/// The [`Busy`] an error carries, if it is one.
226pub fn busy(message: &str) -> Option<Busy> {
227 let at = message.find(BUSY)? + BUSY.len();
228 let mut parts = message[at..].splitn(3, ':');
229 let kind = parts.next()?;
230 let retry_after = parts.next()?.trim().parse().ok()?;
231 Some(Busy { rate_limited: kind == "rate-limited", retry_after, read_only: kind == "read-only" })
232}
233
234/// Seconds to tell a caller to wait, from milliseconds, at least one.
235pub fn seconds(ms: u64) -> u64 {
236 ms.div_ceil(1000).max(1)
237}
238
239#[cfg(test)]
240mod tests {
241 use super::*;
242
243 #[test]
244 fn errors_are_sorted_by_code_and_message() {
245 assert_eq!(classify(Some("INTERNAL_ERROR"), "boom"), Failure::Transient);
246 assert_eq!(classify(Some("UPSTREAM_UNAVAILABLE"), ""), Failure::Transient);
247 assert_eq!(classify(Some("FORK_IN_PROGRESS"), ""), Failure::Transient);
248 assert_eq!(classify(Some("NOT_FOUND"), ""), Failure::Permanent);
249 assert_eq!(classify(Some("ALREADY_EXISTS"), ""), Failure::Permanent);
250 assert_eq!(classify(Some("MEMORY_LIMIT"), ""), Failure::Permanent);
251 assert_eq!(classify(None, "Network connection lost."), Failure::Transient);
252 assert_eq!(classify(Some("INTERNAL_ERROR"), "Rate limit exceeded"), Failure::RateLimited);
253 assert_eq!(classify(None, "Too Many Requests"), Failure::RateLimited);
254 assert_eq!(classify_status(429), Some(Failure::RateLimited));
255 assert_eq!(classify_status(503), Some(Failure::Transient));
256 assert_eq!(classify_status(404), None);
257 assert_eq!(classify_status(200), None);
258 }
259
260 #[test]
261 fn backoff_grows_with_jitter_and_stops() {
262 assert_eq!(backoff_ms(Failure::Transient, 0, 0.0), 40);
263 assert_eq!(backoff_ms(Failure::Transient, 0, 0.999), 79);
264 assert_eq!(backoff_ms(Failure::Transient, 1, 0.0), 80);
265 assert_eq!(backoff_ms(Failure::Transient, 2, 0.5), 240);
266 assert_eq!(backoff_ms(Failure::RateLimited, 0, 0.0), 200);
267 assert!(backoff_ms(Failure::Transient, 30, 1.0) <= MAX_DELAY_MS);
268 assert!(retry(Failure::Transient, 0) && retry(Failure::Transient, 1));
269 assert!(!retry(Failure::Transient, ATTEMPTS - 1));
270 assert!(!retry(Failure::Permanent, 0));
271 }
272
273 #[test]
274 fn the_breaker_opens_after_failures_in_a_row_and_probes_after_a_rest() {
275 let mut breaker = Breaker::default();
276 for n in 0..TRIP_AFTER - 1 {
277 breaker.failed(Failure::Transient, 1_000 + u64::from(n));
278 assert_eq!(breaker.admit(1_000), Admit::Go);
279 }
280 // An answer that will not change is an answer: the count starts over.
281 breaker.failed(Failure::Permanent, 1_000);
282 for _ in 0..TRIP_AFTER - 1 {
283 breaker.failed(Failure::RateLimited, 1_000);
284 }
285 assert_eq!(breaker.admit(1_000), Admit::Go);
286 breaker.failed(Failure::Transient, 2_000);
287 assert!(breaker.is_open(2_000));
288 assert_eq!(breaker.admit(2_500), Admit::Wait(COOL_MS - 500));
289 // After the rest, one probe; the rest wait for it.
290 assert_eq!(breaker.admit(2_000 + COOL_MS), Admit::Probe);
291 assert_eq!(breaker.admit(2_000 + COOL_MS + 1), Admit::Wait(1_000));
292 // A failed probe opens it again.
293 breaker.failed(Failure::Transient, 13_000);
294 assert_eq!(breaker.admit(13_001), Admit::Wait(COOL_MS - 1));
295 // A good one closes it.
296 assert_eq!(breaker.admit(13_000 + COOL_MS), Admit::Probe);
297 breaker.succeeded();
298 assert_eq!(breaker.admit(13_000 + COOL_MS), Admit::Go);
299 // A probe that never reported is replaced.
300 let mut stuck = Breaker::default();
301 for _ in 0..TRIP_AFTER {
302 stuck.failed(Failure::Transient, 0);
303 }
304 assert_eq!(stuck.admit(COOL_MS), Admit::Probe);
305 assert_eq!(stuck.admit(2 * COOL_MS + 1), Admit::Probe);
306 }
307
308 #[test]
309 fn busy_survives_being_passed_along() {
310 let busy_error = Busy { rate_limited: true, retry_after: 7, read_only: false }.error("log failed");
311 let wrapped = format!("divergence failed: {busy_error}");
312 assert_eq!(busy(&wrapped), Some(Busy { rate_limited: true, retry_after: 7, read_only: false }));
313 // Read-only: writes wait, said so in words.
314 let read_only = Busy::read_only();
315 assert_eq!(busy(&format!("land: {}", read_only.error("push"))), Some(read_only));
316 assert_eq!(read_only.status(), 503);
317 assert!(read_only.message().contains("read-only"));
318 let down = Busy { rate_limited: false, retry_after: 10, read_only: false };
319 assert_eq!(busy(&down.error("x").to_string()), Some(down));
320 assert_eq!(down.status(), 503);
321 assert!(down.message().contains("10 seconds"));
322 assert_eq!(busy("NOT_FOUND: no such repository"), None);
323 assert_eq!(seconds(1), 1);
324 assert_eq!(seconds(10_000), 10);
325 assert_eq!(seconds(10_001), 11);
326 }
327}