| 1 | //! How each git store namespace stands: what it holds, how busy and how |
| 2 | //! healthy it has been, and its limits (docs/ARTIFACTS.md, R7). |
| 3 | //! |
| 4 | //! Placing a new repository reads this when there is more than one |
| 5 | //! namespace to choose from (shards.rs `Placement::choose`), kept a minute |
| 6 | //! per isolate; `namespaces` answers it for operators |
| 7 | //! (`scripts/ops/artifacts-namespaces.mjs` reads the same tables itself). |
| 8 | |
| 9 | use std::cell::RefCell; |
| 10 | use std::collections::HashMap; |
| 11 | |
| 12 | use serde::{Deserialize, Serialize}; |
| 13 | use worker::{D1Database, Result}; |
| 14 | |
| 15 | use crate::shards::{self, Load}; |
| 16 | |
| 17 | /// How long loads read for placing are kept in an isolate. |
| 18 | const LOADS_TTL_MS: u64 = 60_000; |
| 19 | /// The minutes a namespace's recent health is judged over. |
| 20 | const FAILING_MINUTES: u64 = 5; |
| 21 | /// The minutes its busiest minute is looked for in. |
| 22 | const PEAK_MINUTES: u64 = 60; |
| 23 | /// A namespace is failing when this share of at least `FAILING_CALLS` of |
| 24 | /// its recent calls failed, as the status page judges it. |
| 25 | const FAILING_SHARE: f64 = 0.25; |
| 26 | const FAILING_CALLS: f64 = 5.0; |
| 27 | |
| 28 | /// What one namespace holds, from the registry. |
| 29 | #[derive(Debug, Default, Deserialize, Clone, PartialEq)] |
| 30 | pub struct Held { |
| 31 | /// The namespace, `''` for the default. |
| 32 | pub ns: String, |
| 33 | pub repos: f64, |
| 34 | pub forks: f64, |
| 35 | pub stored_bytes: f64, |
| 36 | } |
| 37 | |
| 38 | /// How one namespace answered lately, from `store_health`. |
| 39 | #[derive(Debug, Default, Deserialize, Clone, PartialEq)] |
| 40 | pub struct Recent { |
| 41 | pub store: String, |
| 42 | pub peak: f64, |
| 43 | pub calls: f64, |
| 44 | pub errors: f64, |
| 45 | pub rate_limited: f64, |
| 46 | pub rejected: f64, |
| 47 | /// Over the last `FAILING_MINUTES` only. |
| 48 | pub recent_calls: f64, |
| 49 | pub recent_errors: f64, |
| 50 | pub recent_rejected: f64, |
| 51 | } |
| 52 | |
| 53 | impl Recent { |
| 54 | /// Whether its recent calls say it is failing. |
| 55 | pub fn failing(&self) -> bool { |
| 56 | self.recent_rejected > 0.0 || (self.recent_calls >= FAILING_CALLS && self.recent_errors / self.recent_calls >= FAILING_SHARE) |
| 57 | } |
| 58 | } |
| 59 | |
| 60 | /// One namespace, as `namespaces` answers. |
| 61 | #[derive(Debug, Default, Serialize, Clone, PartialEq)] |
| 62 | pub struct Standing { |
| 63 | pub namespace: String, |
| 64 | /// The namespace bound to `ARTIFACTS`, where keys without one live. |
| 65 | pub default: bool, |
| 66 | /// `ARTIFACTS_EU_NAMESPACE`: where EU workspaces' repositories go. |
| 67 | pub eu: bool, |
| 68 | /// Named in `ARTIFACTS_NEW_REPOS`. |
| 69 | pub takes_new_repos: bool, |
| 70 | /// Served from the fallback store now (fallback.rs). |
| 71 | pub on_fallback: bool, |
| 72 | pub writable: bool, |
| 73 | pub repos: u64, |
| 74 | pub forks: u64, |
| 75 | pub stored_bytes: u64, |
| 76 | pub max_repos: Option<u64>, |
| 77 | /// Its busiest minute of calls in the last hour, against Cloudflare's |
| 78 | /// limit for a namespace. |
| 79 | pub peak_per_minute: u64, |
| 80 | pub limit_per_minute: u64, |
| 81 | /// The busiest minute as a share of the limit, 0 to 1 and past it. |
| 82 | pub peak_share: f64, |
| 83 | pub calls_last_hour: u64, |
| 84 | pub errors_last_hour: u64, |
| 85 | pub rate_limited_last_hour: u64, |
| 86 | pub failing: bool, |
| 87 | } |
| 88 | |
| 89 | impl Standing { |
| 90 | pub fn load(&self) -> Load { |
| 91 | Load { |
| 92 | namespace: self.namespace.clone(), |
| 93 | bound: true, |
| 94 | writable: self.writable, |
| 95 | repos: self.repos, |
| 96 | peak_per_minute: self.peak_per_minute, |
| 97 | failing: self.failing, |
| 98 | max_repos: self.max_repos, |
| 99 | } |
| 100 | } |
| 101 | } |
| 102 | |
| 103 | /// What a namespace is configured as, beside what it holds. |
| 104 | pub struct Configured<'a> { |
| 105 | pub bound: &'a [String], |
| 106 | pub default: &'a str, |
| 107 | pub placement: &'a shards::Placement, |
| 108 | pub limits: &'a HashMap<String, u64>, |
| 109 | pub on_fallback: &'a dyn Fn(&str) -> bool, |
| 110 | pub writable: &'a dyn Fn(&str) -> bool, |
| 111 | /// Whether this isolate's breaker for it is open now. |
| 112 | pub breaker_open: &'a dyn Fn(&str) -> bool, |
| 113 | } |
| 114 | |
| 115 | /// Every bound namespace's standing, from what the tables say. |
| 116 | pub fn standings(config: &Configured<'_>, held: &[Held], recent: &[Recent]) -> Vec<Standing> { |
| 117 | config |
| 118 | .bound |
| 119 | .iter() |
| 120 | .map(|namespace| { |
| 121 | let holds = held |
| 122 | .iter() |
| 123 | .filter(|row| (if row.ns.is_empty() { config.default } else { row.ns.as_str() }) == namespace) |
| 124 | .fold(Held::default(), |sum, row| Held { |
| 125 | ns: String::new(), |
| 126 | repos: sum.repos + row.repos, |
| 127 | forks: sum.forks + row.forks, |
| 128 | stored_bytes: sum.stored_bytes + row.stored_bytes, |
| 129 | }); |
| 130 | let health = recent.iter().find(|row| row.store == *namespace).cloned().unwrap_or_default(); |
| 131 | let peak = health.peak.max(0.0) as u64; |
| 132 | Standing { |
| 133 | namespace: namespace.clone(), |
| 134 | default: namespace == config.default, |
| 135 | eu: config.placement.eu.as_deref() == Some(namespace.as_str()), |
| 136 | takes_new_repos: config.placement.new_repos.contains(namespace), |
| 137 | on_fallback: (config.on_fallback)(namespace), |
| 138 | writable: (config.writable)(namespace), |
| 139 | repos: holds.repos.max(0.0) as u64, |
| 140 | forks: holds.forks.max(0.0) as u64, |
| 141 | stored_bytes: holds.stored_bytes.max(0.0) as u64, |
| 142 | max_repos: config.limits.get(namespace).copied(), |
| 143 | peak_per_minute: peak, |
| 144 | limit_per_minute: shards::CONTROL_PLANE_PER_MINUTE, |
| 145 | peak_share: peak as f64 / shards::CONTROL_PLANE_PER_MINUTE as f64, |
| 146 | calls_last_hour: health.calls.max(0.0) as u64, |
| 147 | errors_last_hour: health.errors.max(0.0) as u64, |
| 148 | rate_limited_last_hour: health.rate_limited.max(0.0) as u64, |
| 149 | failing: health.failing() || (config.breaker_open)(namespace), |
| 150 | } |
| 151 | }) |
| 152 | .collect() |
| 153 | } |
| 154 | |
| 155 | /// What the registry holds, by namespace. |
| 156 | pub async fn held(db: &D1Database) -> Result<Vec<Held>> { |
| 157 | db.prepare( |
| 158 | "SELECT CASE WHEN instr(coalesce(store, ''), '/') > 0 THEN substr(store, 1, instr(store, '/') - 1) ELSE '' END AS ns, |
| 159 | count(*) AS repos, |
| 160 | sum(CASE WHEN fork_of IS NULL THEN 0 ELSE 1 END) AS forks, |
| 161 | sum(coalesce(stored_bytes, 0)) AS stored_bytes |
| 162 | FROM repos WHERE deleted_at IS NULL AND retired_at IS NULL |
| 163 | GROUP BY ns", |
| 164 | ) |
| 165 | .all() |
| 166 | .await? |
| 167 | .results::<Held>() |
| 168 | } |
| 169 | |
| 170 | /// How each namespace answered over the last hour, and the last minutes. |
| 171 | pub async fn recent(db: &D1Database, now: u64) -> Result<Vec<Recent>> { |
| 172 | let minute = |ago: u64| g1t_contracts::time::rfc3339(now.saturating_sub(ago * 60_000))[..16].to_owned(); |
| 173 | db.prepare( |
| 174 | "SELECT store, max(calls) AS peak, sum(calls) AS calls, sum(errors) AS errors, |
| 175 | sum(rate_limited) AS rate_limited, sum(rejected) AS rejected, |
| 176 | sum(CASE WHEN minute >= ?2 THEN calls ELSE 0 END) AS recent_calls, |
| 177 | sum(CASE WHEN minute >= ?2 THEN errors ELSE 0 END) AS recent_errors, |
| 178 | sum(CASE WHEN minute >= ?2 THEN rejected ELSE 0 END) AS recent_rejected |
| 179 | FROM store_health WHERE minute >= ?1 GROUP BY store", |
| 180 | ) |
| 181 | .bind(&[minute(PEAK_MINUTES).into(), minute(FAILING_MINUTES).into()])? |
| 182 | .all() |
| 183 | .await? |
| 184 | .results::<Recent>() |
| 185 | } |
| 186 | |
| 187 | thread_local! { |
| 188 | static LOADS: RefCell<Option<(u64, Vec<Load>)>> = const { RefCell::new(None) }; |
| 189 | } |
| 190 | |
| 191 | /// Loads for placing, kept a minute in this isolate. |
| 192 | pub async fn loads(db: &D1Database, config: &Configured<'_>, now: u64) -> Result<Vec<Load>> { |
| 193 | if let Some(loads) = LOADS.with(|kept| { |
| 194 | kept.borrow().as_ref().filter(|(at, _)| now.saturating_sub(*at) < LOADS_TTL_MS).map(|(_, loads)| loads.clone()) |
| 195 | }) { |
| 196 | return Ok(loads); |
| 197 | } |
| 198 | let (held, recent) = futures_util::future::join(held(db), recent(db, now)).await; |
| 199 | let loads: Vec<Load> = standings(config, &held?, &recent?).iter().map(Standing::load).collect(); |
| 200 | LOADS.with(|kept| *kept.borrow_mut() = Some((now, loads.clone()))); |
| 201 | Ok(loads) |
| 202 | } |
| 203 | |
| 204 | #[cfg(test)] |
| 205 | mod tests { |
| 206 | use super::*; |
| 207 | |
| 208 | #[test] |
| 209 | fn standings_add_up_what_each_namespace_holds_and_how_it_answered() { |
| 210 | let bound = vec!["g1t".to_owned(), "g1t-us-1".to_owned(), "g1t-eu".to_owned()]; |
| 211 | let placement = shards::Placement::from_vars(Some("g1t,g1t-us-1"), Some("g1t-eu")); |
| 212 | let limits = HashMap::from([("g1t".to_owned(), 10_000u64)]); |
| 213 | let config = Configured { |
| 214 | bound: &bound, |
| 215 | default: "g1t", |
| 216 | placement: &placement, |
| 217 | limits: &limits, |
| 218 | on_fallback: &|name| name == "g1t-us-1", |
| 219 | writable: &|name| name != "g1t-us-1", |
| 220 | breaker_open: &|name| name == "g1t-eu", |
| 221 | }; |
| 222 | let held = vec![ |
| 223 | // Keys with no namespace, and keys naming the default, are both its. |
| 224 | Held { ns: String::new(), repos: 90.0, forks: 10.0, stored_bytes: 1_000.0 }, |
| 225 | Held { ns: "g1t".into(), repos: 10.0, forks: 0.0, stored_bytes: 24.0 }, |
| 226 | Held { ns: "g1t-us-1".into(), repos: 5.0, forks: 1.0, stored_bytes: 7.0 }, |
| 227 | ]; |
| 228 | let recent = vec![ |
| 229 | Recent { store: "g1t".into(), peak: 6_000.0, calls: 50_000.0, errors: 3.0, recent_calls: 100.0, recent_errors: 1.0, ..Recent::default() }, |
| 230 | Recent { store: "g1t-us-1".into(), peak: 10.0, calls: 20.0, recent_calls: 8.0, recent_errors: 2.0, ..Recent::default() }, |
| 231 | // The fallback's own health is not the namespace's. |
| 232 | Recent { store: "g1t-us-1@fallback".into(), recent_calls: 50.0, recent_errors: 50.0, ..Recent::default() }, |
| 233 | ]; |
| 234 | let out = standings(&config, &held, &recent); |
| 235 | assert_eq!(out.len(), 3); |
| 236 | let g1t = &out[0]; |
| 237 | assert_eq!((g1t.repos, g1t.forks, g1t.stored_bytes), (100, 10, 1_024)); |
| 238 | assert!(g1t.default && g1t.takes_new_repos && !g1t.eu && g1t.writable && !g1t.failing); |
| 239 | assert_eq!(g1t.max_repos, Some(10_000)); |
| 240 | assert_eq!(g1t.peak_per_minute, 6_000); |
| 241 | assert!((g1t.peak_share - 0.5).abs() < 1e-9); |
| 242 | let us1 = &out[1]; |
| 243 | assert!(us1.on_fallback && !us1.writable); |
| 244 | // 2 of 8 recent calls failed: a quarter, so failing. |
| 245 | assert!(us1.failing); |
| 246 | let eu = &out[2]; |
| 247 | assert!(eu.eu && !eu.takes_new_repos); |
| 248 | assert_eq!(eu.repos, 0); |
| 249 | // Its breaker is open here. |
| 250 | assert!(eu.failing); |
| 251 | let load = g1t.load(); |
| 252 | assert_eq!((load.repos, load.peak_per_minute, load.max_repos), (100, 6_000, Some(10_000))); |
| 253 | } |
| 254 | |
| 255 | #[test] |
| 256 | fn a_few_failures_are_not_failing() { |
| 257 | assert!(!Recent { recent_calls: 4.0, recent_errors: 4.0, ..Recent::default() }.failing()); |
| 258 | assert!(!Recent { recent_calls: 100.0, recent_errors: 24.0, ..Recent::default() }.failing()); |
| 259 | assert!(Recent { recent_calls: 100.0, recent_errors: 25.0, ..Recent::default() }.failing()); |
| 260 | assert!(Recent { recent_rejected: 1.0, ..Recent::default() }.failing()); |
| 261 | } |
| 262 | } |