g1t/services/repos/src/pack_cache.rs

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