Commit

Merge branch 'worktree-agent-a1b995daa94e4e1b7'

syntaqxcommitted Parents8445a3d3d1cec8Browse files
14 files+1430−80/14 viewed
+8−0
223223 | Entry | Values |
224224 | --- | --- |
225225 | `refs;desc=` | `hit-colo` or `hit-shared` when the ref listing came from g1t's cache, `miss` when the git store was asked |
226+| `pack;desc=` | Only for a fresh clone: `hit` when its pack came from g1t's cache, `miss` when the git store built it |
226227 | `cred;desc=` | `isolate` or `shared` for a store credential made a moment ago, `mint` for a new one |
227228
228229 The ref listing git asks for first on every clone and fetch is kept for up
231232 tags makes the next fetch ask the git store again. A change can take up to
232233 5 seconds to reach every fetch.
233234
235+A fresh clone, one that has no objects yet (shallow clones such as
236+`git clone --depth=1` included), has its pack kept too, for up to 7 days
237+or until the repository's branches or tags next change. The next clone that
238+asks for the same commits in the same way gets the same pack without the
239+git store building it again. A fetch into a repository you already have,
240+and any pack over 200 MB, always goes to the git store.
241+
234242 Include the header when you report a slow clone, fetch or push.
235243
236244 ## SSH
+3−1
327327 An operation is one clone or fetch (a request that fetches objects) or one
328328 push, and making, forking or deleting a repository: what Cloudflare bills
329329 g1t for. Listing refs, and anything g1t answers from its own cache, is
330−never an operation. g1t meters every request it makes to the store, and
330+never an operation: a repeat clone of the same commit is served from
331+g1t's [pack cache](/guides/git/#where-a-slow-requests-time-went) and is not
332+counted. g1t meters every request it makes to the store, and
331333 what counts follows what Cloudflare confirms it bills; this page changes
332334 with it.
333335
+4−1
7777 "d1": { "database": "g1t-repos", "migrations": "migrations" },
7878 "stage": "core",
7979 "secrets": ["REPOS_KEY"],
80− "setup": ["The Artifacts namespace `g1t` (the ARTIFACTS binding)"],
80+ "setup": [
81+ "The Artifacts namespace `g1t` (the ARTIFACTS binding)",
82+ "The R2 bucket `g1t-git-packs` (GIT_PACKS) with its lifecycle rule: `npx wrangler r2 bucket create g1t-git-packs`, then `npx wrangler r2 bucket lifecycle add g1t-git-packs expire-packs packs/ --expire-days 7 --abort-multipart-days 1`"
83+ ],
8184 "self_host": "run"
8285 },
8386 "work": {
+1−0
364364 | R10 | Built in repos; work unchanged | `divergence` works out the target's side once per target head per isolate (`coalesce.rs`: the head under the refs version, then the history by hash, kept 60 s), and what the target changed between two trees once per pair (10 minutes). `readCommit` and logs by hash come from the cache. Work's fan-out (`after_push`, up to 100 pull requests) is unchanged: its 100 `divergence` calls now cost one walk of the target instead of 100. |
365365 | R7 | Groundwork | `shards.rs`: bindings named in `ARTIFACTS_NAMESPACES` (JSON, binding → namespace; `ARTIFACTS` → `g1t` always there), a repository's namespace kept in its `store` column as `<namespace>/<key>` (no prefix means the `ARTIFACTS` namespace, so every existing key reads the same), new repositories placed by `ARTIFACTS_NEW_REPOS` (comma-separated, spread by an FNV hash of the repository id; names not bound are skipped), forks always in their repository's namespace, `ARTIFACTS_EU_NAMESPACE` reserved for EU residency (no workspace setting yet). Works with only `ARTIFACTS` bound, as today. |
366366 | R8 | Built | `crates/runner/src/clone.rs`: every sandbox clones at `--depth=1` (a full g1t clone took 5.4 s, depth 1 took 3.8 s). Work that merges (catch-up, the merge queue, merge checks, a review's diff) deepens 50, 500, then 5000 commits until the two sides share one, and fetches everything only as the last resort (`share_history`). `G1T_CLONE_DEPTH` (0 or `full` for everything) and `G1T_CLONE_FILTER=blob:none` change it per runner. |
367+| R6 | Built; the bucket must exist before it deploys | `pack_cache.rs`: an upload-pack POST with wants and no `have` or `shallow` lines (a fresh clone, the sandboxes' `deepen 1` ones included), uncompressed and at most 1 MiB, is keyed `packs/<repo id>/<refs_version>/<sha256>` over the request normalized: protocol v2 capabilities without `agent=`/`session-id=` and its arguments, each sorted and deduplicated; v0/v1 wants sorted, the first want's capabilities split off, sorted and without `agent=`, then `deepen`/`filter` lines, a flush and `done`. Only while `refs_cache::usable` (the version known, and no push credential out of g1t's hands), so never across a refs change. Looked up after authorization, alongside the free-workspace limits and the kept refs answer; a hit streams from the bucket (`Server-Timing` `pack;desc=hit`). A miss streams the store's 200 to git through a tee that copies it to a fill in `ctx.wait_until` (at most 5 MiB queued between them, 2 fills per isolate, one per key): under 5 MiB it is one `put` once it all arrived; larger, 5 MiB multipart parts completed only after the last part and a check that it is one whole side-band pack (well-formed pkt-lines, `PACK` on channel 1, no `ERR` or channel 3, a closing flush). Over 200 MB, a queue that falls behind, git going away or the store's stream failing lets the fill go and aborts the upload; nothing partial can be read. Meters `pack_cache.hit` (with the bytes served) and `pack_cache.miss` (counted with `record`, bytes added at the end), neither an operation by default; a hit records no `git.fetch`. Storage is behind the `PackStore` port with an R2 adapter (`GIT_PACKS`, bucket `g1t-git-packs`, lifecycle: packs deleted after 7 days, unfinished uploads after 1); without the binding (self-hosted) nothing is kept. |
367368
368369 ### R1: reading `scripts/ops/artifacts-usage.mjs`
369370
+8−2
493493 - Queues: `npx wrangler queues create <queue>` for each queue in the
494494 manifest: `g1t-events`, `g1t-events-<service>` for every subscriber,
495495 `g1t-search-jobs`, `g1t-context-jobs`.
496−- R2: `npx wrangler r2 bucket create g1t-screenshots` and
496+- R2: `npx wrangler r2 bucket create g1t-screenshots`,
497497 `npx wrangler r2 bucket create g1t-actions-cache` (with its 30-day
498− lifecycle rule, above).
498+ lifecycle rule, above), and `npx wrangler r2 bucket create g1t-git-packs`,
499+ the clone pack cache (`services/repos/src/pack_cache.rs`), with a rule
500+ deleting packs 7 days after they were written and unfinished uploads
501+ after a day:
502+ `npx wrangler r2 bucket lifecycle add g1t-git-packs expire-packs packs/ --expire-days 7 --abort-multipart-days 1`.
503+ A Worker bound to a bucket that does not exist fails to deploy, so make
504+ it before the first deploy of `g1t-repos` that binds it.
499505 - The runner's base image: `node scripts/deploy.mjs build-base`.
500506 - Vectorize, dispatch namespace, DNS, Access, Email Sending, Artifacts: each
501507 unit's `setup`.
+1−1
102102 | `apps/docs` | Static | — | Not run (docs.g1t.sh serves them) |
103103 | `apps/status` | TS Worker | Email Sending, cron; bound only to billing | Runs in a process of its own (`status.sh`), so it stays up when the site does not |
104104 | `services/identity` | Rust | Email Sending, KV `AVATARS` | Runs unchanged; `EMAIL` goes to the mail shim |
105−| `services/repos` | Rust | **Artifacts**, Cache API, optional KV `GIT_CACHE` with `REPOS_KEY` | Runs unchanged; `ARTIFACTS` goes to the git store. Without `GIT_CACHE` and `REPOS_KEY`, credentials and ref listings are kept per isolate only |
105+| `services/repos` | Rust | **Artifacts**, Cache API, optional KV `GIT_CACHE` with `REPOS_KEY`, optional R2 `GIT_PACKS` | Runs unchanged; `ARTIFACTS` goes to the git store. Without `GIT_CACHE` and `REPOS_KEY`, credentials and ref listings are kept per isolate only. `GIT_PACKS` (the clone pack cache, behind the `PackStore` port in `src/pack_cache.rs`) is not given, so every clone goes to the git store; an S3 adapter like packages' would turn it on |
106106 | `services/work` | Rust | Queue consumer | Runs unchanged |
107107 | `services/events` | Rust | Queues (producer and fan-out) | Runs unchanged; the off services' queues are not produced to |
108108 | `services/projects` | TS | Queue consumer | Runs unchanged |
+8−0
1+// The Artifacts binding for running repos on its own (repos.jsonc): the
2+// self-hosted shim, in front of a git store on this machine.
3+{
4+ "name": "g1t-artifacts",
5+ "main": "../../../deploy/self-host/workers/artifacts/index.js",
6+ "compatibility_date": "2026-09-26",
7+ "vars": { "GITSTORE_URL": "http://localhost:8799", "GITSTORE_SECRET": "dev-gitstore-secret-0123" }
8+}
+164−0
1+#!/usr/bin/env node
2+// Clones a repository through the repos service twice, full and shallow,
3+// over protocol v2 and v0, and checks that the second clone of each came
4+// from the pack cache (src/pack_cache.rs) and matches the first. Then
5+// moves the repository's refs version and checks the next clone misses.
6+//
7+// Runs everything on this machine: the self-hosted git store, and
8+// `wrangler dev` with dev/repos.jsonc, dev/artifacts.jsonc and
9+// dev/stubs.jsonc, state kept in a temporary folder. Build first:
10+//
11+// cd services/repos && node ../../scripts/build-rust-worker.mjs
12+// node dev/clone-check.mjs
13+//
14+// Needs node, git (with git-http-backend) and the repository's npm packages.
15+
16+import { spawn, spawnSync } from "node:child_process";
17+import { mkdtempSync, rmSync, writeFileSync } from "node:fs";
18+import { randomBytes } from "node:crypto";
19+import { tmpdir } from "node:os";
20+import { dirname, join, resolve } from "node:path";
21+import { fileURLToPath } from "node:url";
22+import { createRequire } from "node:module";
23+
24+const here = dirname(fileURLToPath(import.meta.url));
25+const service = resolve(here, "..");
26+const root = resolve(service, "../..");
27+const work = mkdtempSync(join(tmpdir(), "g1t-clone-check-"));
28+const persist = join(work, "state");
29+const SECRET = "dev-gitstore-secret-0123";
30+const STORE = "http://localhost:8799";
31+const REPOS = "http://localhost:8791";
32+const KEY = "acme--rocket";
33+const children = [];
34+// Wrangler from the repository's packages, run with node: no shell to quote for.
35+const WRANGLER = join(dirname(createRequire(join(service, "package.json")).resolve("wrangler/package.json")), "bin/wrangler.js");
36+
37+function run(command, args, options = {}) {
38+ const done = spawnSync(command, args, { encoding: "utf8", ...options });
39+ if (done.status !== 0 && !options.allowFail) {
40+ throw new Error(`${command} ${args.join(" ")} failed:\n${done.stdout}\n${done.stderr}`);
41+ }
42+ return done;
43+}
44+
45+const git = (args, cwd = work, env = {}) => run("git", args, { cwd, env: { ...process.env, ...env } });
46+
47+function start(command, args, options) {
48+ const child = spawn(command, args, { ...options });
49+ children.push(child);
50+ return child;
51+}
52+
53+async function waitFor(url, what) {
54+ for (let i = 0; i < 120; i++) {
55+ try {
56+ const response = await fetch(url);
57+ if (response.status < 500) return;
58+ } catch {}
59+ await new Promise((resolve) => setTimeout(resolve, 500));
60+ }
61+ throw new Error(`${what} did not start`);
62+}
63+
64+const sql = (command) =>
65+ run("node", [WRANGLER, "d1", "execute", "g1t-repos", "--local", "--persist-to", persist, "-c", "dev/repos.jsonc", "--command", command], {
66+ cwd: service,
67+ env: { ...process.env, CI: "1" },
68+ });
69+
70+/** Clones with `args`, and says how the pack was found and what came. */
71+function clone(name, args) {
72+ const dir = join(work, name);
73+ const done = git(["clone", ...args, `${REPOS}/acme/rocket.git`, dir], work, { GIT_TRACE_CURL: "1", GIT_TRACE_CURL_NO_DATA: "1" });
74+ const timings = done.stderr.split("\n").filter((line) => /server-timing:/i.test(line) && /pack;desc=/.test(line));
75+ const pack = timings.map((line) => /pack;desc=(\w+)/.exec(line)[1]);
76+ const head = git(["rev-parse", "HEAD"], dir).stdout.trim();
77+ const files = git(["ls-tree", "-r", "HEAD"], dir).stdout;
78+ const count = git(["rev-list", "--count", "HEAD"], dir).stdout.trim();
79+ git(["fsck", "--no-progress"], dir);
80+ return { pack, head, files, count };
81+}
82+
83+const checks = [];
84+function check(what, ok, detail = "") {
85+ checks.push({ what, ok });
86+ console.log(`${ok ? "ok " : "FAIL"} ${what}${detail ? ` (${detail})` : ""}`);
87+}
88+
89+function twice(label, args, depth) {
90+ const first = clone(`${label}-1`, args);
91+ const second = clone(`${label}-2`, args);
92+ check(`${label}: the first clone misses`, first.pack.includes("miss"), first.pack.join(","));
93+ check(`${label}: the second clone hits`, second.pack.includes("hit"), second.pack.join(","));
94+ check(`${label}: both clones are the same`, first.head === second.head && first.files === second.files && first.count === second.count);
95+ if (depth) check(`${label}: ${depth} commit(s) of history`, second.count === String(depth), second.count);
96+ return second;
97+}
98+
99+try {
100+ // The git store, with a repository of a few commits.
101+ start("node", [join(root, "deploy/self-host/gitstore/server.mjs")], {
102+ env: { ...process.env, GITSTORE_ROOT: join(work, "git"), GITSTORE_SECRET: SECRET, GITSTORE_PORT: "8799", GITSTORE_URL: STORE },
103+ stdio: "inherit",
104+ });
105+ await waitFor(`${STORE}/healthz`, "the git store");
106+ const api = (path, body) =>
107+ fetch(`${STORE}/api/repos${path}`, {
108+ method: "POST",
109+ headers: { "x-gitstore-secret": SECRET, "content-type": "application/json" },
110+ body: JSON.stringify(body),
111+ }).then((response) => response.json());
112+ await api("", { name: KEY, defaultBranch: "main" });
113+ const token = (await api(`/${KEY}/tokens`, { scope: "write" })).plaintext;
114+ const seed = join(work, "seed");
115+ git(["init", "-q", "-b", "main", seed]);
116+ for (let i = 1; i <= 5; i++) {
117+ writeFileSync(join(seed, `file-${i}.txt`), `${"line\n".repeat(200 * i)}${i}\n`);
118+ // 12 MB that does not compress, so a pack goes up in multipart parts.
119+ if (i === 5) writeFileSync(join(seed, "noise.bin"), randomBytes(12 * 1024 * 1024));
120+ git(["add", "."], seed);
121+ git(["-c", "user.name=dev", "-c", "user.email=dev@example.com", "commit", "-q", "-m", `commit ${i}`], seed);
122+ }
123+ git(["-c", `http.extraHeader=Authorization: Bearer ${token}`, "push", "-q", `${STORE}/git/${KEY}.git`, "main"], seed);
124+
125+ // The repos service, its database with the repository in it.
126+ run("node", [WRANGLER, "d1", "migrations", "apply", "g1t-repos", "--local", "--persist-to", persist, "-c", "dev/repos.jsonc"], {
127+ cwd: service,
128+ env: { ...process.env, CI: "1" },
129+ });
130+ sql("INSERT INTO repos (id, namespace, name, is_private, owner_id, default_branch, refs_version) VALUES ('rep_rocket', 'acme', 'rocket', 0, 'usr_dev', 'main', 1)");
131+ start(
132+ "node",
133+ [WRANGLER, "dev", "-c", "dev/repos.jsonc", "-c", "dev/artifacts.jsonc", "-c", "dev/stubs.jsonc", "--local", "--persist-to", persist, "--port", "8791"],
134+ { cwd: service, env: { ...process.env, CI: "1" }, stdio: ["ignore", "inherit", "inherit"] },
135+ );
136+ await waitFor(`${REPOS}/acme/rocket.git/info/refs?service=git-upload-pack`, "wrangler dev");
137+
138+ twice("full, v2", [], 5);
139+ twice("shallow, v2", ["--depth=1"], 1);
140+ twice("full, v0", ["-c", "protocol.version=0"], 5);
141+ twice("shallow, v0", ["-c", "protocol.version=0", "--depth=1"], 1);
142+
143+ // A change to the refs: the next clone goes to the store.
144+ sql("UPDATE repos SET refs_version = refs_version + 1 WHERE id = 'rep_rocket'");
145+ // The service keeps a row it read a moment ago for the same clone's next request.
146+ await new Promise((resolve) => setTimeout(resolve, 6000));
147+ const after = clone("after-refs", ["--depth=1"]);
148+ check("after the refs version moves, a clone misses", after.pack.includes("miss"), after.pack.join(","));
149+} catch (error) {
150+ console.error(error);
151+ checks.push({ what: "ran", ok: false });
152+} finally {
153+ for (const child of children) {
154+ if (process.platform === "win32") spawnSync("taskkill", ["/pid", String(child.pid), "/t", "/f"], { stdio: "ignore" });
155+ else child.kill();
156+ }
157+ try {
158+ rmSync(work, { recursive: true, force: true });
159+ } catch {}
160+}
161+
162+const failed = checks.filter((c) => !c.ok);
163+console.log(failed.length ? `\n${failed.length} of ${checks.length} checks failed.` : `\nAll ${checks.length} checks passed.`);
164+process.exit(failed.length ? 1 : 0);
+43−0
1+// The repos service as `wrangler dev` runs it on its own, for git over
2+// HTTPS against a local git store: local D1 and R2 (the clone pack cache,
3+// GIT_PACKS), the Artifacts binding played by the self-hosted shim
4+// (deploy/self-host/workers/artifacts) in front of the self-hosted git store
5+// (deploy/self-host/gitstore/server.mjs), and its other services stubbed
6+// (stubs.js). `node dev/clone-check.mjs` does all of it and clones twice;
7+// by hand, from services/repos:
8+//
9+// node ../../scripts/build-rust-worker.mjs
10+// npx wrangler d1 migrations apply g1t-repos --local -c dev/repos.jsonc
11+// GITSTORE_ROOT=/tmp/g1t-git GITSTORE_SECRET=dev-gitstore-secret-0123 \
12+// GITSTORE_PORT=8799 GITSTORE_URL=http://localhost:8799 \
13+// node ../../deploy/self-host/gitstore/server.mjs &
14+// npx wrangler dev -c dev/repos.jsonc -c dev/artifacts.jsonc -c dev/stubs.jsonc --local --port 8791
15+//
16+// Then make a repository row and its store (clone-check.mjs shows how) and
17+// `git clone http://localhost:8791/acme/rocket.git`.
18+{
19+ "name": "g1t-repos",
20+ "main": "../build/index.js",
21+ "compatibility_date": "2026-09-26",
22+ "d1_databases": [
23+ { "binding": "DB", "database_name": "g1t-repos", "database_id": "local", "migrations_dir": "../migrations" }
24+ ],
25+ "r2_buckets": [{ "binding": "GIT_PACKS", "bucket_name": "g1t-git-packs" }],
26+ "services": [
27+ { "binding": "ARTIFACTS", "service": "g1t-artifacts" },
28+ { "binding": "IDENTITY", "service": "g1t-repos-stubs" },
29+ { "binding": "EVENTS", "service": "g1t-repos-stubs" },
30+ { "binding": "SECURITY", "service": "g1t-repos-stubs" },
31+ { "binding": "BILLING", "service": "g1t-repos-stubs" }
32+ ],
33+ "vars": {
34+ "GIT_OPERATIONS_FREE_CAP": "50000",
35+ "GIT_OPERATIONS_FREE_HOURLY": "100000",
36+ "FREE_PRIVATE_STORAGE_BYTES": "1000000000",
37+ "REPO_STORAGE_LIMIT_BYTES": "950000000",
38+ "LARGE_PUSHES": "unscanned",
39+ "FORK_RETENTION_DAYS": "1",
40+ "ARTIFACTS_NAMESPACES": "{\"ARTIFACTS\":\"g1t\"}",
41+ "ARTIFACTS_NEW_REPOS": ""
42+ }
43+}
+26−0
1+// Stand-ins for identity, events, security and billing, so the repos
2+// service answers anonymous git requests with `wrangler dev` (repos.jsonc).
3+//
4+// - Identity knows nobody: credentials name no one, and no workspace was
5+// renamed. Public repositories can be cloned without signing in.
6+// - Events takes every event and audit entry and logs them.
7+// - Security has allowed no secrets; billing says every workspace is free.
8+
9+export default {
10+ async fetch(request) {
11+ const method = new URL(request.url).pathname.replace(/^\/rpc\//, "");
12+ const args = await request.json().catch(() => ({}));
13+ const json = (value) => Response.json(value);
14+ switch (method) {
15+ case "user_for_git_credentials":
16+ case "resolve_slug":
17+ return json(null);
18+ case "is_free":
19+ case "plan":
20+ return json({ free: true });
21+ default:
22+ console.log(`stub ${method}`, JSON.stringify(args).slice(0, 200));
23+ return json(null);
24+ }
25+ },
26+};
+6−0
1+// The stand-in services for running repos on its own (dev/stubs.js).
2+{
3+ "name": "g1t-repos-stubs",
4+ "main": "stubs.js",
5+ "compatibility_date": "2026-09-26"
6+}
+55−3
2020 mod listing;
2121 mod meters;
2222 mod mirror;
23+mod pack_cache;
2324 mod pack_limits;
2425 mod refs;
2526 mod refs_cache;
194195 placement: shards::Placement,
195196 /// What isolates share: answers that list refs (refs_cache.rs).
196197 shared: Option<Rc<shared::Shared>>,
198+ /// Packs for fresh clones (pack_cache.rs); `None` without the bucket.
199+ packs: Option<Rc<pack_cache::R2Packs>>,
197200 }
198201
199202 impl<S: GitStore> Repos<S> {
15131516 .map(|(kind, version)| {
15141517 refs_cache::Key::new(&repo.id, version, default_branch.as_deref(), protocol, &kind)
15151518 });
1519+ // A fresh clone's pack may have been kept too: see pack_cache.rs.
1520+ // Under the same refs version, so never across a change to them.
1521+ let pack_key = self
1522+ .packs
1523+ .as_ref()
1524+ .and_then(|_| {
1525+ let encoding = request.headers().get("content-encoding").ok().flatten();
1526+ pack_cache::cacheable(git, get, protocol, encoding.as_deref(), body.as_deref())
1527+ })
1528+ .zip(refs_cache::usable(registry::refs_state(&repo.id), now_ms()))
1529+ .map(|(normalized, version)| pack_cache::Key::new(&repo.id, version, &normalized));
15161530 // A kept answer and the free workspace limits, with a kept
15171531 // credential looked up alongside. A kept answer goes back without
15181532 // waiting for the credential, which it does not need.
1519− let ((answer, limited), kept_access) = {
1533+ let ((answer, pack, limited), kept_access) = {
15201534 let shared = self.shared.as_deref();
1521− let answer_and_limits = std::pin::pin!(futures_util::future::join(
1535+ let answer_and_limits = std::pin::pin!(futures_util::future::join3(
15221536 async {
15231537 match &kept_key {
15241538 Some(kept_key) => refs_cache::get(shared, kept_key).await,
15251539 None => None,
15261540 }
15271541 },
1542+ async {
1543+ match (&pack_key, self.packs.as_deref()) {
1544+ (Some(pack_key), Some(packs)) => pack_cache::get(packs, pack_key).await,
1545+ _ => None,
1546+ }
1547+ },
15281548 self.git_limits(call, git, &repo, env),
15291549 ));
15301550 let kept_access = std::pin::pin!(self.store.kept_access(&key, scope));
15311551 match futures_util::future::select(answer_and_limits, kept_access).await {
15321552 futures_util::future::Either::Left((first, kept_access)) => {
1533− let answered = first.0.is_some() || matches!(first.1, Ok(Some(_)) | Err(_));
1553+ let answered = first.0.is_some() || first.1.is_some() || matches!(first.2, Ok(Some(_)) | Err(_));
15341554 (first, if answered { None } else { kept_access.await })
15351555 }
15361556 futures_util::future::Either::Right((kept_access, first)) => (first.await, kept_access),
15421562 after.spawn(env, ctx);
15431563 return Ok(response);
15441564 }
1565+ if let Some(kept) = pack {
1566+ timing.note("pack", "hit");
1567+ // Never reached the store: never an operation.
1568+ let sent = body.as_ref().map_or(0, |body| body.len() as u64);
1569+ meters::record(pack_cache::HIT, &key, sent, kept.size);
1570+ after.ended(200, None);
1571+ after.spawn(env, ctx);
1572+ return kept.response();
1573+ }
1574+ if pack_key.is_some() {
1575+ timing.note("pack", "miss");
1576+ }
15451577 if let (Some((entry, found)), Some(kept_key)) = (answer, &kept_key) {
15461578 timing.note("refs", found.as_str());
15471579 if found == refs_cache::Found::Shared {
16771709 }
16781710 }
16791711 response = Response::from_bytes(body)?.with_headers(headers).with_status(status);
1712+ } else if let (Some(pack_key), Some(packs), true) = (&pack_key, &self.packs, forwarded.from_store) {
1713+ // A fresh clone the bucket did not have: counted, and its pack
1714+ // kept as it streams to git, when it is a whole one.
1715+ meters::record(pack_cache::MISS, &key, forwarded.sent, 0);
1716+ if status == 200 {
1717+ let store_key = key.clone();
1718+ let measured = Box::new(move |bytes: u64| meters::record_bytes(pack_cache::MISS, &store_key, 0, bytes));
1719+ let (teed, filling) = pack_cache::tee(response, packs.clone(), pack_key, measured)?;
1720+ response = teed;
1721+ if let Some(filling) = filling {
1722+ let pack_key = pack_key.clone();
1723+ ctx.wait_until(async move {
1724+ let filled = filling.await;
1725+ if !matches!(filled, pack_cache::Filled::Kept { .. } | pack_cache::Filled::Abandoned) {
1726+ worker::console_warn!("pack {} not kept: {filled:?}", pack_key.as_str());
1727+ }
1728+ });
1729+ }
1730+ }
16801731 }
16811732 after.ended(status, None);
16821733 if status == 200 && (forwarded.pack_bytes > 0 || !forwarded.pushed.is_empty()) {
18791930 registry: Registry { db: env.d1("DB")? },
18801931 store: ArtifactsStore::new(env, shared.clone())?,
18811932 shared,
1933+ packs: pack_cache::R2Packs::from_env(env).map(Rc::new),
18821934 events: env.service("EVENTS")?,
18831935 security: env.service("SECURITY").ok(),
18841936 billing: env.service("BILLING").ok(),
+1095−0
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) 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+
28+use std::cell::RefCell;
29+use std::collections::{BTreeSet, HashSet, VecDeque};
30+use std::pin::Pin;
31+use std::rc::Rc;
32+use std::task::{Context, Poll, Waker};
33+
34+use futures_util::{Stream, StreamExt};
35+use g1t_contracts::repos::GitService;
36+use worker::{Bucket, Env, Headers, Response, ResponseBody, Result, UploadedPart};
37+
38+use 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.
42+macro_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.
54+pub const HIT: &str = "pack_cache.hit";
55+pub const MISS: &str = "pack_cache.miss";
56+
57+/// Packs larger than this pass through without being kept.
58+pub const MAX_PACK_BYTES: u64 = 200 * 1024 * 1024;
59+/// Requests larger than this (a great many wants) are not kept.
60+const 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.
63+const 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.
66+const MAX_QUEUED_BYTES: usize = 5 * 1024 * 1024;
67+/// Fills at once in one isolate; more clones pass through.
68+const MAX_FILLS: usize = 2;
69+pub 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)]
74+pub 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.
88+pub struct R2Packs {
89+ bucket: Bucket,
90+}
91+
92+impl 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+
99+impl 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.
142+fn 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.
165+fn 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).
175+pub 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.
185+fn 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`.
230+fn 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.
267+pub 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)]
279+pub struct Key(String);
280+
281+impl 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.
292+pub struct Kept {
293+ pub size: u64,
294+ pub body: ResponseBody,
295+}
296+
297+impl 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.
308+pub 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)]
322+pub 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`.
335+const START: usize = 5;
336+
337+impl 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)]
397+enum 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)]
406+struct Pipe {
407+ pieces: VecDeque<Piece>,
408+ queued: usize,
409+ closed: bool,
410+ waker: Option<Waker>,
411+}
412+
413+type SharedPipe = Rc<RefCell<Pipe>>;
414+
415+impl 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.
426+struct Drain(SharedPipe);
427+
428+impl 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+
447+type 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.
452+pub 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+
461+impl 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+
486+impl 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+
526+impl 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)]
539+pub 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.
553+async 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+
630+thread_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.
636+struct Filling(String);
637+
638+impl 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+
647+impl 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.
655+pub 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`.
667+pub 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.
672+pub 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)]
700+mod 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+}
+8−0
3030 // hex characters). Without the binding or the secret each isolate keeps
3131 // only its own. Made by `npx wrangler kv namespace create g1t-repos-git-cache`:
3232 "kv_namespaces": [{ "binding": "GIT_CACHE", "id": "be765052d0124c2a935b3db4dff99f1f" }],
33+ // Packs for fresh clones (wants and no haves), kept under the
34+ // repository's refs version so the next clone of the same commit skips
35+ // the git store (src/pack_cache.rs). Without the binding nothing is
36+ // kept. Made with a lifecycle rule that deletes packs after 7 days and
37+ // unfinished uploads after 1:
38+ // npx wrangler r2 bucket create g1t-git-packs
39+ // npx wrangler r2 bucket lifecycle add g1t-git-packs expire-packs packs/ --expire-days 7 --abort-multipart-days 1
40+ "r2_buckets": [{ "binding": "GIT_PACKS", "bucket_name": "g1t-git-packs" }],
3341 "services": [
3442 { "binding": "IDENTITY", "service": "g1t-identity" },
3543 { "binding": "EVENTS", "service": "g1t-events" },