g1t/services/repos/src/namespaces.rs

262 lines10,861 bytesCodeBlame
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
9use std::cell::RefCell;
10use std::collections::HashMap;
11
12use serde::{Deserialize, Serialize};
13use worker::{D1Database, Result};
14
15use crate::shards::{self, Load};
16
17/// How long loads read for placing are kept in an isolate.
18const LOADS_TTL_MS: u64 = 60_000;
19/// The minutes a namespace's recent health is judged over.
20const FAILING_MINUTES: u64 = 5;
21/// The minutes its busiest minute is looked for in.
22const 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.
25const FAILING_SHARE: f64 = 0.25;
26const FAILING_CALLS: f64 = 5.0;
27
28/// What one namespace holds, from the registry.
29#[derive(Debug, Default, Deserialize, Clone, PartialEq)]
30pub 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)]
40pub 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
53impl 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)]
62pub 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
89impl 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.
104pub 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.
116pub 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.
156pub 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.
171pub 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
187thread_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.
192pub 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)]
205mod 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}