g1t/services/repos/src/resilience.rs

327 lines13,057 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 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
Merge branch 'worktree-agent-a2013627e5ea4ab13'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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily173/// 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
Merge branch 'worktree-agent-a2013627e5ea4ab13'178/// breaker open; or read-only for now, served from the fallback store
179/// (fallback.rs). Seconds to wait.
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily180#[derive(Clone, Copy, Debug, PartialEq, Eq)]
181pub struct Busy {
182 pub rate_limited: bool,
183 pub retry_after: u64,
Merge branch 'worktree-agent-a2013627e5ea4ab13'184 /// Reads work; writes wait until the store is back.
185 pub read_only: bool,
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily186}
187
Merge branch 'worktree-agent-a2013627e5ea4ab13'188/// How long a writer is told to wait while the store is read-only.
189pub const READ_ONLY_RETRY_AFTER: u64 = 300;
190
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily191impl Busy {
Merge branch 'worktree-agent-a2013627e5ea4ab13'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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily197 pub fn error(self, detail: &str) -> worker::Error {
Merge branch 'worktree-agent-a2013627e5ea4ab13'198 let kind = if self.read_only {
199 "read-only"
200 } else if self.rate_limited {
201 "rate-limited"
202 } else {
203 "unavailable"
204 };
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily205 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 {
Merge branch 'worktree-agent-a2013627e5ea4ab13'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 {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily218 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, ':');
Merge branch 'worktree-agent-a2013627e5ea4ab13'229 let kind = parts.next()?;
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily230 let retry_after = parts.next()?.trim().parse().ok()?;
Merge branch 'worktree-agent-a2013627e5ea4ab13'231 Some(Busy { rate_limited: kind == "rate-limited", retry_after, read_only: kind == "read-only" })
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily232}
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() {
Merge branch 'worktree-agent-a2013627e5ea4ab13'310 let busy_error = Busy { rate_limited: true, retry_after: 7, read_only: false }.error("log failed");
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily311 let wrapped = format!("divergence failed: {busy_error}");
Merge branch 'worktree-agent-a2013627e5ea4ab13'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 };
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily319 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}

This file's history is long; its oldest lines are credited to the oldest commit read.