g1t/services/repos/src/pack_cache.rs

1,095 lines43,912 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.

Keep fresh clones' packs in R2, so a repeat clone of the same commit skips the git store1//! Packs for fresh clones, kept for the next clone of the same commit.
2//!
3//! An upload-pack request that wants objects and has none (`have` lines)
4//! is a clone: a new checkout, or a sandbox's shallow `deepen 1` clone of
5//! a pull request's head. The git store builds the same pack for it every
6//! time, which takes seconds and is an operation. The answer is kept in a
7//! bucket ([`PackStore`], R2's `GIT_PACKS` on Cloudflare) under the
8//! repository's id, the version of its refs (`refs_version`, see
9//! registry.rs and refs_cache.rs) and a hash of the request with what does
10//! not change the answer taken out ([`normalize`]): the client's `agent`,
11//! its session id, and the order of its lines. A change to the refs moves
12//! the version, so a pack is never served across one; old ones are left for
13//! the bucket's lifecycle rule to delete.
14//!
15//! On a miss the store's answer goes to git as it arrives and, at the same
16//! time, to the bucket ([`Tee`] and [`fill`]): never held whole. Packs over
17//! [`MAX_PACK_BYTES`] pass through. What is kept is only ever a whole pack:
18//! a small one is written in one `put` once it has all arrived, a larger
19//! one in multipart parts that become an object only when the last has
20//! been checked ([`PackCheck`]); a fill cut short is aborted and leaves
21//! nothing to serve.
22//!
23//! Only ever served after the request was authorized, like any answer from
24//! the store: a private repository's packs are read only by whoever may
25//! read it. Without the bucket binding (self-hosted, or before it exists)
26//! nothing is kept and every clone goes to the store, as before.
27
28use std::cell::RefCell;
29use std::collections::{BTreeSet, HashSet, VecDeque};
30use std::pin::Pin;
31use std::rc::Rc;
32use std::task::{Context, Poll, Waker};
33
34use futures_util::{Stream, StreamExt};
35use g1t_contracts::repos::GitService;
36use worker::{Bucket, Env, Headers, Response, ResponseBody, Result, UploadedPart};
37
38use crate::git_http::GitRequest;
39
40/// An error worth seeing in the Worker's logs; on stderr in tests, where
41/// there is no console to call.
42macro_rules! log {
43 ($($arg:tt)*) => {{
44 #[cfg(target_arch = "wasm32")]
45 worker::console_error!($($arg)*);
46 #[cfg(not(target_arch = "wasm32"))]
47 eprintln!($($arg)*);
48 }};
49}
50
51/// The meters (meters.rs): a clone answered from the bucket, and one the
52/// store was asked for. Neither is an operation unless `operation_mapping`
53/// says so; by default they are not.
54pub const HIT: &str = "pack_cache.hit";
55pub const MISS: &str = "pack_cache.miss";
56
57/// Packs larger than this pass through without being kept.
58pub const MAX_PACK_BYTES: u64 = 200 * 1024 * 1024;
59/// Requests larger than this (a great many wants) are not kept.
60const MAX_REQUEST_BYTES: usize = 1024 * 1024;
61/// A multipart part: R2's smallest, so a fill holds as little as it can.
62/// Every part but the last is this size, as R2 requires.
63const PART_BYTES: usize = 5 * 1024 * 1024;
64/// What may wait between git's stream and the bucket. A bucket slower than
65/// git ends the fill rather than holding more.
66const MAX_QUEUED_BYTES: usize = 5 * 1024 * 1024;
67/// Fills at once in one isolate; more clones pass through.
68const MAX_FILLS: usize = 2;
69pub const CONTENT_TYPE: &str = "application/x-git-upload-pack-result";
70
71/// Where packs are kept: the R2 bucket on Cloudflare. A small port, so a
72/// self-hosted installation can put another store behind it.
73#[allow(async_fn_in_trait)]
74pub trait PackStore {
75 /// The object's size and body, if it is there.
76 async fn get(&self, key: &str) -> Result<Option<(u64, ResponseBody)>>;
77 async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()>;
78 /// Starts a multipart upload to `key`; its id.
79 async fn begin(&self, key: &str) -> Result<String>;
80 /// Uploads part `number` (from 1); its etag.
81 async fn part(&self, key: &str, upload: &str, number: u16, bytes: Vec<u8>) -> Result<String>;
82 /// Makes the object from the parts, which until now nobody can read.
83 async fn complete(&self, key: &str, upload: &str, parts: Vec<(u16, String)>) -> Result<()>;
84 async fn abort(&self, key: &str, upload: &str) -> Result<()>;
85}
86
87/// The R2 adapter: the `GIT_PACKS` bucket binding.
88pub struct R2Packs {
89 bucket: Bucket,
90}
91
92impl R2Packs {
93 /// `None` without the binding: nothing is kept.
94 pub fn from_env(env: &Env) -> Option<R2Packs> {
95 env.bucket("GIT_PACKS").ok().map(|bucket| R2Packs { bucket })
96 }
97}
98
99impl PackStore for R2Packs {
100 async fn get(&self, key: &str) -> Result<Option<(u64, ResponseBody)>> {
101 let Some(object) = self.bucket.get(key).execute().await? else {
102 return Ok(None);
103 };
104 let size = object.size();
105 let Some(body) = object.body() else {
106 return Ok(None);
107 };
108 Ok(Some((size, body.response_body()?)))
109 }
110
111 async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()> {
112 self.bucket.put(key, bytes).execute().await?;
113 Ok(())
114 }
115
116 async fn begin(&self, key: &str) -> Result<String> {
117 let upload = self.bucket.create_multipart_upload(key).execute().await?;
118 Ok(upload.upload_id().await)
119 }
120
121 async fn part(&self, key: &str, upload: &str, number: u16, bytes: Vec<u8>) -> Result<String> {
122 let upload = self.bucket.resume_multipart_upload(key, upload)?;
123 Ok(upload.upload_part(number, bytes).await?.etag())
124 }
125
126 async fn complete(&self, key: &str, upload: &str, parts: Vec<(u16, String)>) -> Result<()> {
127 let upload = self.bucket.resume_multipart_upload(key, upload)?;
128 upload
129 .complete(parts.into_iter().map(|(number, etag)| UploadedPart::new(number, etag)))
130 .await?;
131 Ok(())
132 }
133
134 async fn abort(&self, key: &str, upload: &str) -> Result<()> {
135 self.bucket.resume_multipart_upload(key, upload)?.abort().await
136 }
137}
138
139/// The pkt-line payloads of `body`, with flush (`0000`) and delimiter
140/// (`0001`) packets as `None`; `None` if it is not one whole pkt-line
141/// stream.
142fn lines(body: &[u8]) -> Option<Vec<Option<&[u8]>>> {
143 let mut out = Vec::new();
144 let mut position = 0;
145 while position < body.len() {
146 let header = body.get(position..position + 4)?;
147 let length = usize::from_str_radix(std::str::from_utf8(header).ok()?, 16).ok()?;
148 match length {
149 0 | 1 => {
150 out.push(None);
151 position += 4;
152 }
153 2 | 3 => return None,
154 _ => {
155 let payload = body.get(position + 4..position + length)?;
156 out.push(Some(payload.strip_suffix(b"\n").unwrap_or(payload)));
157 position += length;
158 }
159 }
160 }
161 Some(out)
162}
163
164/// Capabilities that say who is asking, not what: left out of the key.
165fn says_who(capability: &str) -> bool {
166 capability.starts_with("agent=") || capability.starts_with("session-id=")
167}
168
169/// A fetch request in a form that is the same for every request whose
170/// answer is the same, or `None` if its answer is not kept: anything but a
171/// fetch of objects, one that sends `have`s (negotiation: what the client
172/// has decides the pack) or `shallow`s (it has commits already), one with
173/// no wants, or one that is not a whole request. `protocol` is the
174/// `Git-Protocol` version (refs_cache::protocol).
175pub fn normalize(protocol: u8, body: &[u8]) -> Option<String> {
176 if body.len() > MAX_REQUEST_BYTES {
177 return None;
178 }
179 let lines = lines(body)?;
180 if protocol == 2 { normalize_v2(&lines) } else { normalize_v0(&lines) }
181}
182
183/// Protocol v2: `command=fetch`, capabilities, a delimiter, then the
184/// arguments and a flush. Neither section's order matters.
185fn normalize_v2(lines: &[Option<&[u8]>]) -> Option<String> {
186 let text = |line: &[u8]| std::str::from_utf8(line).ok().map(str::to_owned);
187 let mut lines = lines.iter();
188 if *lines.next()? != Some(b"command=fetch".as_slice()) {
189 return None;
190 }
191 let mut capabilities = BTreeSet::new();
192 let mut arguments = BTreeSet::new();
193 let mut in_arguments = false;
194 let mut ended = false;
195 for line in lines.by_ref() {
196 match line {
197 None if !in_arguments => in_arguments = true,
198 None => {
199 ended = true;
200 break;
201 }
202 Some(line) => {
203 let line = text(line)?;
204 if !in_arguments {
205 if !says_who(&line) {
206 capabilities.insert(line);
207 }
208 } else if line.starts_with("have ") || line.starts_with("shallow ") {
209 return None;
210 } else {
211 arguments.insert(line);
212 }
213 }
214 }
215 }
216 // A flush ends the request, and nothing follows it.
217 if !ended || lines.next().is_some() {
218 return None;
219 }
220 if !arguments.iter().any(|line| line.starts_with("want ") || line.starts_with("want-ref ")) {
221 return None;
222 }
223 let capabilities = capabilities.into_iter().collect::<Vec<_>>().join("\n");
224 let arguments = arguments.into_iter().collect::<Vec<_>>().join("\n");
225 Some(format!("v2\n{capabilities}\n--\n{arguments}"))
226}
227
228/// Protocol v0 and v1: wants (the first with the capabilities after its
229/// object id), `shallow`, `deepen`, `filter` lines, a flush, then `done`.
230fn normalize_v0(lines: &[Option<&[u8]>]) -> Option<String> {
231 let mut capabilities = BTreeSet::new();
232 let mut requests = BTreeSet::new();
233 let mut lines = lines.iter();
234 let mut flushed = false;
235 for line in lines.by_ref() {
236 let Some(line) = line else {
237 flushed = true;
238 break;
239 };
240 let line = std::str::from_utf8(line).ok()?;
241 if let Some(rest) = line.strip_prefix("want ") {
242 let (oid, rest) = rest.split_once(' ').unwrap_or((rest, ""));
243 capabilities.extend(rest.split(' ').filter(|c| !c.is_empty() && !says_who(c)).map(str::to_owned));
244 requests.insert(format!("want {oid}"));
245 } else if line.starts_with("shallow ") || line.starts_with("have ") {
246 return None;
247 } else {
248 requests.insert(line.to_owned());
249 }
250 }
251 // Then only `done`: a `have` is negotiation, and without `done` the
252 // answer is not a pack.
253 if !flushed || lines.as_slice() != [Some(b"done".as_slice())] {
254 return None;
255 }
256 if !requests.iter().any(|line| line.starts_with("want ")) {
257 return None;
258 }
259 let capabilities = capabilities.into_iter().collect::<Vec<_>>().join(" ");
260 let requests = requests.into_iter().collect::<Vec<_>>().join("\n");
261 Some(format!("v0\n{capabilities}\n--\n{requests}\ndone"))
262}
263
264/// The normalized request, if `git`'s answer may be kept: a POST to
265/// upload-pack with a body read whole and not compressed, that
266/// [`normalize`] accepts.
267pub fn cacheable(git: &GitRequest, get: bool, protocol: u8, encoding: Option<&str>, body: Option<&[u8]>) -> Option<String> {
268 if git.service != GitService::UploadPack || git.endpoint != "git-upload-pack" || get {
269 return None;
270 }
271 if encoding.is_some_and(|encoding| !encoding.trim().is_empty() && !encoding.trim().eq_ignore_ascii_case("identity")) {
272 return None;
273 }
274 normalize(protocol, body?)
275}
276
277/// Where one pack is kept: by repository, refs version and request.
278#[derive(Clone, Debug, PartialEq, Eq)]
279pub struct Key(String);
280
281impl Key {
282 pub fn new(repo_id: &str, refs_version: u64, normalized: &str) -> Key {
283 Key(format!("packs/{repo_id}/{refs_version}/{}", g1t_secrets::sha256_hex(normalized)))
284 }
285
286 pub fn as_str(&self) -> &str {
287 &self.0
288 }
289}
290
291/// A kept pack.
292pub struct Kept {
293 pub size: u64,
294 pub body: ResponseBody,
295}
296
297impl Kept {
298 /// The answer git gets, as the store would give it.
299 pub fn response(self) -> Result<Response> {
300 let headers = Headers::new();
301 headers.set("content-type", CONTENT_TYPE)?;
302 headers.set("cache-control", "no-cache")?;
303 Ok(Response::from_body(self.body)?.with_headers(headers))
304 }
305}
306
307/// The pack kept under `key`. A failure to read is a miss.
308pub async fn get<S: PackStore>(store: &S, key: &Key) -> Option<Kept> {
309 match store.get(key.as_str()).await {
310 Ok(found) => found.map(|(size, body)| Kept { size, body }),
311 Err(error) => {
312 log!("pack cache not read: {error}");
313 None
314 }
315 }
316}
317
318/// Checks, as it streams past, that an upload-pack answer is one whole
319/// pack: every packet well formed, a pack in side-band channel 1, no
320/// error (`ERR`, or channel 3), and a flush at the end.
321#[derive(Debug, Default)]
322pub struct PackCheck {
323 header: Vec<u8>,
324 /// Payload bytes still to come in the current packet.
325 remaining: usize,
326 /// The first bytes of the current packet's payload, until judged.
327 start: Vec<u8>,
328 judged: bool,
329 saw_pack: bool,
330 failed: bool,
331 last_flush: bool,
332}
333
334/// How much of a packet's start says what it is: `\x01PACK`.
335const START: usize = 5;
336
337impl PackCheck {
338 pub fn feed(&mut self, mut bytes: &[u8]) {
339 while !bytes.is_empty() && !self.failed {
340 if self.remaining == 0 {
341 let take = (4 - self.header.len()).min(bytes.len());
342 self.header.extend_from_slice(&bytes[..take]);
343 bytes = &bytes[take..];
344 if self.header.len() < 4 {
345 return;
346 }
347 let length = std::str::from_utf8(&self.header).ok().and_then(|hex| usize::from_str_radix(hex, 16).ok());
348 self.header.clear();
349 self.last_flush = length == Some(0);
350 match length {
351 None | Some(3) => self.failed = true,
352 Some(0..=2 | 4) => {}
353 Some(length) => {
354 self.remaining = length - 4;
355 self.start.clear();
356 self.judged = false;
357 }
358 }
359 continue;
360 }
361 let take = self.remaining.min(bytes.len());
362 if !self.judged {
363 let room = (START - self.start.len()).min(take);
364 self.start.extend_from_slice(&bytes[..room]);
365 }
366 self.remaining -= take;
367 bytes = &bytes[take..];
368 if !self.judged && (self.start.len() >= START || self.remaining == 0) {
369 self.judge();
370 }
371 }
372 }
373
374 /// Looks at the start of a packet, once.
375 fn judge(&mut self) {
376 self.judged = true;
377 let start = self.start.as_slice();
378 if start.first() == Some(&3) || start.starts_with(b"ERR ") || start.starts_with(b"\x01ERR ") {
379 self.failed = true;
380 }
381 if start.starts_with(b"\x01PACK") {
382 self.saw_pack = true;
383 }
384 }
385
386 pub fn failed(&self) -> bool {
387 self.failed
388 }
389
390 pub fn complete(&self) -> bool {
391 !self.failed && self.saw_pack && self.last_flush && self.remaining == 0 && self.header.is_empty()
392 }
393}
394
395/// What passes from git's stream to the fill.
396#[derive(Debug, PartialEq, Eq)]
397enum Piece {
398 Bytes(Vec<u8>),
399 /// The store's answer ended as it should.
400 End,
401}
402
403/// The queue between a [`Tee`] and its [`fill`]. Closed without an
404/// [`Piece::End`], the fill is abandoned.
405#[derive(Default)]
406struct Pipe {
407 pieces: VecDeque<Piece>,
408 queued: usize,
409 closed: bool,
410 waker: Option<Waker>,
411}
412
413type SharedPipe = Rc<RefCell<Pipe>>;
414
415impl Pipe {
416 fn close(pipe: &SharedPipe) {
417 let mut pipe = pipe.borrow_mut();
418 pipe.closed = true;
419 if let Some(waker) = pipe.waker.take() {
420 waker.wake();
421 }
422 }
423}
424
425/// The receiving end of a [`Pipe`], as a stream.
426struct Drain(SharedPipe);
427
428impl Stream for Drain {
429 type Item = Piece;
430
431 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Piece>> {
432 let mut pipe = self.0.borrow_mut();
433 if let Some(piece) = pipe.pieces.pop_front() {
434 if let Piece::Bytes(bytes) = &piece {
435 pipe.queued -= bytes.len();
436 }
437 return Poll::Ready(Some(piece));
438 }
439 if pipe.closed {
440 return Poll::Ready(None);
441 }
442 pipe.waker = Some(cx.waker().clone());
443 Poll::Pending
444 }
445}
446
447type ByteChunks = Pin<Box<dyn Stream<Item = Result<Vec<u8>>>>>;
448
449/// The store's answer on its way to git, copied to a fill while one is
450/// wanted, and measured. The fill is let go, and the answer streams on, if
451/// the pack grows past [`MAX_PACK_BYTES`] or the bucket falls behind.
452pub struct Tee {
453 inner: ByteChunks,
454 pipe: Option<SharedPipe>,
455 seen: u64,
456 done: bool,
457 /// Told the bytes that went to git, once, however the stream ends.
458 on_end: Option<Box<dyn FnOnce(u64)>>,
459}
460
461impl Tee {
462 fn new(inner: ByteChunks, pipe: Option<SharedPipe>, on_end: Box<dyn FnOnce(u64)>) -> Tee {
463 Tee { inner, pipe, seen: 0, done: false, on_end: Some(on_end) }
464 }
465
466 fn let_go(&mut self) {
467 if let Some(pipe) = self.pipe.take() {
468 Pipe::close(&pipe);
469 }
470 }
471
472 fn ended(&mut self) {
473 self.done = true;
474 if let Some(pipe) = self.pipe.take() {
475 let mut held = pipe.borrow_mut();
476 held.pieces.push_back(Piece::End);
477 drop(held);
478 Pipe::close(&pipe);
479 }
480 if let Some(on_end) = self.on_end.take() {
481 on_end(self.seen);
482 }
483 }
484}
485
486impl Stream for Tee {
487 type Item = Result<Vec<u8>>;
488
489 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
490 match self.inner.as_mut().poll_next(cx) {
491 Poll::Ready(Some(Ok(chunk))) => {
492 self.seen += chunk.len() as u64;
493 if self.seen > MAX_PACK_BYTES {
494 self.let_go();
495 }
496 if let Some(pipe) = &self.pipe {
497 let mut held = pipe.borrow_mut();
498 if held.queued + chunk.len() > MAX_QUEUED_BYTES {
499 drop(held);
500 self.let_go();
501 } else {
502 held.queued += chunk.len();
503 held.pieces.push_back(Piece::Bytes(chunk.clone()));
504 if let Some(waker) = held.waker.take() {
505 waker.wake();
506 }
507 }
508 }
509 Poll::Ready(Some(Ok(chunk)))
510 }
511 Poll::Ready(Some(Err(error))) => {
512 self.let_go();
513 Poll::Ready(Some(Err(error)))
514 }
515 Poll::Ready(None) => {
516 if !self.done {
517 self.ended();
518 }
519 Poll::Ready(None)
520 }
521 Poll::Pending => Poll::Pending,
522 }
523 }
524}
525
526impl Drop for Tee {
527 fn drop(&mut self) {
528 // Git went away, or the answer failed, before it ended: whatever the
529 // fill has is not a whole pack.
530 self.let_go();
531 if let Some(on_end) = self.on_end.take() {
532 on_end(self.seen);
533 }
534 }
535}
536
537/// How a fill went.
538#[derive(Debug, PartialEq, Eq)]
539pub enum Filled {
540 Kept { bytes: u64 },
541 /// Let go: too large, the bucket behind, or the answer cut short.
542 Abandoned,
543 /// Not a whole pack (an error the store sent as a 200).
544 NotAPack,
545 /// The bucket refused a write.
546 Failed,
547}
548
549/// Writes what comes down the pipe to `key`, and only a whole pack: in one
550/// put when it is small, in parts otherwise, completed only once the last
551/// part has arrived and the pack has been checked. Anything else is
552/// aborted, and nothing is left to be read.
553async fn fill<S: PackStore>(store: &S, key: &str, mut pieces: impl Stream<Item = Piece> + Unpin) -> Filled {
554 let mut check = PackCheck::default();
555 let mut buffer: Vec<u8> = Vec::new();
556 let mut upload: Option<(String, Vec<(u16, String)>)> = None;
557 let mut total = 0u64;
558 let outcome = 'pieces: loop {
559 match pieces.next().await {
560 Some(Piece::Bytes(bytes)) => {
561 check.feed(&bytes);
562 if check.failed() {
563 break Filled::NotAPack;
564 }
565 total += bytes.len() as u64;
566 buffer.extend_from_slice(&bytes);
567 // Strictly more than a part, so the last part is never empty.
568 while buffer.len() > PART_BYTES {
569 let rest = buffer.split_off(PART_BYTES);
570 let part = std::mem::replace(&mut buffer, rest);
571 if upload.is_none() {
572 match store.begin(key).await {
573 Ok(id) => upload = Some((id, Vec::new())),
574 Err(error) => {
575 log!("pack cache upload not started: {error}");
576 break 'pieces Filled::Failed;
577 }
578 }
579 }
580 let Some((id, parts)) = upload.as_mut() else { break 'pieces Filled::Failed };
581 let number = parts.len() as u16 + 1;
582 match store.part(key, id, number, part).await {
583 Ok(etag) => parts.push((number, etag)),
584 Err(error) => {
585 log!("pack cache part not uploaded: {error}");
586 break 'pieces Filled::Failed;
587 }
588 }
589 }
590 }
591 Some(Piece::End) if check.complete() => {
592 let last = std::mem::take(&mut buffer);
593 let finished = match upload.take() {
594 None => store.put(key, last).await,
595 Some((id, mut parts)) => {
596 let number = parts.len() as u16 + 1;
597 let done = match store.part(key, &id, number, last).await {
598 Ok(etag) => {
599 parts.push((number, etag));
600 store.complete(key, &id, parts).await
601 }
602 Err(error) => Err(error),
603 };
604 if done.is_err() {
605 let _ = store.abort(key, &id).await;
606 }
607 done
608 }
609 };
610 break match finished {
611 Ok(()) => Filled::Kept { bytes: total },
612 Err(error) => {
613 log!("pack not kept: {error}");
614 Filled::Failed
615 }
616 };
617 }
618 Some(Piece::End) => break Filled::NotAPack,
619 None => break Filled::Abandoned,
620 }
621 };
622 if let Some((id, _)) = upload {
623 // Never completed: an upload nobody can read, aborted now, or by
624 // the bucket's lifecycle rule should this fail too.
625 let _ = store.abort(key, &id).await;
626 }
627 outcome
628}
629
630thread_local! {
631 /// Keys being filled in this isolate: one fill each, at most [`MAX_FILLS`].
632 static FILLING: RefCell<HashSet<String>> = RefCell::new(HashSet::new());
633}
634
635/// Holds a key's place among the fills, until dropped.
636struct Filling(String);
637
638impl Filling {
639 fn take(key: &str) -> Option<Filling> {
640 FILLING.with(|filling| {
641 let mut filling = filling.borrow_mut();
642 (filling.len() < MAX_FILLS && filling.insert(key.to_owned())).then(|| Filling(key.to_owned()))
643 })
644 }
645}
646
647impl Drop for Filling {
648 fn drop(&mut self) {
649 FILLING.with(|filling| filling.borrow_mut().remove(&self.0));
650 }
651}
652
653/// Whether `response`, the store's answer to a cacheable request, may be
654/// kept: a 200 of the right type, not known to be too large.
655pub fn keepable(response: &Response) -> bool {
656 let headers = response.headers();
657 let typed = headers
658 .get("content-type")
659 .ok()
660 .flatten()
661 .is_some_and(|value| value.trim().eq_ignore_ascii_case(CONTENT_TYPE));
662 let length = headers.get("content-length").ok().flatten().and_then(|value| value.parse::<u64>().ok());
663 response.status_code() == 200 && typed && length.is_none_or(|length| length <= MAX_PACK_BYTES)
664}
665
666/// A fill under way, for `ctx.wait_until`.
667pub type Fill = Pin<Box<dyn std::future::Future<Output = Filled>>>;
668
669/// The store's answer, streamed to git and, when it may be kept, at once
670/// to `store` under `key`; the fill to wait on, if there is one. `on_end`
671/// is told the bytes that went to git.
672pub fn tee<S: PackStore + 'static>(
673 mut response: Response,
674 store: Rc<S>,
675 key: &Key,
676 on_end: Box<dyn FnOnce(u64)>,
677) -> Result<(Response, Option<Fill>)> {
678 let place = keepable(&response).then(|| Filling::take(key.as_str())).flatten();
679 let status = response.status_code();
680 let headers = response.headers().clone();
681 headers.delete("content-length")?;
682 let inner: ByteChunks = Box::pin(response.stream()?);
683 let pipe = place.as_ref().map(|_| SharedPipe::default());
684 let tee = Tee::new(inner, pipe.clone(), on_end);
685 let filling = match (place, pipe) {
686 (Some(place), Some(pipe)) => {
687 let key = key.as_str().to_owned();
688 let future: Fill = Box::pin(async move {
689 let _place = place;
690 fill(store.as_ref(), &key, Drain(pipe)).await
691 });
692 Some(future)
693 }
694 _ => None,
695 };
696 Ok((Response::from_stream(tee)?.with_headers(headers).with_status(status), filling))
697}
698
699#[cfg(test)]
700mod tests {
701 use super::*;
702 use std::cell::Cell;
703 use std::collections::HashMap;
704 use std::future::Future;
705
706 fn pkt(payload: &str) -> Vec<u8> {
707 format!("{:04x}{payload}", payload.len() + 4).into_bytes()
708 }
709
710 fn pkt_bytes(payload: &[u8]) -> Vec<u8> {
711 let mut out = format!("{:04x}", payload.len() + 4).into_bytes();
712 out.extend_from_slice(payload);
713 out
714 }
715
716 const A: &str = "1111111111111111111111111111111111111111";
717 const B: &str = "2222222222222222222222222222222222222222";
718
719 fn v2(capabilities: &[&str], arguments: &[&str]) -> Vec<u8> {
720 let mut body = pkt("command=fetch\n");
721 for line in capabilities {
722 body.extend(pkt(&format!("{line}\n")));
723 }
724 body.extend(b"0001");
725 for line in arguments {
726 body.extend(pkt(&format!("{line}\n")));
727 }
728 body.extend(b"0000");
729 body
730 }
731
732 fn v0(lines: &[&str], after: &[&str]) -> Vec<u8> {
733 let mut body = Vec::new();
734 for line in lines {
735 body.extend(pkt(&format!("{line}\n")));
736 }
737 body.extend(b"0000");
738 for line in after {
739 body.extend(pkt(&format!("{line}\n")));
740 }
741 body
742 }
743
744 #[test]
745 fn a_v2_clone_is_the_same_whatever_the_order_or_the_agent() {
746 let one = v2(
747 &["agent=git/2.45.0", "object-format=sha1"],
748 &["thin-pack", "ofs-delta", "deepen 1", &format!("want {A}"), &format!("want {B}"), "filter blob:none", "done"],
749 );
750 let two = v2(
751 &["object-format=sha1", "agent=git/2.47.1", "session-id=abc"],
752 &["done", "filter blob:none", &format!("want {B}"), "deepen 1", "ofs-delta", &format!("want {A}"), "thin-pack", &format!("want {A}")],
753 );
754 let first = normalize(2, &one).unwrap();
755 assert_eq!(first, normalize(2, &two).unwrap());
756 assert!(!first.contains("agent="));
757 // Each of deepen, filter and the wants changes the answer.
758 let deeper = v2(&["object-format=sha1"], &["thin-pack", "ofs-delta", "deepen 2", &format!("want {A}"), &format!("want {B}"), "filter blob:none", "done"]);
759 let unfiltered = v2(&["object-format=sha1"], &["thin-pack", "ofs-delta", "deepen 1", &format!("want {A}"), &format!("want {B}"), "done"]);
760 let fewer = v2(&["object-format=sha1"], &["thin-pack", "ofs-delta", "deepen 1", &format!("want {A}"), "filter blob:none", "done"]);
761 for other in [deeper, unfiltered, fewer] {
762 assert_ne!(normalize(2, &other).unwrap(), first);
763 }
764 // Progress is part of the answer.
765 let quiet = v2(&["object-format=sha1"], &["no-progress", "thin-pack", "ofs-delta", "deepen 1", &format!("want {A}"), &format!("want {B}"), "filter blob:none", "done"]);
766 assert_ne!(normalize(2, &quiet).unwrap(), first);
767 }
768
769 #[test]
770 fn a_v0_clone_is_the_same_whatever_the_order_or_the_agent() {
771 let one = v0(
772 &[&format!("want {A} multi_ack_detailed side-band-64k thin-pack ofs-delta agent=git/2.45.0"), &format!("want {B}"), "deepen 1", "filter blob:none"],
773 &["done"],
774 );
775 let two = v0(
776 &[&format!("want {B} ofs-delta thin-pack side-band-64k multi_ack_detailed agent=git/2.47.1"), "filter blob:none", "deepen 1", &format!("want {A}")],
777 &["done"],
778 );
779 let first = normalize(0, &one).unwrap();
780 assert_eq!(first, normalize(1, &two).unwrap());
781 assert!(!first.contains("agent="));
782 let deeper = v0(&[&format!("want {A} side-band-64k"), &format!("want {B}"), "deepen 3"], &["done"]);
783 let shallow = v0(&[&format!("want {A} side-band-64k"), &format!("want {B}"), "deepen 1"], &["done"]);
784 assert_ne!(normalize(0, &deeper), normalize(0, &shallow));
785 // The same request read as v2 is not a v2 fetch.
786 assert_eq!(normalize(2, &one), None);
787 }
788
789 #[test]
790 fn a_request_with_haves_is_never_kept() {
791 let v2_have = v2(&["object-format=sha1"], &[&format!("want {A}"), &format!("have {B}"), "done"]);
792 assert_eq!(normalize(2, &v2_have), None);
793 let v2_negotiating = v2(&["object-format=sha1"], &[&format!("want {A}"), &format!("have {B}")]);
794 assert_eq!(normalize(2, &v2_negotiating), None);
795 let v0_have = v0(&[&format!("want {A} side-band-64k")], &[&format!("have {B}"), "done"]);
796 assert_eq!(normalize(0, &v0_have), None);
797 let v0_more = v0(&[&format!("want {A} side-band-64k")], &[&format!("have {B}")]);
798 assert_eq!(normalize(0, &v0_more), None);
799 // A shallow clone deepened: the client has commits already.
800 let deepened = v2(&["object-format=sha1"], &[&format!("want {A}"), &format!("shallow {B}"), "deepen 5", "done"]);
801 assert_eq!(normalize(2, &deepened), None);
802 }
803
804 #[test]
805 fn only_whole_fetches_with_wants_are_kept() {
806 let ls_refs = [pkt("command=ls-refs\n"), b"0001".to_vec(), pkt("peel\n"), b"0000".to_vec()].concat();
807 assert_eq!(normalize(2, &ls_refs), None);
808 let no_wants = v2(&["object-format=sha1"], &["done"]);
809 assert_eq!(normalize(2, &no_wants), None);
810 let clone = v2(&["object-format=sha1"], &[&format!("want {A}"), "done"]);
811 assert!(normalize(2, &clone).is_some());
812 assert_eq!(normalize(2, &clone[..clone.len() - 2]), None);
813 let mut trailing = clone.clone();
814 trailing.extend(pkt("done\n"));
815 assert_eq!(normalize(2, &trailing), None);
816 // v0 without `done` is not answered with a pack.
817 assert_eq!(normalize(0, &v0(&[&format!("want {A} side-band-64k")], &[])), None);
818 assert_eq!(normalize(0, b"\x1f\x8b\x08\x00gzip"), None);
819 }
820
821 fn git(endpoint: &'static str, service: GitService) -> GitRequest {
822 GitRequest {
823 path: g1t_contracts::repos::RepoPath { namespace: "acme".into(), name: "rocket".into() },
824 endpoint,
825 service,
826 }
827 }
828
829 #[test]
830 fn the_cacheable_decision() {
831 let clone = v2(&["object-format=sha1"], &[&format!("want {A}"), "deepen 1", "done"]);
832 let fetch = git("git-upload-pack", GitService::UploadPack);
833 assert!(cacheable(&fetch, false, 2, None, Some(&clone)).is_some());
834 assert!(cacheable(&fetch, false, 2, Some("identity"), Some(&clone)).is_some());
835 assert_eq!(cacheable(&fetch, false, 2, Some("gzip"), Some(&clone)), None);
836 assert_eq!(cacheable(&fetch, true, 2, None, Some(&clone)), None);
837 assert_eq!(cacheable(&fetch, false, 2, None, None), None);
838 assert_eq!(cacheable(&git("info/refs", GitService::UploadPack), false, 2, None, Some(&clone)), None);
839 assert_eq!(cacheable(&git("git-receive-pack", GitService::ReceivePack), false, 2, None, Some(&clone)), None);
840 // A key moves with the refs version, and differs by repository.
841 let normalized = cacheable(&fetch, false, 2, None, Some(&clone)).unwrap();
842 assert_ne!(Key::new("r1", 4, &normalized), Key::new("r1", 5, &normalized));
843 assert_ne!(Key::new("r1", 4, &normalized), Key::new("r2", 4, &normalized));
844 assert_eq!(Key::new("r1", 4, &normalized), Key::new("r1", 4, &normalized));
845 assert!(Key::new("r1", 4, &normalized).as_str().starts_with("packs/r1/4/"));
846 }
847
848 /// A v2 answer: shallow-info, then the pack in channel 1, progress in 2.
849 fn answer(pack: &[u8]) -> Vec<u8> {
850 let mut out = [pkt("shallow-info\n"), pkt(&format!("shallow {A}\n")), b"0001".to_vec(), pkt("packfile\n")].concat();
851 out.extend(pkt_bytes(b"\x02Enumerating objects: 3, done.\n"));
852 let mut data = b"PACK\0\0\0\x02\0\0\0\x03".to_vec();
853 data.extend_from_slice(pack);
854 for chunk in data.chunks(1000) {
855 let mut packet = vec![1u8];
856 packet.extend_from_slice(chunk);
857 out.extend(pkt_bytes(&packet));
858 }
859 out.extend(b"0000");
860 out
861 }
862
863 #[test]
864 fn a_whole_pack_passes_the_check_in_any_chunks() {
865 let body = answer(&[7u8; 5000]);
866 for size in [1, 3, 4, 7, 64, 1000, body.len()] {
867 let mut check = PackCheck::default();
868 for chunk in body.chunks(size) {
869 check.feed(chunk);
870 }
871 assert!(check.complete(), "chunks of {size}");
872 }
873 // v0: NAK, then the pack in channel 1.
874 let v0 = [pkt("NAK\n"), pkt_bytes(b"\x01PACK\0\0\0\x02\0\0\0\x01xyz"), b"0000".to_vec()].concat();
875 let mut check = PackCheck::default();
876 check.feed(&v0);
877 assert!(check.complete());
878 }
879
880 #[test]
881 fn a_cut_short_or_failed_answer_is_not_a_pack() {
882 let body = answer(&[7u8; 5000]);
883 let mut check = PackCheck::default();
884 check.feed(&body[..body.len() - 4]);
885 assert!(!check.complete());
886 let mut check = PackCheck::default();
887 check.feed(&body[..body.len() - 100]);
888 assert!(!check.complete());
889 let failed = [pkt("packfile\n"), pkt_bytes(b"\x03fatal: out of memory\n"), b"0000".to_vec()].concat();
890 let mut check = PackCheck::default();
891 check.feed(&failed);
892 assert!(check.failed() && !check.complete());
893 let err = [pkt("ERR upload-pack: not our ref\n"), b"0000".to_vec()].concat();
894 let mut check = PackCheck::default();
895 check.feed(&err);
896 assert!(!check.complete());
897 let no_pack = [pkt("acknowledgments\n"), pkt("NAK\n"), b"0000".to_vec()].concat();
898 let mut check = PackCheck::default();
899 check.feed(&no_pack);
900 assert!(!check.complete());
901 let mut check = PackCheck::default();
902 check.feed(b"zzzz");
903 assert!(check.failed());
904 }
905
906 /// Kept objects, and multipart uploads in progress.
907 #[derive(Default)]
908 struct Memory {
909 objects: RefCell<HashMap<String, Vec<u8>>>,
910 uploads: RefCell<HashMap<String, Vec<Vec<u8>>>>,
911 aborted: Cell<u32>,
912 fail_parts_after: Option<u16>,
913 }
914
915 impl PackStore for Memory {
916 async fn get(&self, key: &str) -> Result<Option<(u64, ResponseBody)>> {
917 Ok(self.objects.borrow().get(key).map(|bytes| (bytes.len() as u64, ResponseBody::Body(bytes.clone()))))
918 }
919 async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()> {
920 self.objects.borrow_mut().insert(key.to_owned(), bytes);
921 Ok(())
922 }
923 async fn begin(&self, key: &str) -> Result<String> {
924 let id = format!("upload-{key}");
925 self.uploads.borrow_mut().insert(id.clone(), Vec::new());
926 Ok(id)
927 }
928 async fn part(&self, _key: &str, upload: &str, number: u16, bytes: Vec<u8>) -> Result<String> {
929 if self.fail_parts_after.is_some_and(|after| number > after) {
930 return Err(worker::Error::RustError("part refused".into()));
931 }
932 let mut uploads = self.uploads.borrow_mut();
933 let parts = uploads.get_mut(upload).unwrap();
934 assert_eq!(parts.len() + 1, number as usize);
935 parts.push(bytes);
936 Ok(format!("etag-{number}"))
937 }
938 async fn complete(&self, key: &str, upload: &str, parts: Vec<(u16, String)>) -> Result<()> {
939 let uploaded = self.uploads.borrow_mut().remove(upload).unwrap();
940 assert_eq!(parts.len(), uploaded.len());
941 for (index, part) in uploaded.iter().enumerate() {
942 if index + 1 < uploaded.len() {
943 assert_eq!(part.len(), PART_BYTES, "every part but the last is the same size");
944 }
945 assert!(!part.is_empty());
946 }
947 self.objects.borrow_mut().insert(key.to_owned(), uploaded.concat());
948 Ok(())
949 }
950 async fn abort(&self, _key: &str, upload: &str) -> Result<()> {
951 self.uploads.borrow_mut().remove(upload);
952 self.aborted.set(self.aborted.get() + 1);
953 Ok(())
954 }
955 }
956
957 fn block_on<F: Future>(future: F) -> F::Output {
958 let mut future = std::pin::pin!(future);
959 let mut cx = Context::from_waker(Waker::noop());
960 loop {
961 if let Poll::Ready(output) = future.as_mut().poll(&mut cx) {
962 return output;
963 }
964 }
965 }
966
967 /// Streams `body` through a tee in `chunk`-sized pieces, as git would
968 /// read it, stopping after `read` chunks if given; then runs the fill.
969 fn run(store: &Memory, body: &[u8], chunk: usize, read: Option<usize>, error_at: Option<usize>) -> (Vec<u8>, Filled, u64) {
970 let mut chunks: Vec<Result<Vec<u8>>> = body.chunks(chunk).map(|c| Ok(c.to_vec())).collect();
971 if let Some(at) = error_at {
972 chunks.truncate(at);
973 chunks.push(Err(worker::Error::RustError("store went away".into())));
974 }
975 let pipe = SharedPipe::default();
976 let measured = Rc::new(Cell::new(None));
977 let told = measured.clone();
978 let mut tee = Tee::new(Box::pin(futures_util::stream::iter(chunks)), Some(pipe.clone()), Box::new(move |bytes| told.set(Some(bytes))));
979 // Git reads, and the fill keeps up, a piece at a time.
980 let mut got = Vec::new();
981 let mut drain = Drain(pipe);
982 let mut kept = Vec::new();
983 let mut taken = 0;
984 loop {
985 if read.is_some_and(|read| taken >= read) {
986 break;
987 }
988 match block_on(tee.next()) {
989 Some(Ok(bytes)) => got.extend(bytes),
990 Some(Err(_)) | None => break,
991 }
992 taken += 1;
993 while let Poll::Ready(Some(piece)) = Pin::new(&mut drain).poll_next(&mut Context::from_waker(Waker::noop())) {
994 kept.push(piece);
995 }
996 }
997 drop(tee);
998 while let Poll::Ready(Some(piece)) = Pin::new(&mut drain).poll_next(&mut Context::from_waker(Waker::noop())) {
999 kept.push(piece);
1000 }
1001 let filled = block_on(fill(store, "packs/r/1/k", futures_util::stream::iter(kept)));
1002 (got, filled, measured.get().unwrap())
1003 }
1004
1005 #[test]
1006 fn a_small_pack_is_kept_whole_once_it_has_all_arrived() {
1007 let store = Memory::default();
1008 let body = answer(&[9u8; 20_000]);
1009 let (got, filled, measured) = run(&store, &body, 4096, None, None);
1010 assert_eq!(got, body);
1011 assert_eq!(measured, body.len() as u64);
1012 assert_eq!(filled, Filled::Kept { bytes: body.len() as u64 });
1013 assert_eq!(store.objects.borrow().get("packs/r/1/k"), Some(&body));
1014 }
1015
1016 #[test]
1017 fn a_large_pack_goes_up_in_equal_parts() {
1018 let store = Memory::default();
1019 // Just over two parts, in chunks that do not divide a part.
1020 let body = answer(&vec![5u8; 2 * PART_BYTES + 10_000]);
1021 let (got, filled, _) = run(&store, &body, 64 * 1024 + 3, None, None);
1022 assert_eq!(got, body);
1023 assert_eq!(filled, Filled::Kept { bytes: body.len() as u64 });
1024 assert_eq!(store.objects.borrow().get("packs/r/1/k"), Some(&body));
1025 assert!(store.uploads.borrow().is_empty());
1026 }
1027
1028 #[test]
1029 fn a_fill_cut_short_leaves_nothing_to_serve() {
1030 // Git went away after a few chunks.
1031 let store = Memory::default();
1032 let body = answer(&vec![5u8; 2 * PART_BYTES]);
1033 let (_, filled, measured) = run(&store, &body, 64 * 1024, Some(90), None);
1034 assert_eq!(filled, Filled::Abandoned);
1035 assert_eq!(measured, 90 * 64 * 1024);
1036 assert!(store.objects.borrow().is_empty());
1037 assert!(store.uploads.borrow().is_empty());
1038 assert_eq!(store.aborted.get(), 1);
1039 // The store's answer failed partway.
1040 let store = Memory::default();
1041 let (_, filled, _) = run(&store, &body, 64 * 1024, None, Some(10));
1042 assert_eq!(filled, Filled::Abandoned);
1043 assert!(store.objects.borrow().is_empty());
1044 // A part the bucket refused.
1045 let store = Memory { fail_parts_after: Some(0), ..Memory::default() };
1046 let (got, filled, _) = run(&store, &body, 64 * 1024, None, None);
1047 assert_eq!(got, body, "git has its answer whatever the bucket does");
1048 assert_ne!(filled, Filled::Kept { bytes: body.len() as u64 });
1049 assert!(store.objects.borrow().is_empty());
1050 // An answer that is not a pack.
1051 let store = Memory::default();
1052 let error = [pkt("ERR not our ref\n"), b"0000".to_vec()].concat();
1053 let (_, filled, _) = run(&store, &error, 4096, None, None);
1054 assert_eq!(filled, Filled::NotAPack);
1055 assert!(store.objects.borrow().is_empty());
1056 }
1057
1058 #[test]
1059 fn a_pack_past_the_cap_streams_through_unkept() {
1060 // The cap itself is 200 MB; the tee lets go of a slow bucket the
1061 // same way, which is quicker to show: nothing drains the pipe.
1062 let pipe = SharedPipe::default();
1063 let chunks: Vec<Result<Vec<u8>>> = (0..200).map(|_| Ok(vec![1u8; 64 * 1024])).collect();
1064 let mut tee = Tee::new(Box::pin(futures_util::stream::iter(chunks)), Some(pipe.clone()), Box::new(|_| {}));
1065 let mut sent = 0;
1066 while let Some(Ok(bytes)) = block_on(tee.next()) {
1067 sent += bytes.len();
1068 }
1069 assert_eq!(sent, 200 * 64 * 1024);
1070 let pipe = pipe.borrow();
1071 assert!(pipe.closed);
1072 assert!(pipe.queued <= MAX_QUEUED_BYTES);
1073 assert!(!pipe.pieces.contains(&Piece::End), "a fill let go never hears the end");
1074 }
1075
1076 #[test]
1077 fn the_meters_are_not_operations() {
1078 let mapping = crate::meters::Mapping::defaults();
1079 for meter in [HIT, MISS] {
1080 assert_eq!(mapping.billable(meter), 0.0);
1081 assert_eq!(mapping.cost(meter), 0.0);
1082 }
1083 }
1084
1085 #[test]
1086 fn fills_are_one_per_key_and_few_at_once() {
1087 let first = Filling::take("a").unwrap();
1088 assert!(Filling::take("a").is_none());
1089 let second = Filling::take("b").unwrap();
1090 assert!(Filling::take("c").is_none());
1091 drop(first);
1092 assert!(Filling::take("c").is_some());
1093 drop(second);
1094 }
1095}