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 daily | 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 | ||
| 16 | use std::cell::RefCell; | |
| 17 | use std::collections::HashMap; | |
| 18 | ||
| 19 | /// How a failure is treated. | |
| 20 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] | |
| 21 | pub 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. | |
| 31 | pub 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, ...). | |
| 54 | pub 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. | |
| 63 | pub const ATTEMPTS: u32 = 3; | |
| 64 | const BASE_MS: u64 = 80; | |
| 65 | const RATE_LIMITED_BASE_MS: u64 = 400; | |
| 66 | const 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. | |
| 71 | pub 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). | |
| 78 | pub 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. | |
| 83 | pub const TRIP_AFTER: u32 = 5; | |
| 84 | /// How long an open breaker refuses calls before letting a probe through. | |
| 85 | pub const COOL_MS: u64 = 10_000; | |
| 86 | ||
| 87 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] | |
| 88 | enum 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)] | |
| 97 | pub 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)] | |
| 106 | pub 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 | ||
| 113 | impl Default for Breaker { | |
| 114 | fn default() -> Self { | |
| 115 | Breaker { state: State::Closed { failures: 0 }, probe_at: 0 } | |
| 116 | } | |
| 117 | } | |
| 118 | ||
| 119 | impl 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 | ||
| 159 | thread_local! { | |
| 160 | static BREAKERS: RefCell<HashMap<String, Breaker>> = RefCell::new(HashMap::new()); | |
| 161 | } | |
| 162 | ||
| 163 | /// The breaker of namespace `store`, in this isolate. | |
| 164 | pub 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 | ||
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 168 | /// Whether namespace `store`'s breaker refuses calls in this isolate now. |
| 169 | pub 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 daily | 173 | /// What marks an error as the git store being busy, through every `?` and |
| 174 | /// `format!` it passes. | |
| 175 | const BUSY: &str = "git-store-busy:"; | |
| 176 | ||
| 177 | /// The git store is busy: rate limited, unavailable after retries, or its | |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 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 daily | 180 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] |
| 181 | pub struct Busy { | |
| 182 | pub rate_limited: bool, | |
| 183 | pub retry_after: u64, | |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 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 daily | 186 | } |
| 187 | ||
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 188 | /// How long a writer is told to wait while the store is read-only. |
| 189 | pub const READ_ONLY_RETRY_AFTER: u64 = 300; | |
| 190 | ||
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 191 | impl Busy { |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 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 daily | 197 | pub fn error(self, detail: &str) -> worker::Error { |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 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 daily | 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 { | |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 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 daily | 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. | |
| 226 | pub fn busy(message: &str) -> Option<Busy> { | |
| 227 | let at = message.find(BUSY)? + BUSY.len(); | |
| 228 | let mut parts = message[at..].splitn(3, ':'); | |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 229 | let kind = parts.next()?; |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 230 | let retry_after = parts.next()?.trim().parse().ok()?; |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 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 daily | 232 | } |
| 233 | ||
| 234 | /// Seconds to tell a caller to wait, from milliseconds, at least one. | |
| 235 | pub fn seconds(ms: u64) -> u64 { | |
| 236 | ms.div_ceil(1000).max(1) | |
| 237 | } | |
| 238 | ||
| 239 | #[cfg(test)] | |
| 240 | mod 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() { | |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 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 daily | 311 | let wrapped = format!("divergence failed: {busy_error}"); |
| Repositories shard across git store namespaces, move between them, and can keep to the EU; a namespace can be served read-only from the self-hosted git store, rebuilt from the nightly backups (#20, #25) | 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 daily | 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 | } |