g1t/services/repos/src/shards.rs

359 lines16,810 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.

Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily1//! Which git store namespace a repository lives in.
2//!
3//! Cloudflare limits each Artifacts namespace to 2,000 control-plane
4//! requests per 10 seconds, and fixes its jurisdiction (US or EU) when it
5//! is made. So repositories can live in several namespaces, each reached
6//! through its own binding (`ARTIFACTS`, `ARTIFACTS_1`, ..., `ARTIFACTS_EU`),
7//! named in the `ARTIFACTS_NAMESPACES` variable:
8//!
9//! ```json
10//! { "ARTIFACTS": "g1t", "ARTIFACTS_1": "g1t-us-1", "ARTIFACTS_EU": "g1t-eu" }
11//! ```
12//!
13//! A repository's namespace is part of its store key in the registry's
14//! `store` column: `g1t-us-1/acme--rocket`. A key with no namespace
15//! (`acme--rocket`, every repository made before this) is in the namespace
16//! bound to `ARTIFACTS`. A pull request's working copy is always in its
17//! repository's namespace, since Artifacts forks within a namespace.
18//!
19//! New repositories go where `ARTIFACTS_NEW_REPOS` says (a comma-separated
Merge branch 'worktree-agent-a2013627e5ea4ab13'20//! list of namespaces): the emptier healthy ones, spread by repository id
21//! ([`Placement::choose`], from namespaces.rs's loads). A workspace that
22//! keeps its data in the EU (identity's `data_residency`, offered once
23//! `ARTIFACTS_EU_NAMESPACE` names a bound namespace) gets that namespace,
24//! or nothing. `ARTIFACTS_NAMESPACE_LIMITS` caps how many repositories a
25//! namespace takes. Without any of these, everything stays in
26//! `ARTIFACTS`'s, and nothing extra is read. An existing repository moves
27//! between namespaces only when an operator asks (moves.rs).
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily28
29/// The binding every installation has.
30pub const DEFAULT_BINDING: &str = "ARTIFACTS";
31/// Its namespace, unless `ARTIFACTS_NAMESPACES` says otherwise.
32pub const DEFAULT_NAMESPACE: &str = "g1t";
33
34/// Each binding and the namespace it reaches, the default first.
35pub fn bindings(config: Option<&str>) -> Vec<(String, String)> {
36 let mut out: Vec<(String, String)> = config
37 .and_then(|text| serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(text).ok())
38 .map(|map| {
39 map.into_iter()
40 .filter_map(|(binding, namespace)| Some((binding, namespace.as_str()?.trim().to_owned())))
41 .filter(|(_, namespace)| valid_namespace(namespace))
42 .collect()
43 })
44 .unwrap_or_default();
45 if !out.iter().any(|(binding, _)| binding == DEFAULT_BINDING) {
46 out.push((DEFAULT_BINDING.to_owned(), DEFAULT_NAMESPACE.to_owned()));
47 }
48 out.sort_by_key(|(binding, _)| (binding != DEFAULT_BINDING, binding.clone()));
49 out
50}
51
52/// Cloudflare's rule for names: 2 to 63 letters, digits, `.`, `_`, `-`,
53/// starting with a letter or digit.
54pub fn valid_namespace(name: &str) -> bool {
55 (2..=63).contains(&name.len())
56 && name.chars().next().is_some_and(|c| c.is_ascii_alphanumeric())
57 && name.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'))
58}
59
60/// A store key's namespace (`None`: the default one) and its name there.
61pub fn split(key: &str) -> (Option<&str>, &str) {
62 match key.split_once('/') {
63 Some((namespace, name)) => (Some(namespace), name),
64 None => (None, key),
65 }
66}
67
68/// The store key for `name` in `namespace`; the default one is left out,
69/// so keys made before sharding read the same.
70pub fn compose(namespace: Option<&str>, name: &str, default: &str) -> String {
71 match namespace {
72 Some(namespace) if namespace != default => format!("{namespace}/{name}"),
73 _ => name.to_owned(),
74 }
75}
76
77/// Where a workspace keeps its data.
78#[derive(Clone, Copy, Debug, PartialEq, Eq)]
79pub enum Residency {
80 Anywhere,
81 Eu,
82}
83
Merge branch 'worktree-agent-a2013627e5ea4ab13'84/// Cloudflare's control-plane limit for one namespace, per minute: 2,000
85/// requests per 10 seconds.
86pub const CONTROL_PLANE_PER_MINUTE: u64 = 12_000;
87/// Past this share of the limit in its busiest recent minute, a namespace
88/// is hot and takes no new repositories while another can.
89pub const HOT_SHARE: f64 = 0.7;
90/// New repositories go to the namespaces holding no more than this many
91/// above the emptiest, or this share above it, spread among them by id,
92/// so a burst between counts does not all land in one place.
93const SLACK_REPOS: u64 = 100;
94const SLACK_SHARE: f64 = 0.05;
95
96/// The most each namespace should hold, from `ARTIFACTS_NAMESPACE_LIMITS`
97/// (JSON, `{"g1t": {"max_repos": 50000}}`); a namespace not named has no
98/// limit but the control-plane one.
99pub fn limits(config: Option<&str>) -> std::collections::HashMap<String, u64> {
100 config
101 .and_then(|text| serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(text).ok())
102 .map(|map| {
103 map.into_iter()
104 .filter_map(|(namespace, limit)| Some((namespace, limit.get("max_repos")?.as_u64()?)))
105 .collect()
106 })
107 .unwrap_or_default()
108}
109
110/// How a namespace stands, for placing a repository in it.
111#[derive(Clone, Debug, Default, PartialEq)]
112pub struct Load {
113 pub namespace: String,
114 /// Bound to a binding here.
115 pub bound: bool,
116 /// Takes writes now: not while served from a read-only fallback.
117 pub writable: bool,
118 /// Repositories and working copies in it (not deleted).
119 pub repos: u64,
120 /// Its busiest minute of calls lately (meters.rs `store_health`).
121 pub peak_per_minute: u64,
122 /// Failing now: its breaker open here, or a quarter of its recent
123 /// calls failing.
124 pub failing: bool,
125 /// `ARTIFACTS_NAMESPACE_LIMITS`'s `max_repos` for it.
126 pub max_repos: Option<u64>,
127}
128
129impl Load {
130 /// Whether it can take a repository at all.
131 fn usable(&self) -> bool {
132 self.bound && self.writable
133 }
134
135 /// Whether it should, when another could.
136 fn healthy(&self) -> bool {
137 self.usable()
138 && !self.failing
139 && self.max_repos.is_none_or(|max| self.repos < max)
140 && (self.peak_per_minute as f64) < CONTROL_PLANE_PER_MINUTE as f64 * HOT_SHARE
141 }
142}
143
144/// Why a repository could not be placed.
145#[derive(Debug, PartialEq, Eq)]
146pub enum Unplaced {
147 /// The workspace keeps its data in the EU, and there is no EU
148 /// namespace to take it now.
149 NoEu,
150}
151
152impl Unplaced {
153 pub fn message(&self) -> &'static str {
154 match self {
155 Unplaced::NoEu => {
156 "This workspace keeps its data in the EU, and EU storage cannot take new repositories right now. Try again later, or ask support."
157 }
158 }
159 }
160}
161
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily162/// Where new repositories go.
163#[derive(Clone, Debug, Default, PartialEq, Eq)]
164pub struct Placement {
165 /// Namespaces that take new repositories; empty for the default.
166 pub new_repos: Vec<String>,
167 pub eu: Option<String>,
168}
169
170impl Placement {
171 pub fn from_vars(new_repos: Option<&str>, eu: Option<&str>) -> Self {
172 Placement {
173 new_repos: new_repos
174 .unwrap_or_default()
175 .split(',')
176 .map(str::trim)
177 .filter(|name| valid_namespace(name))
178 .map(str::to_owned)
179 .collect(),
180 eu: eu.map(str::trim).filter(|name| valid_namespace(name)).map(str::to_owned),
181 }
182 }
183
Merge branch 'worktree-agent-a2013627e5ea4ab13'184 /// Whether loads are worth reading: more than one namespace could take
185 /// the repository. With one or none, `choose` needs none.
186 pub fn needs_loads(&self, bound: &[String]) -> bool {
187 self.new_repos.iter().filter(|name| bound.contains(name)).count() > 1
188 }
189
190 /// Whether an EU namespace is there to take repositories.
191 pub fn eu_available(&self, bound: &[String], writable: impl Fn(&str) -> bool) -> bool {
192 self.eu.as_deref().is_some_and(|eu| bound.iter().any(|name| name == eu) && writable(eu))
193 }
194
195 /// The namespace for a new repository with this id, from how each
196 /// stands (`loads`); `Ok(None)` for the default.
197 ///
198 /// - EU residency: the EU namespace, if it is bound and takes writes;
199 /// otherwise refused, never placed elsewhere.
200 /// - Anywhere: among the namespaces named in `ARTIFACTS_NEW_REPOS` that
201 /// are healthy (bound, writable, not failing, under their limit and
202 /// not hot), those within the slack of the emptiest, spread by id. If
203 /// none is healthy, the usable ones the same way; if none is usable,
204 /// the default. A namespace with no load given is passed over, as one
205 /// not bound is.
206 pub fn choose(&self, repo_id: &str, residency: Residency, loads: &[Load]) -> Result<Option<String>, Unplaced> {
207 let load_of = |name: &str| loads.iter().find(|load| load.namespace == name);
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily208 if residency == Residency::Eu {
Merge branch 'worktree-agent-a2013627e5ea4ab13'209 return match self.eu.as_deref().and_then(load_of) {
210 Some(load) if load.usable() => Ok(Some(load.namespace.clone())),
211 _ => Err(Unplaced::NoEu),
212 };
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily213 }
Merge branch 'worktree-agent-a2013627e5ea4ab13'214 let named: Vec<&Load> = self.new_repos.iter().filter_map(|name| load_of(name)).collect();
215 let healthy: Vec<&Load> = named.iter().copied().filter(|load| load.healthy()).collect();
216 let pool = if healthy.is_empty() { named.into_iter().filter(|load| load.usable()).collect() } else { healthy };
217 let Some(emptiest) = pool.iter().map(|load| load.repos).min() else {
218 return Ok(None);
219 };
220 let ceiling = emptiest.saturating_add(SLACK_REPOS.max((emptiest as f64 * SLACK_SHARE) as u64));
221 let eligible: Vec<&Load> = pool.into_iter().filter(|load| load.repos <= ceiling).collect();
222 Ok(Some(eligible[(fnv(repo_id) % eligible.len() as u64) as usize].namespace.clone()))
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily223 }
224}
225
226/// FNV-1a, to spread ids over namespaces the same way every time.
227fn fnv(text: &str) -> u64 {
228 text.bytes().fold(0xcbf2_9ce4_8422_2325, |hash, byte| (hash ^ u64::from(byte)).wrapping_mul(0x0100_0000_01b3))
229}
230
231#[cfg(test)]
232mod tests {
233 use super::*;
234
235 #[test]
236 fn the_default_binding_is_always_there_and_first() {
237 assert_eq!(bindings(None), vec![("ARTIFACTS".to_owned(), "g1t".to_owned())]);
238 let configured = bindings(Some(r#"{"ARTIFACTS_EU":"g1t-eu","ARTIFACTS_1":"g1t-us-1","ARTIFACTS":"g1t","ARTIFACTS_2":"bad name!"}"#));
239 assert_eq!(
240 configured,
241 vec![
242 ("ARTIFACTS".to_owned(), "g1t".to_owned()),
243 ("ARTIFACTS_1".to_owned(), "g1t-us-1".to_owned()),
244 ("ARTIFACTS_EU".to_owned(), "g1t-eu".to_owned()),
245 ]
246 );
247 assert_eq!(bindings(Some("not json")), bindings(None));
248 }
249
250 #[test]
251 fn keys_name_their_namespace_unless_it_is_the_default() {
252 assert_eq!(split("acme--rocket"), (None, "acme--rocket"));
253 assert_eq!(split("g1t-us-1/acme--rocket"), (Some("g1t-us-1"), "acme--rocket"));
254 assert_eq!(compose(None, "acme--rocket", "g1t"), "acme--rocket");
255 assert_eq!(compose(Some("g1t"), "acme--rocket", "g1t"), "acme--rocket");
256 assert_eq!(compose(Some("g1t-us-1"), "pulls--pul_1", "g1t"), "g1t-us-1/pulls--pul_1");
257 }
258
259 #[test]
260 fn new_repositories_spread_over_the_bound_namespaces() {
Merge branch 'worktree-agent-a2013627e5ea4ab13'261 // A load for each bound namespace, all empty.
262 let bound: Vec<Load> = ["g1t", "g1t-us-1", "g1t-us-2", "g1t-eu"].map(|name| load(name, 0)).to_vec();
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily263 // Nothing configured: everything stays where it was.
Merge branch 'worktree-agent-a2013627e5ea4ab13'264 assert_eq!(Placement::default().choose("rep_1", Residency::Anywhere, &bound), Ok(None));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily265 let placement = Placement::from_vars(Some("g1t-us-1, g1t-us-2,g1t-us-3"), Some("g1t-eu"));
266 let mut seen = std::collections::HashMap::new();
267 for n in 0..1000 {
268 let id = format!("rep_{n:05}");
Merge branch 'worktree-agent-a2013627e5ea4ab13'269 let placed = placement.choose(&id, Residency::Anywhere, &bound).unwrap().unwrap();
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily270 // The same id always lands in the same place.
Merge branch 'worktree-agent-a2013627e5ea4ab13'271 assert_eq!(placement.choose(&id, Residency::Anywhere, &bound).unwrap().unwrap(), placed);
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily272 *seen.entry(placed).or_insert(0) += 1;
273 }
274 // g1t-us-3 is named but not bound, so it takes nothing.
275 assert_eq!(seen.len(), 2);
276 assert!(seen.values().all(|count| *count > 400));
Merge branch 'worktree-agent-a2013627e5ea4ab13'277 assert_eq!(placement.choose("rep_1", Residency::Eu, &bound), Ok(Some("g1t-eu".to_owned())));
278 // EU asked for but not bound: nowhere, so the caller refuses.
279 assert_eq!(placement.choose("rep_1", Residency::Eu, &bound[..3]), Err(Unplaced::NoEu));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily280 assert!(valid_namespace("g1t") && !valid_namespace("-g1t") && !valid_namespace("x"));
281 }
Merge branch 'worktree-agent-a2013627e5ea4ab13'282
283 fn load(namespace: &str, repos: u64) -> Load {
284 Load { namespace: namespace.to_owned(), bound: true, writable: true, repos, ..Load::default() }
285 }
286
287 fn spread(placement: &Placement, loads: &[Load]) -> std::collections::HashMap<String, u32> {
288 let mut seen = std::collections::HashMap::new();
289 for n in 0..1000 {
290 let placed = placement.choose(&format!("rep_{n:05}"), Residency::Anywhere, loads).unwrap().unwrap_or_default();
291 *seen.entry(placed).or_insert(0) += 1;
292 }
293 seen
294 }
295
296 #[test]
297 fn new_repositories_go_to_the_emptier_healthy_namespaces() {
298 let placement = Placement::from_vars(Some("g1t,g1t-us-1,g1t-us-2"), Some("g1t-eu"));
299 // Nothing configured: the default, whatever the loads.
300 assert_eq!(Placement::default().choose("rep_1", Residency::Anywhere, &[load("g1t", 0)]), Ok(None));
301 // Close in size: spread over all three.
302 let even = spread(&placement, &[load("g1t", 5_000), load("g1t-us-1", 5_050), load("g1t-us-2", 5_100)]);
303 assert_eq!(even.len(), 3);
304 assert!(even.values().all(|count| *count > 250));
305 // One far emptier: it takes them all until it catches up.
306 let filling = spread(&placement, &[load("g1t", 9_000), load("g1t-us-1", 200), load("g1t-us-2", 9_000)]);
307 assert_eq!(filling.get("g1t-us-1"), Some(&1000));
308 // Failing, hot, full or read-only: passed over while another can.
309 let mut failing = load("g1t-us-1", 0);
310 failing.failing = true;
311 let mut hot = load("g1t-us-2", 0);
312 hot.peak_per_minute = 9_000;
313 let healthy = spread(&placement, &[load("g1t", 9_000), failing.clone(), hot.clone()]);
314 assert_eq!(healthy.get("g1t"), Some(&1000));
315 let mut full = load("g1t", 10);
316 full.max_repos = Some(10);
317 let mut read_only = load("g1t-us-1", 0);
318 read_only.writable = false;
319 let left = spread(&placement, &[full.clone(), read_only.clone(), load("g1t-us-2", 50)]);
320 assert_eq!(left.get("g1t-us-2"), Some(&1000));
321 // Nothing healthy: still somewhere usable, never a read-only one.
322 let worst = spread(&placement, &[full, read_only.clone(), hot]);
323 assert!(!worst.contains_key("g1t-us-1"));
324 // Nothing usable at all: the default, as before sharding.
325 assert_eq!(placement.choose("rep_1", Residency::Anywhere, &[read_only]), Ok(None));
326 // Not bound: no load given, never chosen.
327 assert_eq!(spread(&placement, &[load("g1t", 10)]).get("g1t"), Some(&1000));
328 }
329
330 #[test]
331 fn eu_residency_goes_to_the_eu_namespace_or_nowhere() {
332 let placement = Placement::from_vars(Some("g1t"), Some("g1t-eu"));
333 let mut eu = load("g1t-eu", 0);
334 assert_eq!(placement.choose("rep_1", Residency::Eu, &[load("g1t", 0), eu.clone()]), Ok(Some("g1t-eu".to_owned())));
335 // Failing is the store's to say when asked; read-only or unbound refuses.
336 eu.failing = true;
337 assert_eq!(placement.choose("rep_1", Residency::Eu, &[eu.clone()]), Ok(Some("g1t-eu".to_owned())));
338 eu.writable = false;
339 assert_eq!(placement.choose("rep_1", Residency::Eu, &[eu]), Err(Unplaced::NoEu));
340 assert_eq!(placement.choose("rep_1", Residency::Eu, &[load("g1t", 0)]), Err(Unplaced::NoEu));
341 assert_eq!(Placement::default().choose("rep_1", Residency::Eu, &[]), Err(Unplaced::NoEu));
342 let bound = vec!["g1t".to_owned(), "g1t-eu".to_owned()];
343 assert!(placement.eu_available(&bound, |_| true));
344 assert!(!placement.eu_available(&bound, |name| name != "g1t-eu"));
345 assert!(!placement.eu_available(&bound[..1], |_| true));
346 assert!(!Placement::default().eu_available(&bound, |_| true));
347 }
348
349 #[test]
350 fn loads_are_read_only_when_there_is_a_choice() {
351 let bound = vec!["g1t".to_owned(), "g1t-us-1".to_owned()];
352 assert!(!Placement::default().needs_loads(&bound));
353 assert!(!Placement::from_vars(Some("g1t-us-1"), None).needs_loads(&bound));
354 assert!(!Placement::from_vars(Some("g1t-us-1,g1t-us-9"), None).needs_loads(&bound));
355 assert!(Placement::from_vars(Some("g1t,g1t-us-1"), None).needs_loads(&bound));
356 assert_eq!(limits(Some(r#"{"g1t":{"max_repos":50000},"x":{}}"#)).get("g1t"), Some(&50_000));
357 assert!(limits(Some("nope")).is_empty() && limits(None).is_empty());
358 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily359}

This file's history is long; its oldest lines are credited to the oldest commit read.