flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/repos/src/refs_cache.rs

404 lines17,139 bytesCodeBlame
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
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(",
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}