g1t/services/repos/src/resilience.rs

299 lines11,675 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily1//! 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 #[cfg(test)]
155 pub fn is_open(&self, now: u64) -> bool {
156 matches!(self.state, State::Open { until } if now < until)
157 }
158}
159
160thread_local! {
161 static BREAKERS: RefCell<HashMap<String, Breaker>> = RefCell::new(HashMap::new());
162}
163
164/// The breaker of namespace `store`, in this isolate.
165pub fn with_breaker<T>(store: &str, f: impl FnOnce(&mut Breaker) -> T) -> T {
166 BREAKERS.with(|breakers| f(breakers.borrow_mut().entry(store.to_owned()).or_default()))
167}
168
169/// What marks an error as the git store being busy, through every `?` and
170/// `format!` it passes.
171const BUSY: &str = "git-store-busy:";
172
173/// The git store is busy: rate limited, unavailable after retries, or its
174/// breaker open. Seconds to wait.
175#[derive(Clone, Copy, Debug, PartialEq, Eq)]
176pub struct Busy {
177 pub rate_limited: bool,
178 pub retry_after: u64,
179}
180
181impl Busy {
182 pub fn error(self, detail: &str) -> worker::Error {
183 let kind = if self.rate_limited { "rate-limited" } else { "unavailable" };
184 worker::Error::RustError(format!("{BUSY}{kind}:{}: {detail}", self.retry_after))
185 }
186
187 /// The HTTP status for git: 429 for a rate limit, else 503.
188 pub fn status(self) -> u16 {
189 if self.rate_limited { 429 } else { 503 }
190 }
191
192 /// What people are told.
193 pub fn message(self) -> String {
194 if self.rate_limited {
195 format!("g1t's git storage is handling more requests than it allows right now. Try again in {} seconds.\n", self.retry_after)
196 } else {
197 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)
198 }
199 }
200}
201
202/// The [`Busy`] an error carries, if it is one.
203pub fn busy(message: &str) -> Option<Busy> {
204 let at = message.find(BUSY)? + BUSY.len();
205 let mut parts = message[at..].splitn(3, ':');
206 let rate_limited = parts.next()? == "rate-limited";
207 let retry_after = parts.next()?.trim().parse().ok()?;
208 Some(Busy { rate_limited, retry_after })
209}
210
211/// Seconds to tell a caller to wait, from milliseconds, at least one.
212pub fn seconds(ms: u64) -> u64 {
213 ms.div_ceil(1000).max(1)
214}
215
216#[cfg(test)]
217mod tests {
218 use super::*;
219
220 #[test]
221 fn errors_are_sorted_by_code_and_message() {
222 assert_eq!(classify(Some("INTERNAL_ERROR"), "boom"), Failure::Transient);
223 assert_eq!(classify(Some("UPSTREAM_UNAVAILABLE"), ""), Failure::Transient);
224 assert_eq!(classify(Some("FORK_IN_PROGRESS"), ""), Failure::Transient);
225 assert_eq!(classify(Some("NOT_FOUND"), ""), Failure::Permanent);
226 assert_eq!(classify(Some("ALREADY_EXISTS"), ""), Failure::Permanent);
227 assert_eq!(classify(Some("MEMORY_LIMIT"), ""), Failure::Permanent);
228 assert_eq!(classify(None, "Network connection lost."), Failure::Transient);
229 assert_eq!(classify(Some("INTERNAL_ERROR"), "Rate limit exceeded"), Failure::RateLimited);
230 assert_eq!(classify(None, "Too Many Requests"), Failure::RateLimited);
231 assert_eq!(classify_status(429), Some(Failure::RateLimited));
232 assert_eq!(classify_status(503), Some(Failure::Transient));
233 assert_eq!(classify_status(404), None);
234 assert_eq!(classify_status(200), None);
235 }
236
237 #[test]
238 fn backoff_grows_with_jitter_and_stops() {
239 assert_eq!(backoff_ms(Failure::Transient, 0, 0.0), 40);
240 assert_eq!(backoff_ms(Failure::Transient, 0, 0.999), 79);
241 assert_eq!(backoff_ms(Failure::Transient, 1, 0.0), 80);
242 assert_eq!(backoff_ms(Failure::Transient, 2, 0.5), 240);
243 assert_eq!(backoff_ms(Failure::RateLimited, 0, 0.0), 200);
244 assert!(backoff_ms(Failure::Transient, 30, 1.0) <= MAX_DELAY_MS);
245 assert!(retry(Failure::Transient, 0) && retry(Failure::Transient, 1));
246 assert!(!retry(Failure::Transient, ATTEMPTS - 1));
247 assert!(!retry(Failure::Permanent, 0));
248 }
249
250 #[test]
251 fn the_breaker_opens_after_failures_in_a_row_and_probes_after_a_rest() {
252 let mut breaker = Breaker::default();
253 for n in 0..TRIP_AFTER - 1 {
254 breaker.failed(Failure::Transient, 1_000 + u64::from(n));
255 assert_eq!(breaker.admit(1_000), Admit::Go);
256 }
257 // An answer that will not change is an answer: the count starts over.
258 breaker.failed(Failure::Permanent, 1_000);
259 for _ in 0..TRIP_AFTER - 1 {
260 breaker.failed(Failure::RateLimited, 1_000);
261 }
262 assert_eq!(breaker.admit(1_000), Admit::Go);
263 breaker.failed(Failure::Transient, 2_000);
264 assert!(breaker.is_open(2_000));
265 assert_eq!(breaker.admit(2_500), Admit::Wait(COOL_MS - 500));
266 // After the rest, one probe; the rest wait for it.
267 assert_eq!(breaker.admit(2_000 + COOL_MS), Admit::Probe);
268 assert_eq!(breaker.admit(2_000 + COOL_MS + 1), Admit::Wait(1_000));
269 // A failed probe opens it again.
270 breaker.failed(Failure::Transient, 13_000);
271 assert_eq!(breaker.admit(13_001), Admit::Wait(COOL_MS - 1));
272 // A good one closes it.
273 assert_eq!(breaker.admit(13_000 + COOL_MS), Admit::Probe);
274 breaker.succeeded();
275 assert_eq!(breaker.admit(13_000 + COOL_MS), Admit::Go);
276 // A probe that never reported is replaced.
277 let mut stuck = Breaker::default();
278 for _ in 0..TRIP_AFTER {
279 stuck.failed(Failure::Transient, 0);
280 }
281 assert_eq!(stuck.admit(COOL_MS), Admit::Probe);
282 assert_eq!(stuck.admit(2 * COOL_MS + 1), Admit::Probe);
283 }
284
285 #[test]
286 fn busy_survives_being_passed_along() {
287 let busy_error = Busy { rate_limited: true, retry_after: 7 }.error("log failed");
288 let wrapped = format!("divergence failed: {busy_error}");
289 assert_eq!(busy(&wrapped), Some(Busy { rate_limited: true, retry_after: 7 }));
290 let down = Busy { rate_limited: false, retry_after: 10 };
291 assert_eq!(busy(&down.error("x").to_string()), Some(down));
292 assert_eq!(down.status(), 503);
293 assert!(down.message().contains("10 seconds"));
294 assert_eq!(busy("NOT_FOUND: no such repository"), None);
295 assert_eq!(seconds(1), 1);
296 assert_eq!(seconds(10_000), 10);
297 assert_eq!(seconds(10_001), 11);
298 }
299}