g1t/services/repos/src/namespaces.rs

262 lines10,861 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.

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)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}