g1t/services/repos/src/refs_cache.rs
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.
| Mission control shows where you are needed and what agents landed without you; git answers in about 200ms | 1 | //! The answers that list a repository's refs, kept for a moment: the ref |
| 2 | //! advertisement git asks for first on every clone and fetch | |
| 3 | //! (`info/refs?service=git-upload-pack`), and protocol v2's `ls-refs`. | |
| 4 | //! Asking the git store for one takes hundreds of milliseconds; a kept one | |
| 5 | //! is a few. | |
| 6 | //! | |
| 7 | //! An answer is kept under the repository's id and the version of its refs | |
| 8 | //! (`refs_version`, see registry.rs), which goes up after everything g1t | |
| 9 | //! does that changes them, so a change leaves the old answer behind rather | |
| 10 | //! than having to find and remove it. The key also holds the default branch | |
| 11 | //! (the answer's `HEAD` is rewritten to it, see `git_http::with_head`), the | |
| 12 | //! protocol version, and for `ls-refs` the whole request. Answers are kept | |
| 13 | //! in this colo's cache and, sealed, in the cache isolates share | |
| 14 | //! (shared.rs), each for [`TTL_SECONDS`] at most, which bounds how stale | |
| 15 | //! one can be should a change ever fail to move the version. | |
| 16 | //! | |
| 17 | //! Only ever served after the request was authorized, like any answer | |
| 18 | //! from the store: a private repository's answers are kept like any | |
| 19 | //! other's, and read only by whoever may read it. | |
| 20 | ||
| 21 | use crate::git_http::GitRequest; | |
| 22 | use crate::registry::RefsState; | |
| 23 | use crate::shared::Shared; | |
| 24 | use g1t_contracts::repos::GitService; | |
| 25 | use worker::{Headers, Response, Result}; | |
| 26 | ||
| 27 | /// How long an answer is kept, at most. | |
| 28 | pub const TTL_SECONDS: u64 = 60; | |
| 29 | /// Larger answers (repositories with tens of thousands of refs) are not kept. | |
| 30 | const MAX_KEPT_BYTES: usize = 1024 * 1024; | |
| 31 | /// An `ls-refs` request larger than this (a great many ref prefixes) is | |
| 32 | /// passed through. | |
| 33 | const MAX_REQUEST_BYTES: usize = 64 * 1024; | |
| 34 | /// Where answers live in this colo's cache. Not reachable from outside. | |
| 35 | const COLO_CACHE: &str = "https://refs.g1t.internal/"; | |
| 36 | ||
| 37 | /// What a request asks for that can be kept. | |
| 38 | #[derive(Debug, PartialEq, Eq)] | |
| 39 | pub enum Kind { | |
| 40 | /// `GET info/refs?service=git-upload-pack`: the ref advertisement, or | |
| 41 | /// for protocol v2 the capabilities. | |
| 42 | Advertisement, | |
| 43 | /// A protocol v2 `ls-refs` command, by the SHA-256 of the whole request, | |
| 44 | /// so only the same question gets the same answer. | |
| 45 | LsRefs { request: String }, | |
| 46 | } | |
| 47 | ||
| 48 | /// What `git` asks that can be kept, if anything: only fetches, and of | |
| 49 | /// those only the answers that list refs. `body` is a POST's. | |
| 50 | pub fn kind(git: &GitRequest, get: bool, protocol: u8, body: Option<&[u8]>) -> Option<Kind> { | |
| 51 | if git.service != GitService::UploadPack { | |
| 52 | return None; | |
| 53 | } | |
| 54 | match (get, git.endpoint, body) { | |
| 55 | (true, "info/refs", _) => Some(Kind::Advertisement), | |
| 56 | (false, "git-upload-pack", Some(body)) if protocol == 2 && is_ls_refs(body) => Some(Kind::LsRefs { | |
| 57 | request: g1t_secrets::sha256_hex_bytes(body), | |
| 58 | }), | |
| 59 | _ => None, | |
| 60 | } | |
| 61 | } | |
| 62 | ||
| 63 | /// Whether `body` is a whole protocol v2 request whose command is `ls-refs`. | |
| 64 | fn is_ls_refs(body: &[u8]) -> bool { | |
| 65 | if body.len() > MAX_REQUEST_BYTES { | |
| 66 | return false; | |
| 67 | } | |
| 68 | let (lines, end) = crate::land::read_pkt_lines(body); | |
| 69 | end == body.len() | |
| 70 | && lines | |
| 71 | .first() | |
| 72 | .is_some_and(|line| line.strip_suffix(b"\n").unwrap_or(line) == b"command=ls-refs") | |
| 73 | } | |
| 74 | ||
| 75 | /// The protocol version a `Git-Protocol` header asks for: `version=2` | |
| 76 | /// among its colon-separated parameters. 0 without one. | |
| 77 | pub fn protocol(header: Option<&str>) -> u8 { | |
| 78 | header | |
| 79 | .into_iter() | |
| 80 | .flat_map(|value| value.split(':')) | |
| 81 | .filter_map(|parameter| parameter.trim().strip_prefix("version=")) | |
| 82 | .filter_map(|version| version.parse().ok()) | |
| 83 | .next_back() | |
| 84 | .unwrap_or(0) | |
| 85 | } | |
| 86 | ||
| 87 | /// The version to keep answers under, if they may be kept now: not before | |
| 88 | /// the version is known, nor while a credential that could push is out of | |
| 89 | /// g1t's hands. | |
| 90 | pub fn usable(state: Option<RefsState>, now: u64) -> Option<u64> { | |
| 91 | state.filter(|state| now >= state.open_until).map(|state| state.version) | |
| 92 | } | |
| 93 | ||
| 94 | /// Where one answer is kept. | |
| 95 | #[derive(Debug, PartialEq, Eq, Clone)] | |
| 96 | pub struct Key { | |
| 97 | repo_id: String, | |
| 98 | hash: String, | |
| 99 | } | |
| 100 | ||
| 101 | impl Key { | |
| 102 | /// `head` is the default branch `HEAD` is rewritten to, if it is. | |
| 103 | pub fn new(repo_id: &str, version: u64, head: Option<&str>, protocol: u8, kind: &Kind) -> Key { | |
| 104 | let what = match kind { | |
| 105 | Kind::Advertisement => "advertisement".to_owned(), | |
| 106 | Kind::LsRefs { request } => format!("ls-refs {request}"), | |
| 107 | }; | |
| 108 | // Branch names can hold characters a URL cannot, and git forbids | |
| 109 | // newlines in them. | |
| 110 | let head = head.map_or_else(|| "-".to_owned(), |branch| format!("refs/heads/{branch}")); | |
| 111 | Key { | |
| 112 | repo_id: repo_id.to_owned(), | |
| 113 | hash: g1t_secrets::sha256_hex(&format!("{version}\nv{protocol}\n{what}\n{head}")), | |
| 114 | } | |
| 115 | } | |
| 116 | ||
| 117 | fn colo_url(&self) -> String { | |
| 118 | format!("{COLO_CACHE}{}/{}", self.repo_id, self.hash) | |
| 119 | } | |
| 120 | ||
| 121 | fn shared_key(&self) -> String { | |
| 122 | format!("refs:{}:{}", self.repo_id, self.hash) | |
| 123 | } | |
| 124 | } | |
| 125 | ||
| 126 | /// A kept answer. | |
| 127 | #[derive(Debug, PartialEq, Eq, Clone)] | |
| 128 | pub struct Entry { | |
| 129 | pub content_type: String, | |
| 130 | pub body: Vec<u8>, | |
| 131 | } | |
| 132 | ||
| 133 | impl Entry { | |
| 134 | /// Whether it is small enough to keep. | |
| 135 | pub fn keepable(&self) -> bool { | |
| 136 | // Every answer that lists refs ends with a flush packet; one that | |
| 137 | // does not (an error the store sent as a 200) is not kept. | |
| 138 | self.body.len() <= MAX_KEPT_BYTES && self.body.ends_with(b"0000") && !self.content_type.contains('\n') | |
| 139 | } | |
| 140 | ||
| 141 | fn encode(&self) -> Vec<u8> { | |
| 142 | let mut out = self.content_type.as_bytes().to_vec(); | |
| 143 | out.push(b'\n'); | |
| 144 | out.extend_from_slice(&self.body); | |
| 145 | out | |
| 146 | } | |
| 147 | ||
| 148 | fn decode(bytes: &[u8]) -> Option<Entry> { | |
| 149 | let at = bytes.iter().position(|byte| *byte == b'\n')?; | |
| 150 | Some(Entry { | |
| 151 | content_type: String::from_utf8(bytes[..at].to_vec()).ok()?, | |
| 152 | body: bytes[at + 1..].to_vec(), | |
| 153 | }) | |
| 154 | } | |
| 155 | ||
| 156 | /// The answer, as the git store gives it. | |
| 157 | pub fn response(&self) -> Result<Response> { | |
| 158 | let headers = Headers::new(); | |
| 159 | headers.set("content-type", &self.content_type)?; | |
| 160 | headers.set("cache-control", "no-cache")?; | |
| 161 | Ok(Response::from_bytes(self.body.clone())?.with_headers(headers)) | |
| 162 | } | |
| 163 | } | |
| 164 | ||
| 165 | /// Where a kept answer was found, for `Server-Timing`. | |
| 166 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] | |
| 167 | pub enum Found { | |
| 168 | Colo, | |
| 169 | Shared, | |
| 170 | } | |
| 171 | ||
| 172 | impl Found { | |
| 173 | pub fn as_str(self) -> &'static str { | |
| 174 | match self { | |
| 175 | Found::Colo => "hit-colo", | |
| 176 | Found::Shared => "hit-shared", | |
| 177 | } | |
| 178 | } | |
| 179 | } | |
| 180 | ||
| 181 | /// The answer kept under `key`: this colo's first, then the shared one. | |
| 182 | pub async fn get(shared: Option<&Shared>, key: &Key) -> Option<(Entry, Found)> { | |
| 183 | if let Ok(Some(mut response)) = worker::Cache::default().get(key.colo_url(), false).await { | |
| 184 | let content_type = response.headers().get("content-type").ok().flatten(); | |
| 185 | if let (Some(content_type), Ok(body)) = (content_type, response.bytes().await) { | |
| 186 | return Some((Entry { content_type, body }, Found::Colo)); | |
| 187 | } | |
| 188 | } | |
| 189 | let bytes = shared?.get(&key.shared_key()).await?; | |
| 190 | Some((Entry::decode(&bytes)?, Found::Shared)) | |
| 191 | } | |
| 192 | ||
| 193 | /// Keeps `entry` in this colo's cache. A failure only costs a later miss. | |
| 194 | pub async fn keep_in_colo(key: &Key, entry: &Entry) { | |
| 195 | let headers = Headers::new(); | |
| 196 | let _ = headers.set("content-type", &entry.content_type); | |
| 197 | let _ = headers.set("cache-control", &format!("public, max-age={TTL_SECONDS}")); | |
| 198 | let Ok(response) = Response::from_bytes(entry.body.clone()) else { | |
| 199 | return; | |
| 200 | }; | |
| 201 | let _ = worker::Cache::default() | |
| 202 | .put(key.colo_url(), response.with_headers(headers)) | |
| 203 | .await; | |
| 204 | } | |
| 205 | ||
| 206 | /// Keeps `entry` in this colo's cache and the shared one. | |
| 207 | pub async fn keep(shared: Option<&Shared>, key: &Key, entry: &Entry) { | |
| 208 | if !entry.keepable() { | |
| 209 | return; | |
| 210 | } | |
| 211 | let shared_put = async { | |
| 212 | if let Some(shared) = shared { | |
| 213 | shared.put(&key.shared_key(), &entry.encode(), TTL_SECONDS).await; | |
| 214 | } | |
| 215 | }; | |
| 216 | futures_util::future::join(keep_in_colo(key, entry), shared_put).await; | |
| 217 | } | |
| 218 | ||
| 219 | #[cfg(test)] | |
| 220 | mod tests { | |
| 221 | use super::*; | |
| 222 | use g1t_contracts::repos::RepoPath; | |
| 223 | ||
| 224 | fn git(endpoint: &'static str, service: GitService) -> GitRequest { | |
| 225 | GitRequest { | |
| 226 | path: RepoPath { | |
| 227 | namespace: "acme".into(), | |
| 228 | name: "rocket".into(), | |
| 229 | }, | |
| 230 | endpoint, | |
| 231 | service, | |
| 232 | } | |
| 233 | } | |
| 234 | ||
| 235 | fn pkt(payload: &str) -> Vec<u8> { | |
| 236 | format!("{:04x}{payload}", payload.len() + 4).into_bytes() | |
| 237 | } | |
| 238 | ||
| 239 | fn ls_refs(prefix: &str) -> Vec<u8> { | |
| 240 | [ | |
| 241 | pkt("command=ls-refs\n"), | |
| 242 | pkt("agent=git/2.45.0\n"), | |
| 243 | pkt("object-format=sha1\n"), | |
| 244 | b"0001".to_vec(), | |
| 245 | pkt("peel\n"), | |
| 246 | pkt("symrefs\n"), | |
| 247 | pkt(&format!("ref-prefix {prefix}\n")), | |
| 248 | b"0000".to_vec(), | |
| 249 | ] | |
| 250 | .concat() | |
| 251 | } | |
| 252 | ||
| 253 | #[test] | |
| 254 | fn only_the_answers_that_list_refs_are_kept() { | |
| 255 | let refs = git("info/refs", GitService::UploadPack); | |
| 256 | assert_eq!(kind(&refs, true, 0, None), Some(Kind::Advertisement)); | |
| 257 | assert_eq!(kind(&refs, true, 2, None), Some(Kind::Advertisement)); | |
| 258 | // A push's advertisement is never kept: a push needs the refs as | |
| 259 | // they are. | |
| 260 | assert_eq!(kind(&git("info/refs", GitService::ReceivePack), true, 0, None), None); | |
| 261 | let pack = git("git-upload-pack", GitService::UploadPack); | |
| 262 | let listing = ls_refs("refs/heads/"); | |
| 263 | assert!(matches!(kind(&pack, false, 2, Some(&listing)), Some(Kind::LsRefs { .. }))); | |
| 264 | // Only under protocol v2, and never a fetch of objects. | |
| 265 | assert_eq!(kind(&pack, false, 0, Some(&listing)), None); | |
| 266 | let fetch = [pkt("command=fetch\n"), b"0001".to_vec(), pkt("want 1111111111111111111111111111111111111111\n"), pkt("done\n"), b"0000".to_vec()].concat(); | |
| 267 | assert_eq!(kind(&pack, false, 2, Some(&fetch)), None); | |
| 268 | let v0 = [pkt("want 1111111111111111111111111111111111111111 side-band-64k\n"), b"0000".to_vec(), pkt("done\n")].concat(); | |
| 269 | assert_eq!(kind(&pack, false, 0, Some(&v0)), None); | |
| 270 | // Not a whole request, or one too large: passed through. | |
| 271 | assert_eq!(kind(&pack, false, 2, Some(&listing[..listing.len() - 2])), None); | |
| 272 | let huge = [pkt("command=ls-refs\n"), b"0001".to_vec(), (0..3000).flat_map(|n| pkt(&format!("ref-prefix refs/heads/branch-{n}\n"))).collect(), b"0000".to_vec()].concat(); | |
| 273 | assert_eq!(kind(&pack, false, 2, Some(&huge)), None); | |
| 274 | assert_eq!(kind(&pack, false, 2, None), None); | |
| 275 | } | |
| 276 | ||
| 277 | #[test] | |
| 278 | fn the_protocol_version_is_read_from_the_header() { | |
| 279 | assert_eq!(protocol(None), 0); | |
| 280 | assert_eq!(protocol(Some("version=2")), 2); | |
| 281 | assert_eq!(protocol(Some("version=1")), 1); | |
| 282 | assert_eq!(protocol(Some("other=x:version=2")), 2); | |
| 283 | assert_eq!(protocol(Some("version=banana")), 0); | |
| 284 | } | |
| 285 | ||
| 286 | #[test] | |
| 287 | fn nothing_is_kept_without_a_version_or_while_a_push_credential_is_out() { | |
| 288 | assert_eq!(usable(None, 1_000), None); | |
| 289 | assert_eq!(usable(Some(RefsState { version: 4, open_until: 0 }), 1_000), Some(4)); | |
| 290 | assert_eq!(usable(Some(RefsState { version: 4, open_until: 2_000 }), 1_000), None); | |
| 291 | assert_eq!(usable(Some(RefsState { version: 4, open_until: 2_000 }), 2_000), Some(4)); | |
| 292 | } | |
| 293 | ||
| 294 | #[test] | |
| 295 | fn every_part_of_the_question_is_in_the_key() { | |
| 296 | let base = Key::new("rep_1", 3, Some("main"), 0, &Kind::Advertisement); | |
| 297 | assert_eq!(base, Key::new("rep_1", 3, Some("main"), 0, &Kind::Advertisement)); | |
| 298 | // A change to the refs moves the version and leaves the answer behind. | |
| 299 | assert_ne!(base, Key::new("rep_1", 4, Some("main"), 0, &Kind::Advertisement)); | |
| 300 | // So does a new default branch, which HEAD is rewritten to. | |
| 301 | assert_ne!(base, Key::new("rep_1", 3, Some("trunk"), 0, &Kind::Advertisement)); | |
| 302 | assert_ne!(base, Key::new("rep_1", 3, None, 0, &Kind::Advertisement)); | |
| 303 | // v0 and v2 answer differently. | |
| 304 | assert_ne!(base, Key::new("rep_1", 3, Some("main"), 2, &Kind::Advertisement)); | |
| 305 | assert_ne!(base, Key::new("rep_2", 3, Some("main"), 0, &Kind::Advertisement)); | |
| 306 | let heads = Kind::LsRefs { request: g1t_secrets::sha256_hex_bytes(&ls_refs("refs/heads/")) }; | |
| 307 | let tags = Kind::LsRefs { request: g1t_secrets::sha256_hex_bytes(&ls_refs("refs/tags/")) }; | |
| 308 | assert_ne!(Key::new("rep_1", 3, Some("main"), 2, &heads), Key::new("rep_1", 3, Some("main"), 2, &tags)); | |
| 309 | assert_ne!(Key::new("rep_1", 3, Some("main"), 2, &heads), Key::new("rep_1", 3, Some("main"), 2, &Kind::Advertisement)); | |
| 310 | // Odd branch names still make a usable address. | |
| 311 | let odd = Key::new("rep_1", 3, Some("fix/#12 %20"), 0, &Kind::Advertisement); | |
| 312 | assert!(odd.colo_url().starts_with("https://refs.g1t.internal/rep_1/")); | |
| 313 | assert!(!odd.colo_url().contains('#') && !odd.colo_url().contains('%')); | |
| 314 | assert!(odd.shared_key().starts_with("refs:rep_1:")); | |
| 315 | } | |
| 316 | ||
| 317 | /// The functions in `source` (outside its tests) that call any of | |
| 318 | /// `writes`, each with whether it also records the change. | |
| 319 | fn writers(source: &str, writes: &[&str]) -> Vec<(String, bool)> { | |
| 320 | let code = source.split("#[cfg(test)]").next().unwrap_or_default(); | |
| 321 | let lines: Vec<&str> = code.lines().collect(); | |
| 322 | let starts: Vec<usize> = lines | |
| 323 | .iter() | |
| 324 | .enumerate() | |
| 325 | .filter(|(_, line)| { | |
| 326 | let indent = line.len() - line.trim_start().len(); | |
| 327 | let rest = line.trim_start(); | |
| 328 | indent <= 4 | |
| 329 | && ["fn ", "async fn ", "pub fn ", "pub async fn ", "pub(crate) fn ", "pub(crate) async fn "] | |
| 330 | .iter() | |
| 331 | .any(|prefix| rest.starts_with(prefix)) | |
| 332 | }) | |
| 333 | .map(|(at, _)| at) | |
| 334 | .collect(); | |
| 335 | let mut found = Vec::new(); | |
| 336 | for (n, start) in starts.iter().enumerate() { | |
| 337 | let end = starts.get(n + 1).copied().unwrap_or(lines.len()); | |
| 338 | let body = lines[*start..end].join("\n"); | |
| 339 | if writes.iter().any(|write| body.contains(write)) { | |
| 340 | found.push((lines[*start].trim().to_owned(), body.contains("refs_moved("))); | |
| 341 | } | |
| 342 | } | |
| 343 | found | |
| 344 | } | |
| 345 | ||
| 346 | /// Everything that changes a repository's refs moves its version, or a | |
| 347 | /// kept answer would list them as they were. A new way of writing refs | |
| 348 | /// belongs in this list, and its caller must call `refs_moved`. | |
| 349 | #[test] | |
| 350 | fn every_ref_writer_records_the_change() { | |
| 351 | let writes = [ | |
| 352 | "land::push_pack(", | |
| 353 | "land::delete_ref(", | |
| 354 | "land::fast_forward(", | |
| 355 | "mirror::copy(", | |
| 356 | "copy(&theirs, &ours", | |
| 357 | ]; | |
| 358 | let sources = [ | |
| 359 | ("lib.rs", include_str!("lib.rs")), | |
| 360 | ("catch_up.rs", include_str!("catch_up.rs")), | |
| 361 | ("lifecycle.rs", include_str!("lifecycle.rs")), | |
| 362 | ("mirror.rs", include_str!("mirror.rs")), | |
| 363 | ("import.rs", include_str!("import.rs")), | |
| 364 | ("transfer.rs", include_str!("transfer.rs")), | |
| 365 | ("secret_scan.rs", include_str!("secret_scan.rs")), | |
| 366 | ("run_access.rs", include_str!("run_access.rs")), | |
| 367 | ("git_http.rs", include_str!("git_http.rs")), | |
| 368 | ]; | |
| 369 | let mut all = Vec::new(); | |
| 370 | for (file, source) in sources { | |
| 371 | for (function, records) in writers(source, &writes) { | |
| 372 | assert!(records, "{file}: `{function}` changes refs without calling refs_moved"); | |
| 373 | all.push(function); | |
| 374 | } | |
| 375 | } | |
| 376 | // The writers known today, so that the check is seen to find them. | |
| 377 | for expected in ["create", "delete_branch", "land", "update_pull_branch", "mirror", "rename_branch", "forks_follow"] { | |
| 378 | assert!( | |
| 379 | all.iter().any(|function| function.contains(&format!("fn {expected}("))), | |
| 380 | "{expected} not found among {all:?}" | |
| 381 | ); | |
| 382 | } | |
| 383 | // A push through git over HTTPS, and the default branch, which HEAD follows. | |
| 384 | let forwards = writers(include_str!("lib.rs"), &["git_http::forward("]); | |
| 385 | assert!(forwards.iter().any(|(function, _)| function.contains("fn answer_git("))); | |
| 386 | assert!(forwards.iter().all(|(_, records)| *records)); | |
| 387 | let defaults = writers(include_str!("lifecycle.rs"), &["registry.set_default_branch("]); | |
| 388 | assert_eq!(defaults.len(), 3); | |
| 389 | assert!(defaults.iter().all(|(_, records)| *records)); | |
| 390 | } | |
| 391 | ||
| 392 | #[test] | |
| 393 | fn an_entry_survives_being_kept() { | |
| 394 | let entry = Entry { | |
| 395 | content_type: "application/x-git-upload-pack-advertisement".into(), | |
| 396 | body: b"001e# service=git-upload-pack\n0000".to_vec(), | |
| 397 | }; | |
| 398 | assert_eq!(Entry::decode(&entry.encode()), Some(entry.clone())); | |
| 399 | assert!(entry.keepable()); | |
| 400 | let large = Entry { body: vec![b'0'; MAX_KEPT_BYTES + 1], ..entry }; | |
| 401 | assert!(!large.keepable()); | |
| 402 | assert_eq!(Entry::decode(b"no newline"), None); | |
| 403 | } | |
| 404 | } |