g1t/services/repos/src/shards.rs

359 lines16,810 bytesCodeBlame
1//! 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
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).
28
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
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
162/// 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
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);
208 if residency == Residency::Eu {
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 };
213 }
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()))
223 }
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() {
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();
263 // Nothing configured: everything stays where it was.
264 assert_eq!(Placement::default().choose("rep_1", Residency::Anywhere, &bound), Ok(None));
265 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}");
269 let placed = placement.choose(&id, Residency::Anywhere, &bound).unwrap().unwrap();
270 // The same id always lands in the same place.
271 assert_eq!(placement.choose(&id, Residency::Anywhere, &bound).unwrap().unwrap(), placed);
272 *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));
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));
280 assert!(valid_namespace("g1t") && !valid_namespace("-g1t") && !valid_namespace("x"));
281 }
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 }
359}