g1t/services/repos/src/refs_cache.rs
| 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 | } |