| 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 | #[cfg(test)] |
| 155 | pub fn is_open(&self, now: u64) -> bool { |
| 156 | matches!(self.state, State::Open { until } if now < until) |
| 157 | } |
| 158 | } |
| 159 | |
| 160 | thread_local! { |
| 161 | static BREAKERS: RefCell<HashMap<String, Breaker>> = RefCell::new(HashMap::new()); |
| 162 | } |
| 163 | |
| 164 | /// The breaker of namespace `store`, in this isolate. |
| 165 | pub 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. |
| 171 | const 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)] |
| 176 | pub struct Busy { |
| 177 | pub rate_limited: bool, |
| 178 | pub retry_after: u64, |
| 179 | } |
| 180 | |
| 181 | impl 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. |
| 203 | pub 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. |
| 212 | pub fn seconds(ms: u64) -> u64 { |
| 213 | ms.div_ceil(1000).max(1) |
| 214 | } |
| 215 | |
| 216 | #[cfg(test)] |
| 217 | mod 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 | } |