g1t/services/repos/src/refs_cache.rs

406 lines17,245 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.

Mission control shows where you are needed and what agents landed without you; git answers in about 200ms1//! 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
21use crate::git_http::GitRequest;
22use crate::registry::RefsState;
23use crate::shared::Shared;
24use g1t_contracts::repos::GitService;
25use worker::{Headers, Response, Result};
26
27/// How long an answer is kept, at most.
28pub const TTL_SECONDS: u64 = 60;
29/// Larger answers (repositories with tens of thousands of refs) are not kept.
30const MAX_KEPT_BYTES: usize = 1024 * 1024;
31/// An `ls-refs` request larger than this (a great many ref prefixes) is
32/// passed through.
33const MAX_REQUEST_BYTES: usize = 64 * 1024;
34/// Where answers live in this colo's cache. Not reachable from outside.
35const COLO_CACHE: &str = "https://refs.g1t.internal/";
36
37/// What a request asks for that can be kept.
38#[derive(Debug, PartialEq, Eq)]
39pub 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.
50pub 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`.
64fn 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.
77pub 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.
90pub 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)]
96pub struct Key {
97 repo_id: String,
98 hash: String,
99}
100
101impl 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)]
128pub struct Entry {
129 pub content_type: String,
130 pub body: Vec<u8>,
131}
132
133impl 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)]
167pub enum Found {
168 Colo,
169 Shared,
170}
171
172impl 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.
182pub 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.
194pub 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.
207pub 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)]
220mod 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(",
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily355 "land::push_ref(",
Mission control shows where you are needed and what agents landed without you; git answers in about 200ms356 "mirror::copy(",
357 "copy(&theirs, &ours",
358 ];
359 let sources = [
360 ("lib.rs", include_str!("lib.rs")),
361 ("catch_up.rs", include_str!("catch_up.rs")),
362 ("lifecycle.rs", include_str!("lifecycle.rs")),
363 ("mirror.rs", include_str!("mirror.rs")),
364 ("import.rs", include_str!("import.rs")),
365 ("transfer.rs", include_str!("transfer.rs")),
366 ("secret_scan.rs", include_str!("secret_scan.rs")),
367 ("run_access.rs", include_str!("run_access.rs")),
368 ("git_http.rs", include_str!("git_http.rs")),
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily369 ("forks.rs", include_str!("forks.rs")),
Mission control shows where you are needed and what agents landed without you; git answers in about 200ms370 ];
371 let mut all = Vec::new();
372 for (file, source) in sources {
373 for (function, records) in writers(source, &writes) {
374 assert!(records, "{file}: `{function}` changes refs without calling refs_moved");
375 all.push(function);
376 }
377 }
378 // The writers known today, so that the check is seen to find them.
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily379 for expected in ["create", "delete_branch", "land", "update_pull_branch", "mirror", "rename_branch", "forks_follow", "keep_head", "revive"] {
Mission control shows where you are needed and what agents landed without you; git answers in about 200ms380 assert!(
381 all.iter().any(|function| function.contains(&format!("fn {expected}("))),
382 "{expected} not found among {all:?}"
383 );
384 }
385 // A push through git over HTTPS, and the default branch, which HEAD follows.
386 let forwards = writers(include_str!("lib.rs"), &["git_http::forward("]);
387 assert!(forwards.iter().any(|(function, _)| function.contains("fn answer_git(")));
388 assert!(forwards.iter().all(|(_, records)| *records));
389 let defaults = writers(include_str!("lifecycle.rs"), &["registry.set_default_branch("]);
390 assert_eq!(defaults.len(), 3);
391 assert!(defaults.iter().all(|(_, records)| *records));
392 }
393
394 #[test]
395 fn an_entry_survives_being_kept() {
396 let entry = Entry {
397 content_type: "application/x-git-upload-pack-advertisement".into(),
398 body: b"001e# service=git-upload-pack\n0000".to_vec(),
399 };
400 assert_eq!(Entry::decode(&entry.encode()), Some(entry.clone()));
401 assert!(entry.keepable());
402 let large = Entry { body: vec![b'0'; MAX_KEPT_BYTES + 1], ..entry };
403 assert!(!large.keepable());
404 assert_eq!(Entry::decode(b"no newline"), None);
405 }
406}