From e85dac83192d9f1356c7f8653a21e5dd15e1d992 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:12:53 -0500 Subject: [PATCH 01/24] Bound a pack entry's size varint The size varint in a pack entry header had no length limit. After about 147 continuation bytes the multiplier overflows to Infinity, and a zero digit then makes the size NaN. NaN passes both the check against maxObjectSize and the running check during inflation, so the entry inflates with no cap and fails only at its end, as "smaller than its declared size". With a 1,024-byte cap, an entry carrying such a size inflated 8 MiB in full. decodePackStream now rejects a continuation byte whose multiplier is over the cap, since no size under the cap has a digit that high. The arithmetic was the same in decodePackBytes on main. Co-Authored-By: Claude Code --- .../workshop-backend/__tests__/git-codec.test.ts | 12 ++++++++++++ packages/workshop-backend/src/git-codec.ts | 6 ++++++ 2 files changed, 18 insertions(+) diff --git a/packages/workshop-backend/__tests__/git-codec.test.ts b/packages/workshop-backend/__tests__/git-codec.test.ts index e2ccc59f17..f149834c99 100644 --- a/packages/workshop-backend/__tests__/git-codec.test.ts +++ b/packages/workshop-backend/__tests__/git-codec.test.ts @@ -194,6 +194,18 @@ describe("pack decoding", () => { .rejects.toThrow(/exceeds the 64-byte limit/); }); + it("rejects an entry size longer than any size under the cap needs", async () => { + // Enough continuation bytes overflow the size to NaN, which no later comparison rejects. + let pack = concatBytes(await buildPackBytes([{ type: "blob", payload: new Uint8Array(3) }])); + let size = new Uint8Array(162).fill(0x80); + size[0] |= pack[12]; + size[161] = 0; + let body = concatBytes([pack.subarray(0, 12), size, pack.subarray(13, -20)]); + let padded = concatBytes([body, new Uint8Array(await crypto.subtle.digest("SHA-1", body))]); + await expect(decodePack(padded, { maxObjectSize: 64 })) + .rejects.toThrow(/entry size exceeds the 64-byte limit/); + }); + it("enforces the pack size cap", async () => { let pack = b64Bytes(PACK_NO_DELTA); await expect(decodePack(pack, { maxPackSize: pack.length - 1 })) diff --git a/packages/workshop-backend/src/git-codec.ts b/packages/workshop-backend/src/git-codec.ts index 58f553ffee..b74dc8b267 100644 --- a/packages/workshop-backend/src/git-codec.ts +++ b/packages/workshop-backend/src/git-codec.ts @@ -351,6 +351,12 @@ export async function* decodePackStream( let typeCode = (byte >> 4) & 0x07; let size = byte & 0x0f; for (let multiplier = 16; byte & 0x80; multiplier *= 128) { + // No size under the cap has a digit this high. Left to run, the multiplier overflows and + // the size becomes NaN, which every comparison below would let through. + if (multiplier > options.maxObjectSize) { + throw new Error( + `invalid packfile: entry size exceeds the ${options.maxObjectSize}-byte limit`); + } byte = await reader.byte(); size += (byte & 0x7f) * multiplier; } From 953cf79ee31de961791368f792c5198d16e1bb9f Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:12:55 -0500 Subject: [PATCH 02/24] Read a pack in full 64 KiB buffers PackReader asked for 64 KiB per read, but a BYOB read returns as soon as one of the source's chunks arrives. GitHub sent a vscode mount pack (38.6 MiB) as 3,358 sideband packets, 67% of them 8 KiB and 26% 16 KiB, and the gatekeeper forwards each as its own chunk. Through a Workers RPC hop in workerd, a default stream cut the same way took 3,330 reads. Reads now use readAtLeast, which waits for a full buffer: 590 reads for that pack. Each read that follows a storage write costs an implicit commit, so fewer reads mean fewer commits. At the end of the stream readAtLeast returns the short remainder with done unset, then done, on JS byte streams, native streams and RPC-carried ones alike, so more() needs no other change. Not measured end to end: locally the per-read cost was about 0.2 ms and the difference was inside the run-to-run spread. Co-Authored-By: Claude Code --- packages/workshop-backend/src/git-codec.ts | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/packages/workshop-backend/src/git-codec.ts b/packages/workshop-backend/src/git-codec.ts index b74dc8b267..fcb6cf9d47 100644 --- a/packages/workshop-backend/src/git-codec.ts +++ b/packages/workshop-backend/src/git-codec.ts @@ -428,7 +428,8 @@ const PACK_READ_SIZE = 64 << 10; // Workers RPC has carried it to the overseer (verified for the gatekeepers' pull-based stream // shape): a default reader gets 4 KiB chunks there, and each read after one of the caller's // storage writes costs an implicit commit (a TypeScript-size pack took 7.0 s of reads instead of -// 2.8 s, in workerd). +// 2.8 s, in workerd). Each read waits for a full buffer, because a BYOB read otherwise returns as +// soon as one of the source's chunks arrives, and GitHub sends a pack mostly in 8 KiB pieces. class PackReader { #reader: ReadableStreamBYOBReader; #maxSize: number; @@ -452,7 +453,7 @@ class PackReader { /** Whether any bytes remain, buffering at least one if so. */ async more(): Promise { while (this.#pos === this.#chunk.byteLength) { - let next = await this.#reader.read(new Uint8Array(PACK_READ_SIZE)); + let next = await this.#reader.readAtLeast(PACK_READ_SIZE, new Uint8Array(PACK_READ_SIZE)); if (next.done) return false; this.#received += next.value.byteLength; if (this.#received > this.#maxSize) { From 25fa791c45f26ce29ff035e0040f8371d2afac56 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:13:22 -0500 Subject: [PATCH 03/24] Skip the pull-routing hint a gatekeeper has already proven consumePack stores a pack's blobs as they arrive and its trees at the end. Each tree entry then added the gatekeeper to the blob's pullableFrom, rewriting a row that already had it in onRemote. Nothing reads that hint apart from the proof: pull routing takes the union of the two lists, the marking walk tests either, and no path removes an onRemote entry. #recordPullable now writes nothing when the gatekeeper is already in onRemote. Consuming a vscode mount pack (23,280 objects) makes 28,428 metadata writes instead of 47,153, and a TypeScript one (65,044 objects) 66,023 instead of 86,975. Co-Authored-By: Claude Code --- packages/workshop-backend/__tests__/git-cache.test.ts | 11 +++++++++++ packages/workshop-backend/src/git-cache.ts | 8 ++++++-- 2 files changed, 17 insertions(+), 2 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index 7eea12662b..7c59f36923 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -802,6 +802,17 @@ describe("consumePack", () => { expect(t.storage.gitObjectMetadata.get(GITLINK_TARGET)).toBeUndefined(); }); + it("records no pull-routing hint for a blob the pack itself delivered", async () => { + // The blobs are stored before the trees naming them, and proof is already a pull source. + let t = makeCache(); + await new GitCacheImpl(t.cache, G1).consumePack(byteStream(b64Bytes(PACK_OFS_DELTA))); + for (let oid of PACKED_OIDS.filter(o => fixture(o).type === "blob")) { + let meta = t.storage.gitObjectMetadata.get(oid)!; + expect(meta.onRemote).toStrictEqual([G1]); + expect(meta.pullableFrom).toStrictEqual([]); + } + }); + it("rejects corrupt input without storing its commits or trees", async () => { let t = makeCache(); let bytes = b64Bytes(PACK_OFS_DELTA).slice(); diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index a4f872967c..f700e04ba2 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -1158,10 +1158,14 @@ export class WorkspaceGitCache { return { meta, dirty }; } - // Records an assertion-grade pull-routing hint (advertisement or put-referent). + // Records an assertion-grade pull-routing hint (advertisement or put-referent). A gatekeeper + // already in `onRemote` needs none: proof routes pulls and bounds the marking walk by itself. + // (consumePack stores a pack's blobs before the trees naming them, so that is most referents.) #recordPullable(gatekeeperId: WorkpieceId, oid: GitOid, type: GitObjectType): void { let { meta, dirty } = this.#metaFor(gatekeeperId, oid, type, "asserted"); - if (addUnique(meta.pullableFrom, gatekeeperId) || dirty) { + let hinted = !meta.onRemote.includes(gatekeeperId) && + addUnique(meta.pullableFrom, gatekeeperId); + if (hinted || dirty) { this.storage.gitObjectMetadata.put(meta); } } From bba045a5418bc935efb7a49236917b73676c6b96 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:13:40 -0500 Subject: [PATCH 04/24] Leave an object the gatekeeper already stored untouched A pull sends no haves, so the remote answers with everything the commit needs whether or not the cache has it. A retried pull, or one for the next commit of a repository that is already mounted, is therefore mostly objects the same gatekeeper stored before, and each was deflated and written again along with its metadata row. #storeVerifiedObject now returns early for an object that is present, measured, and already proven for this gatekeeper. The check reads the metadata row first, so an object new to the cache pays nothing extra. Because consumePack keeps a failed pack's blobs, a pull that was cut off (the Overseer reset on CPU, for one) also resumes from them rather than repeating the same work. The same pack consumed a second time, local workerd on Durable Object SQLite, two runs each on a busy machine: - vscode (23,280 objects): 6.7-7.9 s -> 2.8-3.2 s, no writes - TypeScript (65,044 objects): 11.6-15.1 s -> 4.8-7.1 s, two writes (its two oversized trees are measured again) Co-Authored-By: Claude Code --- .../workshop-backend/__tests__/git-cache.test.ts | 16 +++++++++++++++- packages/workshop-backend/src/git-cache.ts | 10 ++++++++-- 2 files changed, 23 insertions(+), 3 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index 7c59f36923..a4809e54ad 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; import { deflate } from "pako"; import type { GitPullHints, GitOid } from "@gadgets/workshop-shared/gatekeeper"; import { READ_FILES_RESPONSE_BUDGET } from "@gadgets/workshop-shared/api"; @@ -813,6 +813,20 @@ describe("consumePack", () => { } }); + it("writes nothing when the same gatekeeper sends a pack again", async () => { + // A pull carries no `have`s, so a retry, or the next commit of a mounted repository, is + // mostly objects already stored. + let t = makeCache(); + let stub = new GitCacheImpl(t.cache, G1); + let first = await stub.consumePack(byteStream(b64Bytes(PACK_OFS_DELTA))); + let objectPuts = vi.spyOn(t.storage.gitObjects, "put"); + let metadataPuts = vi.spyOn(t.storage.gitObjectMetadata, "put"); + let again = await stub.consumePack(byteStream(b64Bytes(PACK_OFS_DELTA))); + expect(new Set(again)).toStrictEqual(new Set(first)); + expect(objectPuts).not.toHaveBeenCalled(); + expect(metadataPuts).not.toHaveBeenCalled(); + }); + it("rejects corrupt input without storing its commits or trees", async () => { let t = makeCache(); let bytes = b64Bytes(PACK_OFS_DELTA).slice(); diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index f700e04ba2..6749ff4831 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -1183,15 +1183,21 @@ export class WorkspaceGitCache { // The shared put()-equivalent store step (callers wrap in a transaction): an object over // MAX_GIT_OBJECT_SIZE is only measured, returning false. Anything else is stored, with proof of // possession and referent pull-routing rows recorded and pending-push marks propagated to the - // referents now that they are visible. + // referents now that they are visible. An object this gatekeeper already stored is left alone: + // a pull sends no `have`s, so a retried pull, or one for the next commit of a mounted + // repository, is mostly such objects. #storeVerifiedObject(gatekeeperId: WorkpieceId, oid: GitOid, { type, payload }: PackableObject) : boolean { if (payload.byteLength > MAX_GIT_OBJECT_SIZE) { this.#recordOversized(gatekeeperId, oid, type, payload.byteLength); return false; } - this.storage.gitObjects.put({ oid, data: encodeLooseObject(type, payload) }); let { meta } = this.#metaFor(gatekeeperId, oid, type, "measured"); + if (meta.size !== undefined && meta.onRemote.includes(gatekeeperId) && + this.hasLocalObject(oid)) { + return true; + } + this.storage.gitObjects.put({ oid, data: encodeLooseObject(type, payload) }); addUnique(meta.onRemote, gatekeeperId); meta.size = payload.byteLength; this.storage.gitObjectMetadata.put(meta); From 63553377547864bc1db9c769c9eddfdf81da08b7 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:37:31 -0500 Subject: [PATCH 05/24] Measure an oversized pack entry as it arrives An object over MAX_GIT_OBJECT_SIZE was held with the commits and trees and only measured once the whole pack had verified. A pull that ended early (a bad trailer, a stream error, the Overseer reset) recorded nothing for it, so the next read asked for the same object again and met the same end. A blob fetch sends no filter, so that is every large file a batch read touches. consumePack now records the size when the object arrives. The oid is the hash of the bytes in hand, the same grade of evidence as the small blobs already stored before the trailer, and it is what lets ensureGitObjects fail fast on the object afterwards. The object is still kept to the end of the pack as a possible delta base, in its own map, so the final loop stores only commits, trees and tags. Co-Authored-By: Claude Code --- .../__tests__/git-cache.test.ts | 12 ++++++ packages/workshop-backend/src/git-cache.ts | 38 ++++++++++--------- 2 files changed, 33 insertions(+), 17 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index a4809e54ad..364f992ade 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -869,6 +869,18 @@ describe("consumePack", () => { expect(meta.onRemote).toStrictEqual([G1]); }); + it("measures an oversized entry even when the pack then fails", async () => { + // The recorded size is what lets a later read fail fast instead of pulling the object again. + let t = makeCache(); + let big = new Uint8Array(MAX_GIT_OBJECT_SIZE + 5).fill(0x7a); + let pack = concatBytes(await buildPackBytes([{ type: "blob", payload: big }])); + pack[pack.length - 1] ^= 0xff; + await expect(new GitCacheImpl(t.cache, G1).consumePack(byteStream(pack))) + .rejects.toThrow(/trailer SHA-1 mismatch/); + expect(t.storage.gitObjectMetadata.get(await gitObjectOid("blob", big))?.size) + .toBe(big.byteLength); + }); + it("resolves a delta against an oversized base it declines to store", async () => { // How git packs a file similar to a large one (e.g. a second lockfile in a whole-tree blob // pull), hand-built: the large blob, then a ref-delta copying its first 16 bytes. diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index 6749ff4831..24156e2bcf 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -255,40 +255,44 @@ export class WorkspaceGitCache { * absent from the returned list, which is how a gitPull implementation notices). Returns the * stored oids. * - * The pack streams through: small blobs, the bulk of a checkout, are stored as they arrive. - * Everything else -- oversized blobs included, as a later delta may name one as its base -- is - * held until the whole pack has verified, then stored the same way, one object per - * transaction, with commits last: a commit's local presence is what lets `fetchCommit` mount - * it and skip ever pulling it again, so no commit is stored before every other held object is. - * A store that throws rolls back only its own object, so a failure partway can leave verified - * trees, the root tree included, with no commit: nothing treats those as mounted, and lazy - * reads fault around them. No await separates these stores, so they still reach disk - * together; one transaction around them all would also undo the earlier objects when a later - * one throws, but in production it nearly doubled a vscode-size mount's CPU. + * The pack streams through. Small blobs, the bulk of a checkout, are stored as they arrive. + * An oversized object is measured as it arrives, so a pull that is cut short still leaves the + * size for later reads to fail fast on, and is then kept to the end, as a later delta may name + * it as its base. Commits, trees and tags are held until the whole pack has verified, then + * stored, one object per transaction, with commits last: a commit's local presence is what + * lets `fetchCommit` mount it and skip ever pulling it again, so no commit is stored before + * every other held object is. A store that throws rolls back only its own object, so a + * failure partway can leave verified trees, the root tree included, with no commit: nothing + * treats those as mounted, and lazy reads fault around them. No await separates these stores, + * so they still reach disk together; one transaction around them all would also undo the + * earlier objects when a later one throws, but in production it nearly doubled a vscode-size + * mount's CPU. */ async consumePackFromGatekeeper(gatekeeperId: WorkpieceId, pack: ReadableStream) : Promise { let held = new Map(); + let oversized = new Map(); let stored: GitOid[] = []; let objects = decodePackStream(pack, { maxPackSize: MAX_GIT_PACK_BYTES, maxObjectSize: MAX_GIT_PACK_BYTES, - resolveBase: oid => held.get(oid) ?? this.readLocalObject(oid), + resolveBase: oid => held.get(oid) ?? oversized.get(oid) ?? this.readLocalObject(oid), }); for await (let { oid, ...object } of objects) { - if (object.type === "blob" && object.payload.byteLength <= MAX_GIT_OBJECT_SIZE) { - this.storage.transaction(() => this.#storeVerifiedObject(gatekeeperId, oid, object)); + if (object.type !== "blob" && object.payload.byteLength <= MAX_GIT_OBJECT_SIZE) { + held.set(oid, object); + } else if (this.storage.transaction( + () => this.#storeVerifiedObject(gatekeeperId, oid, object))) { stored.push(oid); } else { - held.set(oid, object); + oversized.set(oid, object); } } let commitsLast = [...held].toSorted(([, a], [, b]) => Number(a.type === "commit") - Number(b.type === "commit")); for (let [oid, object] of commitsLast) { - if (this.storage.transaction(() => this.#storeVerifiedObject(gatekeeperId, oid, object))) { - stored.push(oid); - } + this.storage.transaction(() => this.#storeVerifiedObject(gatekeeperId, oid, object)); + stored.push(oid); } return stored; } From 4b4e0e92857a055b68dfa75f9f4dd0970d17cf30 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:42:26 -0500 Subject: [PATCH 06/24] Store a streamed blob without a transaction of its own Each blob consumePack stores as it arrives sat in its own storage.transaction, and the metadata put inside opens a nested one for its index. A transaction is there to undo a store that fails after its first write, and what can fail there is parsing the objects the stored one names, in #referentEntries and the marking walk. A blob names none, and an oversized object is only measured, so both are now stored without the wrapper. Commits, trees and tags keep theirs. If the metadata write itself failed, the blob's row would stay with no metadata: a hash-verified object no gatekeeper can read, which the next pull completes. First pull in local workerd on Durable Object SQLite, best of three, on a machine that was busy for part of the "before" runs: - vscode (23,280 objects): 8.2 s -> 7.3 s - TypeScript (65,044 objects): 14.5 s -> 11.2 s Not measured in production, where the third commit of #656 found the nesting far more expensive than it is locally. Co-Authored-By: Claude Code --- packages/workshop-backend/src/git-cache.ts | 31 +++++++++++----------- 1 file changed, 16 insertions(+), 15 deletions(-) diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index 24156e2bcf..ecc0d0dcb8 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -255,18 +255,19 @@ export class WorkspaceGitCache { * absent from the returned list, which is how a gitPull implementation notices). Returns the * stored oids. * - * The pack streams through. Small blobs, the bulk of a checkout, are stored as they arrive. - * An oversized object is measured as it arrives, so a pull that is cut short still leaves the - * size for later reads to fail fast on, and is then kept to the end, as a later delta may name - * it as its base. Commits, trees and tags are held until the whole pack has verified, then - * stored, one object per transaction, with commits last: a commit's local presence is what - * lets `fetchCommit` mount it and skip ever pulling it again, so no commit is stored before - * every other held object is. A store that throws rolls back only its own object, so a - * failure partway can leave verified trees, the root tree included, with no commit: nothing - * treats those as mounted, and lazy reads fault around them. No await separates these stores, - * so they still reach disk together; one transaction around them all would also undo the - * earlier objects when a later one throws, but in production it nearly doubled a vscode-size - * mount's CPU. + * The pack streams through. Small blobs, the bulk of a checkout, are stored as they arrive, + * with no transaction around each: what can fail partway through a store is parsing the + * objects it names, and a blob names none. An oversized object is measured as it arrives, so + * a pull that is cut short still leaves the size for later reads to fail fast on, and is then + * kept to the end, as a later delta may name it as its base. Commits, trees and tags are held + * until the whole pack has verified, then stored, one object per transaction, with commits + * last: a commit's local presence is what lets `fetchCommit` mount it and skip ever pulling it + * again, so no commit is stored before every other held object is. A store that throws rolls + * back only its own object, so a failure partway can leave verified trees, the root tree + * included, with no commit: nothing treats those as mounted, and lazy reads fault around + * them. No await separates these stores, so they still reach disk together; one transaction + * around them all would also undo the earlier objects when a later one throws, but in + * production it nearly doubled a vscode-size mount's CPU. */ async consumePackFromGatekeeper(gatekeeperId: WorkpieceId, pack: ReadableStream) : Promise { @@ -281,8 +282,7 @@ export class WorkspaceGitCache { for await (let { oid, ...object } of objects) { if (object.type !== "blob" && object.payload.byteLength <= MAX_GIT_OBJECT_SIZE) { held.set(oid, object); - } else if (this.storage.transaction( - () => this.#storeVerifiedObject(gatekeeperId, oid, object))) { + } else if (this.#storeVerifiedObject(gatekeeperId, oid, object)) { stored.push(oid); } else { oversized.set(oid, object); @@ -1184,7 +1184,8 @@ export class WorkspaceGitCache { this.storage.gitObjectMetadata.put(meta); } - // The shared put()-equivalent store step (callers wrap in a transaction): an object over + // The shared put()-equivalent store step (callers wrap it in a transaction, unless the object + // is a blob or oversized, which leaves nothing to parse after the first write): an object over // MAX_GIT_OBJECT_SIZE is only measured, returning false. Anything else is stored, with proof of // possession and referent pull-routing rows recorded and pending-push marks propagated to the // referents now that they are visible. An object this gatekeeper already stored is left alone: From e1a748a4581764c98e037f9e41410863db04b73f Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:44:10 -0500 Subject: [PATCH 07/24] Deflate loose git objects at zlib's fastest level encodeLooseObject is the largest single CPU cost of storing a pack: about 2.0 s of 8 s for a vscode mount in local workerd, at zlib's default level 6. It now uses level 1, which is git's own default for loose objects (core.looseCompression). The output is still standard zlib, so isomorphic-git and records written before read the same way. This trades stored bytes for CPU. First pull in local workerd on Durable Object SQLite, best of three: - vscode (23,280 objects): 7.3 s -> 6.2 s, 38.8 MiB -> 43.3 MiB stored - TypeScript (65,044 objects): 11.2 s -> 10.4 s, 26.4 MiB -> 28.9 MiB Co-Authored-By: Claude Code --- packages/workshop-backend/src/git-codec.ts | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/packages/workshop-backend/src/git-codec.ts b/packages/workshop-backend/src/git-codec.ts index fcb6cf9d47..4d93a24a3c 100644 --- a/packages/workshop-backend/src/git-codec.ts +++ b/packages/workshop-backend/src/git-codec.ts @@ -19,7 +19,7 @@ // streams with no recorded lengths -- finding where one ends requires a streaming inflater that // reports unconsumed input, which DecompressionStream cannot do. -import { deflateSync, inflateSync } from "node:zlib"; +import { constants, deflateSync, inflateSync } from "node:zlib"; import { Inflate, deflate } from "pako"; import type { GitObjectType, GitOid } from "@gadgets/workshop-shared/gatekeeper"; @@ -77,10 +77,14 @@ export async function gitObjectOid(type: GitObjectType, payload: Uint8Array): Pr return toHex(new Uint8Array(digest)); } -/** Encodes a loose object record's `data` bytes from a type and headerless payload. */ +/** + * Encodes a loose object record's `data` bytes from a type and headerless payload, deflated at + * zlib's fastest level: git's own default for loose objects (`core.looseCompression`), and this + * deflate is the largest single cost of storing a pack. + */ export function encodeLooseObject(type: GitObjectType, payload: Uint8Array): Uint8Array { let header = ENCODER.encode(`${type} ${payload.byteLength}\0`); - return deflateSync(concatBytes([header, payload])); + return deflateSync(concatBytes([header, payload]), { level: constants.Z_BEST_SPEED }); } /** Decodes a loose object record's `data` bytes into its type and headerless payload. */ From 579aa3c06c2894528e0661c463c92d44757f804c Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:49:37 -0500 Subject: [PATCH 08/24] Inflate pack entries with native zlib The pack decoder kept pako because an entry's zlib stream has no recorded length, and finding its end takes an inflater that reports unconsumed input. node:zlib's inflateSync does that when asked for `info`: the returned engine's bytesWritten is the input the stream took. In workerd it tolerates the bytes that follow and enforces maxOutputLength. PackReader now keeps a window of unread bytes instead of one chunk, and inflate() hands inflateSync as much input as zlib itself could turn the declared size into. A stream that runs past that (another deflater's, or a padded one) is inflated again from twice as much; one that ends with the pack still short is truncated. The window drops consumed bytes into the hash when it runs out of room, doubles while one fill keeps reading, and never grows past what the fill asked for. The pako Inflate import, its internals interface and that cast are gone; the `info` result needs a cast of its own, since @types/node types it as a Buffer. The cost is memory for very large entries: the window has to hold an entry's whole compressed stream, and reads ahead up to its declared size, where pako needed one 64 KiB chunk of input. A mount pack's blobs are under 64 KiB, so there the window is about two read buffers except around a large tree. A result smaller than zlib's 16 KiB output chunk comes back as a view on the whole chunk, so it is copied before it is held. Decode only, no storage, local workerd, five alternating runs each: - vscode (23,280 objects): 1.57-1.93 s -> 0.87-1.27 s - TypeScript (65,044 objects): 2.35-3.40 s -> 1.58-2.29 s Both decoders yield the same objects, and the four new tests (declared size mismatch, corrupt data, an entry spanning many reads, a padded stream) pass on either. Co-Authored-By: Claude Code --- .../__tests__/git-codec.test.ts | 67 +++++++ packages/workshop-backend/src/git-codec.ts | 178 +++++++++--------- 2 files changed, 161 insertions(+), 84 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-codec.test.ts b/packages/workshop-backend/__tests__/git-codec.test.ts index f149834c99..63be0aac78 100644 --- a/packages/workshop-backend/__tests__/git-codec.test.ts +++ b/packages/workshop-backend/__tests__/git-codec.test.ts @@ -206,6 +206,73 @@ describe("pack decoding", () => { .rejects.toThrow(/entry size exceeds the 64-byte limit/); }); + // A one-blob pack with its entry (header byte and zlib stream) rewritten and the trailer redone. + async function packOfEntry(payload: Uint8Array, + rewrite: (entry: Uint8Array) => Uint8Array): Promise { + let pack = concatBytes(await buildPackBytes([{ type: "blob", payload }])); + let body = concatBytes([pack.subarray(0, 12), rewrite(pack.slice(12, -20))]); + return concatBytes([body, new Uint8Array(await crypto.subtle.digest("SHA-1", body))]); + } + + it("rejects an entry whose data is not the size it declares", async () => { + let payload = new TextEncoder().encode("ten bytes."); + let declaring = (size: number) => packOfEntry(payload, entry => { + entry[0] = (3 << 4) | size; + return entry; + }); + await expect(decodePack(await declaring(9))).rejects.toThrow(/larger than its declared/); + await expect(decodePack(await declaring(11))).rejects.toThrow(/smaller than its declared/); + expect((await decodePack(await declaring(10)))[0].payload).toStrictEqual(payload); + let oneByteDeclaredEmpty = await packOfEntry(new Uint8Array(1), entry => { + entry[0] = 3 << 4; + return entry; + }); + await expect(decodePack(oneByteDeclaredEmpty)).rejects.toThrow(/larger than its declared/); + }); + + it("rejects an entry whose data is not a zlib stream", async () => { + let pack = await packOfEntry(new Uint8Array(10), entry => { + entry[1] ^= 0xff; + return entry; + }); + await expect(decodePack(pack)).rejects.toThrow(/corrupt object data/); + }); + + it("decodes an entry that spans many reads, however the pack is chunked", async () => { + // Random bytes do not compress, so the entry's stream is three read buffers long. + let payload = new Uint8Array(200_000); + for (let pos = 0; pos < payload.length; pos += 50_000) { + crypto.getRandomValues(payload.subarray(pos, pos + 50_000)); + } + let after = new TextEncoder().encode("the entry after it"); + let pack = concatBytes(await buildPackBytes( + [{ type: "blob", payload }, { type: "blob", payload: after }])); + for (let step of [undefined, 1000]) { + let objects = await decodePack(pack, { step }); + expect(objects.map(o => o.payload)).toStrictEqual([payload, after]); + } + }); + + it("decodes an entry whose stream is longer than zlib would make it", async () => { + // Fifty empty stored blocks, then the payload in a final stored block and its Adler-32: a + // valid stream several times the length the first attempt allows for. + let payload = new TextEncoder().encode("padded out"); + let a = 1, b = 0; + for (let byte of payload) { + a = (a + byte) % 65521; + b = (b + a) % 65521; + } + let pack = await packOfEntry(payload, entry => concatBytes([ + entry.subarray(0, 1), + new Uint8Array([0x78, 0x01]), + ...Array.from({ length: 50 }, () => new Uint8Array([0, 0, 0, 0xff, 0xff])), + new Uint8Array([1, payload.length, 0, ~payload.length & 0xff, 0xff]), + payload, + new Uint8Array([b >> 8, b & 0xff, a >> 8, a & 0xff]), + ])); + expect((await decodePack(pack))[0].payload).toStrictEqual(payload); + }); + it("enforces the pack size cap", async () => { let pack = b64Bytes(PACK_NO_DELTA); await expect(decodePack(pack, { maxPackSize: pack.length - 1 })) diff --git a/packages/workshop-backend/src/git-codec.ts b/packages/workshop-backend/src/git-codec.ts index 4d93a24a3c..bd6618bad3 100644 --- a/packages/workshop-backend/src/git-codec.ts +++ b/packages/workshop-backend/src/git-codec.ts @@ -13,14 +13,15 @@ // commit *writes* (git-store.ts); tests cross-verify the two codecs over the same store. // // Everything here is pure computation over bytes (the pack decoder reads a stream): no storage, -// no RPC. Loose objects use workerd's native node:zlib: storing a mount pack deflates every object -// it carries, and pako's deflate takes about twice the CPU of the native one. The pack decoder -// uses pako (the same library isomorphic-git bundles) because pack entries are concatenated zlib -// streams with no recorded lengths -- finding where one ends requires a streaming inflater that -// reports unconsumed input, which DecompressionStream cannot do. +// no RPC. Storing a mount pack inflates every object it carries and deflates it again, and both +// use workerd's native node:zlib: pako's deflate takes about twice the CPU of the native one, and +// its inflate allocates some 100 KiB per stream. Pack entries are concatenated zlib streams with +// no recorded lengths, so the decoder has to learn where each one ends -- inflateSync reports the +// input it consumed when asked for `info`, which DecompressionStream cannot do. Only writing a +// pack (buildPackBytes, the push path) still uses pako, the library isomorphic-git bundles. import { constants, deflateSync, inflateSync } from "node:zlib"; -import { Inflate, deflate } from "pako"; +import { deflate } from "pako"; import type { GitObjectType, GitOid } from "@gadgets/workshop-shared/gatekeeper"; const ENCODER = new TextEncoder(); @@ -412,37 +413,36 @@ export async function* decodePackStream( if (trailer !== digest) throw new Error("invalid packfile: trailer SHA-1 mismatch"); } -// The Inflate internals this codec relies on beyond @types/pako's declarations, all stable pako -// API in practice (isomorphic-git's own pack parser relies on `strm.avail_in` the same way): -// `ended` flips when the zlib stream completes mid-input, and `strm.avail_in` is how many bytes -// of the last push() the stream did not consume -- together they locate the entry boundary. -interface InflateWithInternals { - ended: boolean; - err: number; - msg: string; - strm: { avail_in: number }; - onData: (chunk: Uint8Array) => void; - push(data: Uint8Array, flush: boolean): void; +// What inflateSync returns when asked for `info`: the output, and the engine whose +// `bytesWritten` is how much of the input the zlib stream took. (@types/node declares the result +// a Buffer whatever the options.) +interface InflatedEntry { + buffer: Uint8Array; + engine: { bytesWritten: number }; } const PACK_READ_SIZE = 64 << 10; -// decodePackStream's reads, in order, hashing every byte before `endBody()` (the trailer's SHA-1 -// input) a chunk at a time. Reads are BYOB, which a gatekeeper facet's pack stream supports once -// Workers RPC has carried it to the overseer (verified for the gatekeepers' pull-based stream -// shape): a default reader gets 4 KiB chunks there, and each read after one of the caller's -// storage writes costs an implicit commit (a TypeScript-size pack took 7.0 s of reads instead of -// 2.8 s, in workerd). Each read waits for a full buffer, because a BYOB read otherwise returns as -// soon as one of the source's chunks arrives, and GitHub sends a pack mostly in 8 KiB pieces. +// decodePackStream's reads, in order, from a window over the bytes that have arrived and are not +// yet consumed, hashing every byte before `endBody()` (the trailer's SHA-1 input) as it leaves +// the window. Reads are BYOB, which a gatekeeper facet's pack stream supports once Workers RPC +// has carried it to the overseer (verified for the gatekeepers' pull-based stream shape): a +// default reader gets 4 KiB chunks there, and each read after one of the caller's storage writes +// costs an implicit commit (a TypeScript-size pack took 7.0 s of reads instead of 2.8 s, in +// workerd). Each read waits for a full buffer, because a BYOB read otherwise returns as soon as +// one of the source's chunks arrives, and GitHub sends a pack mostly in 8 KiB pieces. class PackReader { #reader: ReadableStreamBYOBReader; #maxSize: number; #digest = new crypto.DigestStream("SHA-1"); #hash: WritableStreamDefaultWriter | undefined = this.#digest.getWriter(); - #chunk = new Uint8Array(0); + // Unread bytes are #window[#pos, #end); those before #pos are consumed and not yet hashed. + #window = new Uint8Array(0); #pos = 0; + #end = 0; #received = 0; + #ended = false; constructor(stream: ReadableStream, maxSize: number) { this.#reader = stream.getReader({ mode: "byob" }); @@ -451,87 +451,71 @@ class PackReader { /** The pack offset of the next unread byte. */ get offset(): number { - return this.#received - this.#chunk.byteLength + this.#pos; + return this.#received - this.#end + this.#pos; } /** Whether any bytes remain, buffering at least one if so. */ async more(): Promise { - while (this.#pos === this.#chunk.byteLength) { - let next = await this.#reader.readAtLeast(PACK_READ_SIZE, new Uint8Array(PACK_READ_SIZE)); - if (next.done) return false; - this.#received += next.value.byteLength; - if (this.#received > this.#maxSize) { - throw new Error(`packfile exceeds the ${this.#maxSize}-byte limit`); - } - await this.#hash?.write(this.#chunk); - this.#chunk = next.value; - this.#pos = 0; - } - return true; + return await this.#fill(1) > 0; } async byte(): Promise { - await this.#fill(); - return this.#chunk[this.#pos++]; + if (await this.#fill(1) === 0) throw new Error("invalid packfile: truncated"); + return this.#window[this.#pos++]; } async bytes(n: number): Promise { - let out = new Uint8Array(n); - for (let filled = 0; filled < n;) { - await this.#fill(); - let part = this.#chunk.subarray(this.#pos, this.#pos + n - filled); - out.set(part, filled); - filled += part.byteLength; - this.#pos += part.byteLength; - } - return out; + if (await this.#fill(n) < n) throw new Error("invalid packfile: truncated"); + return this.#window.slice(this.#pos, this.#pos += n); } // Inflates the zlib stream at the read position. `size` comes from the (untrusted) entry - // header; it was pre-checked against the object-size cap, and is enforced again here *during* - // inflation so a lying header cannot cause a larger allocation than it claimed. + // header; it was pre-checked against the object-size cap, and bounds the output here, so a + // lying header cannot cause a larger allocation than it claimed. + // + // inflateSync needs the whole stream in one piece and nothing records its length, so it is + // first given as much input as zlib itself could turn `size` bytes into. A stream that runs + // past that (another deflater's, or a padded one) is inflated again from twice as much. async inflate(size: number): Promise { - let inflator = new Inflate() as unknown as InflateWithInternals; - let chunks: Uint8Array[] = []; - let total = 0; - let overflow = false; - inflator.onData = (chunk: Uint8Array) => { - total += chunk.byteLength; - if (total > size) { - overflow = true; - // pako offers no abort; raising here unwinds through push() below. - throw new Error("pack entry exceeds declared size"); - } - chunks.push(chunk); - }; - - try { - while (!inflator.ended) { - await this.#fill(); - inflator.push(this.#chunk.subarray(this.#pos), false); - if (inflator.err) { - throw new Error(`invalid packfile: corrupt object data (${inflator.msg || inflator.err})`); + for (let want = size + (size >>> 12) + (size >>> 14) + 13;; want *= 2) { + let buffered = await this.#fill(want); + let input = this.#window.subarray(this.#pos, this.#pos + Math.min(buffered, want)); + let inflated: InflatedEntry; + try { + inflated = inflateSync(input, { info: true, maxOutputLength: size || 1 }) as + unknown as InflatedEntry; + } catch (err) { + let code = err instanceof Error && "code" in err ? err.code : undefined; + if (code === "Z_BUF_ERROR") { + if (buffered < want) throw new Error("invalid packfile: truncated", { cause: err }); + continue; } - this.#pos = this.#chunk.byteLength - inflator.strm.avail_in; + let detail = err instanceof Error ? err.message : String(err); + throw new Error( + code === "ERR_BUFFER_TOO_LARGE" + ? "invalid packfile: object larger than its declared size" + : `invalid packfile: corrupt object data (${detail})`, + { cause: err }); } - } catch (err) { - if (overflow) { - throw new Error("invalid packfile: object larger than its declared size", { cause: err }); + let { buffer, engine } = inflated; + if (buffer.byteLength !== size) { + // The output cap above cannot be 0, so a stream declared empty gets this far with a byte. + let how = buffer.byteLength < size ? "smaller" : "larger"; + throw new Error(`invalid packfile: object ${how} than its declared size`); } - throw err; + this.#pos += engine.bytesWritten; + // zlib inflates into 16 KiB chunks, and a result that fits one comes back as a view on + // the whole chunk: copied, so an object the caller holds retains only its own bytes. + return buffer.byteLength < buffer.buffer.byteLength + ? new Uint8Array(buffer) : new Uint8Array(buffer.buffer); } - - if (total !== size) { - throw new Error("invalid packfile: object smaller than its declared size"); - } - return concatBytes(chunks); } /** Ends the hashed body at the read position, returning its SHA-1 (hex). */ async endBody(): Promise { let hash = this.#hash!; this.#hash = undefined; - await hash.write(this.#chunk.subarray(0, this.#pos)); + await hash.write(this.#window.subarray(0, this.#pos)); await hash.close(); return toHex(new Uint8Array(await this.#digest.digest)); } @@ -542,8 +526,34 @@ class PackReader { this.#reader.cancel().catch(() => {}); } - async #fill(): Promise { - if (!await this.more()) throw new Error("invalid packfile: truncated"); + // Buffers at least `n` unread bytes, or all that remain, and returns how many are buffered. + async #fill(n: number): Promise { + while (this.#end - this.#pos < n && !this.#ended) { + let next = await this.#reader.readAtLeast(PACK_READ_SIZE, new Uint8Array(PACK_READ_SIZE)); + if (next.done) { + this.#ended = true; + break; + } + this.#received += next.value.byteLength; + if (this.#received > this.#maxSize) { + throw new Error(`packfile exceeds the ${this.#maxSize}-byte limit`); + } + if (this.#end + next.value.byteLength > this.#window.byteLength) { + // Out of room. What has been consumed goes to the hash and is dropped; the rest moves + // to a window that doubles while this fill keeps reading, and stops at what it asked for. + await this.#hash?.write(this.#window.subarray(0, this.#pos)); + let unread = this.#window.subarray(this.#pos, this.#end); + let window = new Uint8Array( + Math.min(2 * (unread.byteLength + next.value.byteLength), n + PACK_READ_SIZE)); + window.set(unread); + this.#window = window; + this.#end = unread.byteLength; + this.#pos = 0; + } + this.#window.set(next.value, this.#end); + this.#end += next.value.byteLength; + } + return this.#end - this.#pos; } } From 08146e029f78e600f904aac8eba01d13a52b79ad Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:52:10 -0500 Subject: [PATCH 09/24] Test a pack arriving over RPC, and pin the delta order rule Two things #656 depends on had no test. Every consumePack test builds a byte stream by hand, the one shape that is certain to support BYOB reads. Production gets a default stream that gatekeeper-github's demuxGitFetchResponse makes in another Worker and Workers RPC carries over; handed to the decoder directly, such a stream throws "This ReadableStream does not support BYOB reads". The new test loads a small Worker through the LOADER binding, which calls consumePack on a GitCacheImpl stub with a default pull stream, as a gatekeeper does. A ref-delta whose base comes later in the pack decoded on main and is now rejected. Git writes a base before its deltas, so nothing real should send one, but the rule was only stated in a comment. Co-Authored-By: Claude Code --- .../__tests__/git-cache.test.ts | 35 +++++++++++++++++++ .../__tests__/git-codec.test.ts | 22 ++++++++++-- 2 files changed, 54 insertions(+), 3 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index 364f992ade..19f5ca66a1 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -1,4 +1,5 @@ import { describe, expect, it, vi } from "vitest"; +import { env, type WorkerEntrypoint } from "cloudflare:workers"; import { deflate } from "pako"; import type { GitPullHints, GitOid } from "@gadgets/workshop-shared/gatekeeper"; import { READ_FILES_RESPONSE_BUDGET } from "@gadgets/workshop-shared/api"; @@ -802,6 +803,40 @@ describe("consumePack", () => { expect(t.storage.gitObjectMetadata.get(GITLINK_TARGET)).toBeUndefined(); }); + it("consumes a pack another Worker streams in, as a gatekeeper's arrives", async () => { + // What production sees is not the byte stream the other tests build: gatekeeper-github's + // demuxGitFetchResponse makes a default stream, and Workers RPC carries it to the overseer. + type PackSender = WorkerEntrypoint & { + send(cache: GitCacheImpl, pack: Uint8Array, piece: number): Promise; + }; + let t = makeCache(); + let sender = env.LOADER.get(null, () => ({ + compatibilityDate: "2026-02-01", + mainModule: "sender.js", + modules: { + "sender.js": ` + import { WorkerEntrypoint } from "cloudflare:workers"; + export class Sender extends WorkerEntrypoint { + async send(cache, pack, piece) { + let pos = 0; + return await cache.consumePack(new ReadableStream({ + pull(controller) { + if (pos < pack.byteLength) controller.enqueue(pack.slice(pos, pos += piece)); + else controller.close(); + }, + })); + } + }`, + }, + globalOutbound: null, + })).getEntrypoint("Sender"); + let stored = await sender.send(new GitCacheImpl(t.cache, G1), b64Bytes(PACK_OFS_DELTA), 100); + expect(new Set(stored)).toStrictEqual(new Set(PACKED_OIDS)); + for (let oid of PACKED_OIDS) { + expect(t.cache.readLocalObject(oid)).toStrictEqual(fixture(oid)); + } + }); + it("records no pull-routing hint for a blob the pack itself delivered", async () => { // The blobs are stored before the trees naming them, and proof is already a pull source. let t = makeCache(); diff --git a/packages/workshop-backend/__tests__/git-codec.test.ts b/packages/workshop-backend/__tests__/git-codec.test.ts index 63be0aac78..4d0220e4a8 100644 --- a/packages/workshop-backend/__tests__/git-codec.test.ts +++ b/packages/workshop-backend/__tests__/git-codec.test.ts @@ -206,11 +206,13 @@ describe("pack decoding", () => { .rejects.toThrow(/entry size exceeds the 64-byte limit/); }); - // A one-blob pack with its entry (header byte and zlib stream) rewritten and the trailer redone. - async function packOfEntry(payload: Uint8Array, - rewrite: (entry: Uint8Array) => Uint8Array): Promise { + // A one-blob pack with its entry (header byte and zlib stream) rewritten, declaring `count` + // entries, and the trailer redone. + async function packOfEntry(payload: Uint8Array, rewrite: (entry: Uint8Array) => Uint8Array, + count = 1): Promise { let pack = concatBytes(await buildPackBytes([{ type: "blob", payload }])); let body = concatBytes([pack.subarray(0, 12), rewrite(pack.slice(12, -20))]); + new DataView(body.buffer).setUint32(8, count); return concatBytes([body, new Uint8Array(await crypto.subtle.digest("SHA-1", body))]); } @@ -292,6 +294,20 @@ describe("pack decoding", () => { expect(cancelled).toBe(true); }); + it("rejects a ref-delta whose base comes later in the pack", async () => { + // The format allows it, but git writes a base before its deltas, and one pass cannot wait. + let base = new TextEncoder().encode("hello base content"); + let baseOid = await gitObjectOid("blob", base); + let delta = new Uint8Array([base.length, base.length, 0x91, 0, base.length]); + let pack = await packOfEntry(base, entry => concatBytes([ + new Uint8Array([(7 << 4) | delta.length]), + Uint8Array.from(baseOid.match(/../g)!, h => parseInt(h, 16)), + deflate(delta), + entry, + ]), 2); + await expect(decodePack(pack)).rejects.toThrow(`delta base ${baseOid} is unavailable`); + }); + it("fails a ref-delta whose base is nowhere, and resolves it via resolveBase", async () => { // Hand-build a one-entry thin pack: a ref-delta against an external base. let base = new TextEncoder().encode("hello base content"); From bb3aff86433b4f7307faf62fcb297fb177ca5711 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:52:30 -0500 Subject: [PATCH 10/24] Say what consumePack returns and leaves behind on failure The GitCache.consumePack doc called it exactly equivalent to a put() of each object. It differs in ways a gatekeeper can see: an oversized object is left out of the result where put() throws and, since #656, a delta has to follow its base and a pack that fails partway can leave some of its objects stored. The return value was not described at all. Co-Authored-By: Claude Code --- packages/workshop-shared/src/gatekeeper.ts | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/packages/workshop-shared/src/gatekeeper.ts b/packages/workshop-shared/src/gatekeeper.ts index f90c40fbec..09d5d196e8 100644 --- a/packages/workshop-shared/src/gatekeeper.ts +++ b/packages/workshop-shared/src/gatekeeper.ts @@ -1635,9 +1635,19 @@ export interface GitCache extends RpcTarget { buildPack(): Promise>; /** - * Consumes a standard git packfile and inserts all the objects within into the git cache. This - * is exactly equivalent to if the Gatekeeper decoded the packfile itself and `put()` each object - * into the cache. + * Consumes a standard git packfile, storing each object it holds as `put()` would: the oid is + * computed from the bytes, and the same proof of possession is recorded. Returns the oids of + * those objects that are now in the cache, which a `gitPull()` should check its requested + * oids against; an oid repeats if the pack repeats the object. + * + * Unlike a sequence of `put()`s: + * - An object over the cache's size limit is measured and left out of the result, where + * `put()` would throw. + * - A delta must come after its base, as in every pack `git upload-pack` sends. + * - The pack is stored as it streams, so when the call rejects (the pack proves invalid, or + * the stream fails) objects read before that may stay stored. Commits are stored last, once + * the whole pack has verified and the rest of it is stored, so a failed pack does not leave + * a commit without the trees that came with it. */ consumePack(pack: ReadableStream): Promise; From 2163ddb18ae85323dcbbab3d7145e75ecbbd59fe Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 14:18:23 -0500 Subject: [PATCH 11/24] Bound the oversized objects kept as delta bases consumePack keeps every object over MAX_GIT_OBJECT_SIZE inflated until the pack ends, because a later delta may name one as its base. A blob fetch sends no filter, so a batch read can carry any number of them, and the pack's 64 MiB cap bounds only their compressed size. Thirty 5 MiB files that compress well are 150 MiB against the Overseer's 128 MB. #656 kept them the same way, in `held`. They are now kept up to MAX_OVERSIZED_BASE_BYTES (32 MiB), the oldest let go first and the newest always kept, since git writes a delta soon after its base. A delta that names one already let go fails the pack with "delta base ... is unavailable". Every oversized object before it has been measured by then, so the next pull does not ask for them, and without them in the pack the blob that needed the base should arrive whole. That retry is reasoned, not exercised. Nothing here measures Worker memory; the test checks which base is still on hand once two blobs pass the budget. Co-Authored-By: Claude Code --- .../__tests__/git-cache.test.ts | 33 ++++++++++++++++ packages/workshop-backend/src/git-cache.ts | 38 ++++++++++++++----- 2 files changed, 61 insertions(+), 10 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index 19f5ca66a1..d5327318b2 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -9,6 +9,7 @@ import { GitCacheImpl, GitObjectTooLargeError, MAX_GIT_OBJECT_SIZE, + MAX_OVERSIZED_BASE_BYTES, WorkspaceGitCache, } from "../src/git-cache"; import { GitStore, blobOid } from "../src/git-store"; @@ -937,6 +938,38 @@ describe("consumePack", () => { expect(t.cache.readLocalObject(stored[0])!.payload).toStrictEqual(target); expect(t.storage.gitObjectMetadata.get(bigOid)!.size).toBe(big.byteLength); }); + + it("lets the oldest oversized base go once they pass the budget", async () => { + // Two blobs that pass it together, so only the second is still on hand for a delta. + let size = MAX_OVERSIZED_BASE_BYTES / 2 + 1; + let bigs = [1, 2].map(fill => new Uint8Array(size).fill(fill)); + let prefix = concatBytes(await buildPackBytes( + bigs.map(payload => ({ type: "blob" as const, payload })))).slice(0, -20); + new DataView(prefix.buffer).setUint32(8, 3); + // Delta: the base size as a varint, target size 16, then a 16-byte copy from offset 0. + let delta: number[] = []; + for (let rest = size; rest > 0; rest >>>= 7) { + delta.push(rest > 0x7f ? rest & 0x7f | 0x80 : rest); + } + delta.push(16, 0x90, 16); + let consumeWithDeltaOn = async (base: Uint8Array) => { + let oid = await gitObjectOid("blob", base); + let body = concatBytes([prefix, new Uint8Array([(7 << 4) | delta.length]), + Uint8Array.from(oid.match(/../g)!, h => parseInt(h, 16)), deflate(new Uint8Array(delta))]); + let pack = concatBytes([body, new Uint8Array(await crypto.subtle.digest("SHA-1", body))]); + let t = makeCache(); + return { t, oid, stored: new GitCacheImpl(t.cache, G1).consumePack(byteStream(pack)) }; + }; + + let recent = await consumeWithDeltaOn(bigs[1]); + let target = await gitObjectOid("blob", bigs[1].subarray(0, 16)); + expect(await recent.stored).toStrictEqual([target]); + + let dropped = await consumeWithDeltaOn(bigs[0]); + await expect(dropped.stored).rejects.toThrow(`delta base ${dropped.oid} is unavailable`); + // Both were measured before that, which is what keeps them out of the next pull. + expect(dropped.t.storage.gitObjectMetadata.get(dropped.oid)?.size).toBe(size); + }); }); // ======================================================================================= diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index ecc0d0dcb8..4134b56c99 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -86,6 +86,13 @@ export const MAX_GIT_OBJECT_SIZE = 1 << 20; */ export const MAX_GIT_PACK_BYTES = 64 << 20; +/** + * How many bytes of oversized objects `consumePack()` keeps at once as bases a later delta may + * name. They are never stored, and a blob fetch can carry any number of them inflated, which + * the pack's own size cap does not bound. + */ +export const MAX_OVERSIZED_BASE_BYTES = 32 << 20; + /** * Blob size fetched eagerly when pulling a worktree base: the pull requests the base commit with * `filterBlobSize: EAGER_BLOB_LIMIT`, so the commit, its full tree structure, and every blob @@ -259,20 +266,24 @@ export class WorkspaceGitCache { * with no transaction around each: what can fail partway through a store is parsing the * objects it names, and a blob names none. An oversized object is measured as it arrives, so * a pull that is cut short still leaves the size for later reads to fail fast on, and is then - * kept to the end, as a later delta may name it as its base. Commits, trees and tags are held - * until the whole pack has verified, then stored, one object per transaction, with commits - * last: a commit's local presence is what lets `fetchCommit` mount it and skip ever pulling it - * again, so no commit is stored before every other held object is. A store that throws rolls - * back only its own object, so a failure partway can leave verified trees, the root tree - * included, with no commit: nothing treats those as mounted, and lazy reads fault around - * them. No await separates these stores, so they still reach disk together; one transaction - * around them all would also undo the earlier objects when a later one throws, but in - * production it nearly doubled a vscode-size mount's CPU. + * kept as a base a later delta may name. Those are let go oldest first once they pass + * MAX_OVERSIZED_BASE_BYTES, git writing a delta soon after its base; a delta that does name + * one let go fails the pack, and the sizes already recorded keep those objects out of the + * next pull. Commits, trees and tags are held until the whole pack has verified, then stored, + * one object per transaction, with commits last: a commit's local presence is what lets + * `fetchCommit` mount it and skip ever pulling it again, so no commit is stored before every + * other held object is. A store that throws rolls back only its own object, so a failure + * partway can leave verified trees, the root tree included, with no commit: nothing treats + * those as mounted, and lazy reads fault around them. No await separates these stores, so + * they still reach disk together; one transaction around them all would also undo the + * earlier objects when a later one throws, but in production it nearly doubled a vscode-size + * mount's CPU. */ async consumePackFromGatekeeper(gatekeeperId: WorkpieceId, pack: ReadableStream) : Promise { let held = new Map(); let oversized = new Map(); + let oversizedBytes = 0; let stored: GitOid[] = []; let objects = decodePackStream(pack, { maxPackSize: MAX_GIT_PACK_BYTES, @@ -284,8 +295,15 @@ export class WorkspaceGitCache { held.set(oid, object); } else if (this.#storeVerifiedObject(gatekeeperId, oid, object)) { stored.push(oid); - } else { + } else if (!oversized.has(oid)) { oversized.set(oid, object); + oversizedBytes += object.payload.byteLength; + // Oldest out first, and never the one just added. + for (let [oldest, { payload }] of oversized) { + if (oversizedBytes <= MAX_OVERSIZED_BASE_BYTES || oldest === oid) break; + oversized.delete(oldest); + oversizedBytes -= payload.byteLength; + } } } let commitsLast = [...held].toSorted(([, a], [, b]) => From 3f75573174c9466eea3d515772fbc8d244ae5d17 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 14:31:25 -0500 Subject: [PATCH 12/24] Drop oversized bases only from a pack of blobs alone The budget on oversized delta bases failed a pack whenever a delta named a base it had let go, on the reasoning that the sizes recorded by then keep those objects out of the next request. That holds only for a request that named each blob. A mount wants a commit and sends no haves, so the remote answers a retry with the same pack, and it fails the same way: a valid pack could never be mounted. Devin's review caught this. Bases are now let go only from a pack that has carried nothing but blobs, which is what the answer to explicit blob wants looks like. A pack with a commit, tree or tag in it keeps every oversized base, as #656 does. For the blob case the retry now happens inside the same read. ensureGitObjects checked recorded sizes once, before its first pull; it now checks on every round, so after a pull fails on a dropped base the objects it measured surface as GitObjectTooLargeError, ensureBlobs takes them out of the batch, and the next pull brings the dependent blob whole. This also turns "every connection failed" into the too-large error whenever a failed pull was what measured the object. Two tests fail on the previous commit: a pack with a commit and two blobs past the budget is consumed, and a batch read whose pack fails on a dropped base completes in one ensureBlobs call with two pulls. Co-Authored-By: Claude Code --- .../__tests__/git-cache.test.ts | 92 ++++++++++++++----- packages/workshop-backend/src/git-cache.ts | 56 +++++------ 2 files changed, 99 insertions(+), 49 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index d5327318b2..e04dcb6f48 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -939,36 +939,84 @@ describe("consumePack", () => { expect(t.storage.gitObjectMetadata.get(bigOid)!.size).toBe(big.byteLength); }); - it("lets the oldest oversized base go once they pass the budget", async () => { - // Two blobs that pass it together, so only the second is still on hand for a delta. - let size = MAX_OVERSIZED_BASE_BYTES / 2 + 1; - let bigs = [1, 2].map(fill => new Uint8Array(size).fill(fill)); - let prefix = concatBytes(await buildPackBytes( - bigs.map(payload => ({ type: "blob" as const, payload })))).slice(0, -20); - new DataView(prefix.buffer).setUint32(8, 3); + // Two blobs that pass MAX_OVERSIZED_BASE_BYTES together, as pack entries. Deflating them is + // most of what these tests cost, so it is done once. + type BigBlobs = { bigs: Uint8Array[], oids: GitOid[], entries: Uint8Array }; + let twoBigBlobs: Promise | undefined; + function bigBlobs(): Promise { + return twoBigBlobs ??= (async () => { + let bigs = [1, 2].map( + fill => new Uint8Array(MAX_OVERSIZED_BASE_BYTES / 2 + 1).fill(fill)); + let pack = concatBytes(await buildPackBytes( + bigs.map(payload => ({ type: "blob" as const, payload })))); + let oids = await Promise.all(bigs.map(big => gitObjectOid("blob", big))); + return { bigs, oids, entries: pack.slice(12, -20) }; + })(); + } + + // A pack of any `leading` objects, the two big blobs, and a ref-delta copying the first 16 + // bytes of one of them (the `target`). + async function packOfTwoBigBlobs(deltaOn: 0 | 1, leading: PackableObject[] = []) { + let { bigs, oids, entries } = await bigBlobs(); + let size = bigs[0].byteLength; // Delta: the base size as a varint, target size 16, then a 16-byte copy from offset 0. let delta: number[] = []; for (let rest = size; rest > 0; rest >>>= 7) { delta.push(rest > 0x7f ? rest & 0x7f | 0x80 : rest); } delta.push(16, 0x90, 16); - let consumeWithDeltaOn = async (base: Uint8Array) => { - let oid = await gitObjectOid("blob", base); - let body = concatBytes([prefix, new Uint8Array([(7 << 4) | delta.length]), - Uint8Array.from(oid.match(/../g)!, h => parseInt(h, 16)), deflate(new Uint8Array(delta))]); - let pack = concatBytes([body, new Uint8Array(await crypto.subtle.digest("SHA-1", body))]); - let t = makeCache(); - return { t, oid, stored: new GitCacheImpl(t.cache, G1).consumePack(byteStream(pack)) }; - }; + let body = concatBytes([ + concatBytes(await buildPackBytes(leading)).slice(0, -20), + entries, + new Uint8Array([(7 << 4) | delta.length]), + Uint8Array.from(oids[deltaOn].match(/../g)!, h => parseInt(h, 16)), + deflate(new Uint8Array(delta)), + ]); + new DataView(body.buffer).setUint32(8, leading.length + 3); + let pack = concatBytes([body, new Uint8Array(await crypto.subtle.digest("SHA-1", body))]); + let target = bigs[deltaOn].subarray(0, 16); + return { pack, oids, size, target, targetOid: await gitObjectOid("blob", target) }; + } - let recent = await consumeWithDeltaOn(bigs[1]); - let target = await gitObjectOid("blob", bigs[1].subarray(0, 16)); - expect(await recent.stored).toStrictEqual([target]); + it("lets the oldest oversized base go from a pack of blobs past the budget", async () => { + // Only the second blob is still on hand by the time a delta names a base. + let recent = await packOfTwoBigBlobs(1); + expect(await new GitCacheImpl(makeCache().cache, G1).consumePack(byteStream(recent.pack))) + .toStrictEqual([recent.targetOid]); - let dropped = await consumeWithDeltaOn(bigs[0]); - await expect(dropped.stored).rejects.toThrow(`delta base ${dropped.oid} is unavailable`); - // Both were measured before that, which is what keeps them out of the next pull. - expect(dropped.t.storage.gitObjectMetadata.get(dropped.oid)?.size).toBe(size); + let t = makeCache(); + let dropped = await packOfTwoBigBlobs(0); + await expect(new GitCacheImpl(t.cache, G1).consumePack(byteStream(dropped.pack))) + .rejects.toThrow(`delta base ${dropped.oids[0]} is unavailable`); + expect(t.storage.gitObjectMetadata.get(dropped.oids[0])?.size).toBe(dropped.size); + }); + + it("keeps every oversized base in a pack that carries a commit", async () => { + // Such a pack answers a want for the commit and would be sent again unchanged, so failing + // it on a dropped base would fail the mount for good. + let t = makeCache(); + let commit = commitPayload(await gitObjectOid("tree", new Uint8Array(0)), [], "big files\n"); + let { pack, targetOid } = await packOfTwoBigBlobs(0, [{ type: "commit", payload: commit }]); + let stored = await new GitCacheImpl(t.cache, G1).consumePack(byteStream(pack)); + expect(stored).toStrictEqual([targetOid, await gitObjectOid("commit", commit)]); + }); + + it("reads a batch of blobs whose pack failed on a dropped base", async () => { + // The failed pack measured the two big blobs, so the same read asks again without them, and + // the blob that was a delta on one arrives whole. + let t = makeCache(); + let { pack, oids, target, targetOid } = await packOfTwoBigBlobs(0); + let wanted = [...oids, targetOid]; + await t.cache.putFromGatekeeper(G1, "tree", treePayload( + wanted.map((oid, i) => ({ mode: "100644", name: `file-${i}`, oid })))); + t.sources.set(G1, async asked => { + let answer = asked.includes(oids[0]) ? pack + : concatBytes(await buildPackBytes([{ type: "blob", payload: target }])); + await t.cache.consumePackFromGatekeeper(G1, byteStream(answer)); + }); + expect(await t.cache.ensureBlobs(wanted)).toStrictEqual(new Set(oids)); + expect(t.cache.readLocalObject(targetOid)?.payload).toStrictEqual(target); + expect(t.pulls.map(pull => pull.oids)).toStrictEqual([wanted, [targetOid]]); }); }); diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index 4134b56c99..bcbbab3bd2 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -87,9 +87,9 @@ export const MAX_GIT_OBJECT_SIZE = 1 << 20; export const MAX_GIT_PACK_BYTES = 64 << 20; /** - * How many bytes of oversized objects `consumePack()` keeps at once as bases a later delta may - * name. They are never stored, and a blob fetch can carry any number of them inflated, which - * the pack's own size cap does not bound. + * How many bytes of oversized blobs `consumePack()` keeps at once, from a pack of blobs alone, + * as bases a later delta may name. They are never stored, and a blob fetch can carry any number + * of them inflated, which the pack's own size cap does not bound. */ export const MAX_OVERSIZED_BASE_BYTES = 32 << 20; @@ -266,24 +266,26 @@ export class WorkspaceGitCache { * with no transaction around each: what can fail partway through a store is parsing the * objects it names, and a blob names none. An oversized object is measured as it arrives, so * a pull that is cut short still leaves the size for later reads to fail fast on, and is then - * kept as a base a later delta may name. Those are let go oldest first once they pass - * MAX_OVERSIZED_BASE_BYTES, git writing a delta soon after its base; a delta that does name - * one let go fails the pack, and the sizes already recorded keep those objects out of the - * next pull. Commits, trees and tags are held until the whole pack has verified, then stored, - * one object per transaction, with commits last: a commit's local presence is what lets - * `fetchCommit` mount it and skip ever pulling it again, so no commit is stored before every - * other held object is. A store that throws rolls back only its own object, so a failure - * partway can leave verified trees, the root tree included, with no commit: nothing treats - * those as mounted, and lazy reads fault around them. No await separates these stores, so - * they still reach disk together; one transaction around them all would also undo the - * earlier objects when a later one throws, but in production it nearly doubled a vscode-size - * mount's CPU. + * kept as a base a later delta may name. A pack of blobs alone lets them go, oldest first, + * once they pass MAX_OVERSIZED_BASE_BYTES: it answers a request that named each blob, so + * when a delta does name one let go and the pack fails, the sizes already recorded keep them + * out of the next request (see ensureGitObjects). A pack with a commit, tree or tag in it + * would be sent again unchanged, so it keeps every base. Commits, trees and tags are held + * until the whole pack has verified, then stored, one object per transaction, with commits + * last: a commit's local presence is what lets `fetchCommit` mount it and skip ever pulling it + * again, so no commit is stored before every other held object is. A store that throws rolls + * back only its own object, so a failure partway can leave verified trees, the root tree + * included, with no commit: nothing treats those as mounted, and lazy reads fault around + * them. No await separates these stores, so they still reach disk together; one transaction + * around them all would also undo the earlier objects when a later one throws, but in + * production it nearly doubled a vscode-size mount's CPU. */ async consumePackFromGatekeeper(gatekeeperId: WorkpieceId, pack: ReadableStream) : Promise { let held = new Map(); let oversized = new Map(); let oversizedBytes = 0; + let blobsOnly = true; let stored: GitOid[] = []; let objects = decodePackStream(pack, { maxPackSize: MAX_GIT_PACK_BYTES, @@ -291,6 +293,7 @@ export class WorkspaceGitCache { resolveBase: oid => held.get(oid) ?? oversized.get(oid) ?? this.readLocalObject(oid), }); for await (let { oid, ...object } of objects) { + blobsOnly &&= object.type === "blob"; if (object.type !== "blob" && object.payload.byteLength <= MAX_GIT_OBJECT_SIZE) { held.set(oid, object); } else if (this.#storeVerifiedObject(gatekeeperId, oid, object)) { @@ -300,7 +303,7 @@ export class WorkspaceGitCache { oversizedBytes += object.payload.byteLength; // Oldest out first, and never the one just added. for (let [oldest, { payload }] of oversized) { - if (oversizedBytes <= MAX_OVERSIZED_BASE_BYTES || oldest === oid) break; + if (!blobsOnly || oversizedBytes <= MAX_OVERSIZED_BASE_BYTES || oldest === oid) break; oversized.delete(oldest); oversizedBytes -= payload.byteLength; } @@ -357,23 +360,22 @@ export class WorkspaceGitCache { * a gatekeeper bug that wrongly omits a blob self-heals instead of wedging the file. */ async ensureGitObjects(oids: GitOid[], hints: GitPullHints): Promise { - let missing = [...new Set(oids)].filter(oid => !this.hasLocalObject(oid)); - if (missing.length === 0) return; - - // Fail fast on objects whose measured size already proves them unstorable. - for (let oid of missing) { - let size = this.storage.gitObjectMetadata.get(oid)?.size; - if (size !== undefined && size > MAX_GIT_OBJECT_SIZE) { - throw new GitObjectTooLargeError(oid, size); - } - } - + let missing = [...new Set(oids)]; let triedSources = new Map>(); let lastError: unknown; while (true) { missing = missing.filter(oid => !this.hasLocalObject(oid)); if (missing.length === 0) return; + // Fail fast on objects whose measured size already proves them unstorable. Checked on + // every round: a pull that failed may still have measured them (see consumePack). + for (let oid of missing) { + let size = this.storage.gitObjectMetadata.get(oid)?.size; + if (size !== undefined && size > MAX_GIT_OBJECT_SIZE) { + throw new GitObjectTooLargeError(oid, size); + } + } + // Group the still-missing objects by each one's next untried recorded source. let groups = new Map(); for (let oid of missing) { From 0f45e11cb67e1d0570fa8ad0eaf42cfcd454d439 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 14:39:39 -0500 Subject: [PATCH 13/24] Drop only the oversized bases a pull asked for by name The last commit let oversized bases go only from a pack that had carried nothing but blobs, taking that to mean the request named each blob. It inferred the request from the pack, and the inference fails when blobs come before the commit: a base is let go before the commit shows the pack is a mount's, a later delta names it, and the retry gets the same pack. Devin's review caught this too. consumePack is now told what the pull asked for. #pullGitObjects mints the stub it hands to gitPull() with the pull's oids, and only an oversized object among them can be let go. That is safe for exactly the reason the budget needs: a pull that names an oversized object can only end in GitObjectTooLargeError for it, and once measured it is left out of the next request, so the pack that failed is not the one the retry gets. An object the pull did not name, such as anything a commit's traversal brought, is never let go, in whatever order the pack carries it. A stub minted for anything else names nothing and lets nothing go. The test for the mount case now puts the two blobs before the commit, and fails on the last commit. The line in #pullGitObjects that passes the oids has no test of its own. Co-Authored-By: Claude Code --- .../__tests__/git-cache.test.ts | 34 +++++++----- packages/workshop-backend/src/git-cache.ts | 55 ++++++++++--------- packages/workshop-backend/src/overseer.ts | 4 +- 3 files changed, 53 insertions(+), 40 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index e04dcb6f48..660d8b6a77 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -954,9 +954,9 @@ describe("consumePack", () => { })(); } - // A pack of any `leading` objects, the two big blobs, and a ref-delta copying the first 16 - // bytes of one of them (the `target`). - async function packOfTwoBigBlobs(deltaOn: 0 | 1, leading: PackableObject[] = []) { + // A pack of the two big blobs, a ref-delta copying the first 16 bytes of one of them (the + // `target`), and then any `trailing` objects. + async function packOfTwoBigBlobs(deltaOn: 0 | 1, trailing: PackableObject[] = []) { let { bigs, oids, entries } = await bigBlobs(); let size = bigs[0].byteLength; // Delta: the base size as a varint, target size 16, then a 16-byte copy from offset 0. @@ -965,40 +965,46 @@ describe("consumePack", () => { delta.push(rest > 0x7f ? rest & 0x7f | 0x80 : rest); } delta.push(16, 0x90, 16); + let after = concatBytes(await buildPackBytes(trailing)); let body = concatBytes([ - concatBytes(await buildPackBytes(leading)).slice(0, -20), + after.subarray(0, 12), entries, new Uint8Array([(7 << 4) | delta.length]), Uint8Array.from(oids[deltaOn].match(/../g)!, h => parseInt(h, 16)), deflate(new Uint8Array(delta)), + after.subarray(12, -20), ]); - new DataView(body.buffer).setUint32(8, leading.length + 3); + new DataView(body.buffer).setUint32(8, trailing.length + 3); let pack = concatBytes([body, new Uint8Array(await crypto.subtle.digest("SHA-1", body))]); let target = bigs[deltaOn].subarray(0, 16); return { pack, oids, size, target, targetOid: await gitObjectOid("blob", target) }; } - it("lets the oldest oversized base go from a pack of blobs past the budget", async () => { + it("lets go of the oldest oversized base the pull asked for, past the budget", async () => { // Only the second blob is still on hand by the time a delta names a base. let recent = await packOfTwoBigBlobs(1); - expect(await new GitCacheImpl(makeCache().cache, G1).consumePack(byteStream(recent.pack))) + let consume = (t: TestCache, pack: Uint8Array, asked: GitOid[]) => + new GitCacheImpl(t.cache, G1, undefined, asked).consumePack(byteStream(pack)); + expect(await consume(makeCache(), recent.pack, recent.oids)) .toStrictEqual([recent.targetOid]); let t = makeCache(); let dropped = await packOfTwoBigBlobs(0); - await expect(new GitCacheImpl(t.cache, G1).consumePack(byteStream(dropped.pack))) + await expect(consume(t, dropped.pack, dropped.oids)) .rejects.toThrow(`delta base ${dropped.oids[0]} is unavailable`); expect(t.storage.gitObjectMetadata.get(dropped.oids[0])?.size).toBe(dropped.size); }); - it("keeps every oversized base in a pack that carries a commit", async () => { - // Such a pack answers a want for the commit and would be sent again unchanged, so failing - // it on a dropped base would fail the mount for good. + it("keeps an oversized base the pull did not ask for, wherever the commit comes", async () => { + // A mount asks for a commit and would be sent this pack again unchanged, so failing it on + // a dropped base would fail the mount for good. Here the blobs even precede the commit. let t = makeCache(); let commit = commitPayload(await gitObjectOid("tree", new Uint8Array(0)), [], "big files\n"); + let commitOid = await gitObjectOid("commit", commit); let { pack, targetOid } = await packOfTwoBigBlobs(0, [{ type: "commit", payload: commit }]); - let stored = await new GitCacheImpl(t.cache, G1).consumePack(byteStream(pack)); - expect(stored).toStrictEqual([targetOid, await gitObjectOid("commit", commit)]); + let stored = await new GitCacheImpl(t.cache, G1, undefined, [commitOid]) + .consumePack(byteStream(pack)); + expect(stored).toStrictEqual([targetOid, commitOid]); }); it("reads a batch of blobs whose pack failed on a dropped base", async () => { @@ -1012,7 +1018,7 @@ describe("consumePack", () => { t.sources.set(G1, async asked => { let answer = asked.includes(oids[0]) ? pack : concatBytes(await buildPackBytes([{ type: "blob", payload: target }])); - await t.cache.consumePackFromGatekeeper(G1, byteStream(answer)); + await t.cache.consumePackFromGatekeeper(G1, byteStream(answer), asked); }); expect(await t.cache.ensureBlobs(wanted)).toStrictEqual(new Set(oids)); expect(t.cache.readLocalObject(targetOid)?.payload).toStrictEqual(target); diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index bcbbab3bd2..e4ffb65e86 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -87,9 +87,10 @@ export const MAX_GIT_OBJECT_SIZE = 1 << 20; export const MAX_GIT_PACK_BYTES = 64 << 20; /** - * How many bytes of oversized blobs `consumePack()` keeps at once, from a pack of blobs alone, - * as bases a later delta may name. They are never stored, and a blob fetch can carry any number - * of them inflated, which the pack's own size cap does not bound. + * How many bytes of oversized objects `consumePack()` keeps at once as bases a later delta may + * name, before it lets go of those the pull asked for by name. They are never stored, and a + * blob fetch can carry any number of them inflated, which the pack's own size cap does not + * bound. */ export const MAX_OVERSIZED_BASE_BYTES = 32 << 20; @@ -266,26 +267,28 @@ export class WorkspaceGitCache { * with no transaction around each: what can fail partway through a store is parsing the * objects it names, and a blob names none. An oversized object is measured as it arrives, so * a pull that is cut short still leaves the size for later reads to fail fast on, and is then - * kept as a base a later delta may name. A pack of blobs alone lets them go, oldest first, - * once they pass MAX_OVERSIZED_BASE_BYTES: it answers a request that named each blob, so - * when a delta does name one let go and the pack fails, the sizes already recorded keep them - * out of the next request (see ensureGitObjects). A pack with a commit, tree or tag in it - * would be sent again unchanged, so it keeps every base. Commits, trees and tags are held - * until the whole pack has verified, then stored, one object per transaction, with commits - * last: a commit's local presence is what lets `fetchCommit` mount it and skip ever pulling it - * again, so no commit is stored before every other held object is. A store that throws rolls - * back only its own object, so a failure partway can leave verified trees, the root tree - * included, with no commit: nothing treats those as mounted, and lazy reads fault around - * them. No await separates these stores, so they still reach disk together; one transaction - * around them all would also undo the earlier objects when a later one throws, but in - * production it nearly doubled a vscode-size mount's CPU. + * kept as a base a later delta may name. Past MAX_OVERSIZED_BASE_BYTES of them, those the + * pull asked for by name (`asked`, from the stub `gitPull()` was handed) are let go, oldest + * first. Such a pull can only end in GitObjectTooLargeError for that object, and once + * measured it is left out of the next request, so a pack that fails because a delta named it + * is not the pack the retry gets (see ensureGitObjects). An object the pull did not ask for + * -- whatever a commit's traversal brought -- would arrive again in the same pack, so it is + * kept. Commits, trees and tags are held until the whole pack has verified, then stored, one + * object per transaction, with commits last: a commit's local presence is what lets + * `fetchCommit` mount it and skip ever pulling it again, so no commit is stored before every + * other held object is. A store that throws rolls back only its own object, so a failure + * partway can leave verified trees, the root tree included, with no commit: nothing treats + * those as mounted, and lazy reads fault around them. No await separates these stores, so + * they still reach disk together; one transaction around them all would also undo the + * earlier objects when a later one throws, but in production it nearly doubled a vscode-size + * mount's CPU. */ - async consumePackFromGatekeeper(gatekeeperId: WorkpieceId, pack: ReadableStream) - : Promise { + async consumePackFromGatekeeper(gatekeeperId: WorkpieceId, pack: ReadableStream, + asked: readonly GitOid[] = []): Promise { let held = new Map(); let oversized = new Map(); let oversizedBytes = 0; - let blobsOnly = true; + let droppable = new Set(asked); let stored: GitOid[] = []; let objects = decodePackStream(pack, { maxPackSize: MAX_GIT_PACK_BYTES, @@ -293,7 +296,6 @@ export class WorkspaceGitCache { resolveBase: oid => held.get(oid) ?? oversized.get(oid) ?? this.readLocalObject(oid), }); for await (let { oid, ...object } of objects) { - blobsOnly &&= object.type === "blob"; if (object.type !== "blob" && object.payload.byteLength <= MAX_GIT_OBJECT_SIZE) { held.set(oid, object); } else if (this.#storeVerifiedObject(gatekeeperId, oid, object)) { @@ -301,9 +303,10 @@ export class WorkspaceGitCache { } else if (!oversized.has(oid)) { oversized.set(oid, object); oversizedBytes += object.payload.byteLength; - // Oldest out first, and never the one just added. + // Oldest out first: never the one just added, nor one the pull did not ask for. for (let [oldest, { payload }] of oversized) { - if (!blobsOnly || oversizedBytes <= MAX_OVERSIZED_BASE_BYTES || oldest === oid) break; + if (oversizedBytes <= MAX_OVERSIZED_BASE_BYTES || oldest === oid) break; + if (!droppable.has(oldest)) continue; oversized.delete(oldest); oversizedBytes -= payload.byteLength; } @@ -1249,12 +1252,14 @@ export class WorkspaceGitCache { * (put/advertise record this gatekeeper as the source) and the read view. The overseer * additionally binds the stub passed to `applyAction()` to the applying action, which is what * makes `buildPack()` available; session-scoped stubs (from - * `ObservationAuthorizer.getGitCache()`) have no action and `buildPack()` throws. + * `ObservationAuthorizer.getGitCache()`) have no action and `buildPack()` throws. The stub + * passed to `gitPull()` carries the oids that pull asked for, which `consumePack()` needs to + * bound what it keeps in memory. */ @validateRpc() export class GitCacheImpl extends RpcTarget implements GitCache { constructor(private cache: WorkspaceGitCache, private gatekeeperId: WorkpieceId, - private actionId?: number) { + private actionId?: number, private pulling: readonly GitOid[] = []) { super(); } @@ -1291,7 +1296,7 @@ export class GitCacheImpl extends RpcTarget implements GitCache { } async consumePack(pack: ReadableStream): Promise { - return this.cache.consumePackFromGatekeeper(this.gatekeeperId, pack); + return this.cache.consumePackFromGatekeeper(this.gatekeeperId, pack, this.pulling); } async isAncestor(ancestor: GitOid, descendant: GitOid): Promise { diff --git a/packages/workshop-backend/src/overseer.ts b/packages/workshop-backend/src/overseer.ts index baf06fb0ee..9c5a25d093 100644 --- a/packages/workshop-backend/src/overseer.ts +++ b/packages/workshop-backend/src/overseer.ts @@ -4261,7 +4261,9 @@ class OverseerImpl implements AgentHooks { // call, which the pull driver treats as this source failing. let facet = this.getGatekeeperFacet(gatekeeperId) as unknown as Fetcher & Required, "gitPull">>>; - await facet.gitPull(oids, new GitCacheImpl(this.gitCache, gatekeeperId), hints); + // The stub knows what this pull asks for: see consumePackFromGatekeeper on oversized bases. + await facet.gitPull( + oids, new GitCacheImpl(this.gitCache, gatekeeperId, undefined, oids), hints); } // Apply a single pending action: invoke the gatekeeper, mark it approved, and persist (the put From b3dd592ddcc1929077406c821571b0bef16be5bd Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:26:28 -0500 Subject: [PATCH 14/24] Report why a git fetch failed, not the disconnect the cache saw The pack reaches consumePack() as a stream over RPC. When the gatekeeper's end of it fails (the transfer limit, an error line from the server, a truncated response, the fetch timing out), the overseer's reader learns only that the stream ended early. consumePack() rejects with "ReadableStream received over RPC disconnected prematurely", and gitPull() passed that on. Mounting torvalds/linux showed it. Its depth-1 mount pack is 190.9 MiB against the 64 MiB transfer limit, and the agent was told the connection had dropped, so it tried again. demuxGitFetchResponse now reports what its stream failed with, and pullGitObjectsIntoCache throws that in place of consumePack()'s rejection. Only a failure of the stream's own source is kept, so a rejection that is consumePack()'s own, an invalid pack for one, still comes through. Checked two ways besides the tests. In workerd, a default stream whose pull throws reaches a reader in another Worker as that disconnect error. And this transport run against github.com for that commit throws "git fetch response exceeded the 67108864-byte transfer limit" after 67,067,903 pack bytes. Co-Authored-By: Claude Code --- .../__tests__/git-transport.test.ts | 25 ++++++++++++++++++ .../gatekeeper-github/src/git-transport.ts | 26 ++++++++++++++++--- 2 files changed, 47 insertions(+), 4 deletions(-) diff --git a/packages/gatekeeper-github/__tests__/git-transport.test.ts b/packages/gatekeeper-github/__tests__/git-transport.test.ts index 604b6566f0..7267e9c350 100644 --- a/packages/gatekeeper-github/__tests__/git-transport.test.ts +++ b/packages/gatekeeper-github/__tests__/git-transport.test.ts @@ -344,6 +344,31 @@ describe("pullGitObjectsIntoCache", () => { expect(requestLines(requests[0])).toContain(`want ${oid(1)}`); }); + it("reports why the fetch failed, not the disconnect the cache saw", async () => { + // Across RPC the overseer's reader learns only that the stream ended early. + const cache = { + async consumePack(pack: ReadableStream): Promise { + await collect(pack).catch(() => {}); + throw new Error("ReadableStream received over RPC disconnected prematurely."); + }, + }; + await expect(pullGitObjectsIntoCache( + fakeFetch([], [encodePktLine("ERR upload-pack: not our ref")]), [oid(1)], HINTS, cache)) + .rejects.toThrow("git fetch failed: upload-pack: not our ref"); + }); + + it("reports the cache's own failure when the fetch was sound", async () => { + const cache = { + async consumePack(pack: ReadableStream): Promise { + await collect(pack); + throw new Error("invalid packfile: bad magic"); + }, + }; + await expect(pullGitObjectsIntoCache( + fakeFetch([], packfileResponse()), [oid(1)], HINTS, cache)) + .rejects.toThrow("invalid packfile: bad magic"); + }); + it("throws when a requested non-blob object is missing from the stored list", async () => { const cache = fakeCache([oid(1)]); await expect(pullGitObjectsIntoCache( diff --git a/packages/gatekeeper-github/src/git-transport.ts b/packages/gatekeeper-github/src/git-transport.ts index a441f57047..6d63fad98d 100644 --- a/packages/gatekeeper-github/src/git-transport.ts +++ b/packages/gatekeeper-github/src/git-transport.ts @@ -266,16 +266,24 @@ const BAND_ERROR = 3; * discarded, band 3 = fatal server error). An `ERR` pkt or band-3 message fails the stream with * the server's message; `maxBytes` bounds the raw body (see MAX_GIT_FETCH_BYTES); a response * that ends without a flush-pkt, or without ever reaching a packfile section, is an error -- - * a truncated pack must never look like a short success. + * a truncated pack must never look like a short success. `onFailure` is told what the stream + * failed with, which a reader on the far side of an RPC hop does not learn. */ export function demuxGitFetchResponse( body: ReadableStream, maxBytes: number, + onFailure?: (error: unknown) => void, ): ReadableStream { let iterator = demuxPackData(body, maxBytes); return new ReadableStream({ async pull(controller) { - let next = await iterator.next(); + let next: IteratorResult; + try { + next = await iterator.next(); + } catch (error) { + onFailure?.(error); + throw error; + } if (next.done) controller.close(); else controller.enqueue(next.value); }, @@ -381,8 +389,18 @@ export async function pullGitObjectsIntoCache( if (response.body === null) { throw new Error("git fetch failed: response had no body"); } - let stored = new Set(await cache.consumePack( - demuxGitFetchResponse(response.body, MAX_GIT_FETCH_BYTES))); + // When the pack stream itself fails, the overseer's reader sees only that it ended early, and + // consumePack() rejects with that. What the stream failed with says why: the transfer limit, + // the server's own error, a truncated response. + let failure: unknown; + let pack = demuxGitFetchResponse( + response.body, MAX_GIT_FETCH_BYTES, error => { failure = error; }); + let stored: Set; + try { + stored = new Set(await cache.consumePack(pack)); + } catch (error) { + throw failure ?? error; + } let missing = oids.filter(oid => !stored.has(oid)); if (missing.length > 0 && !(hints.type === "blob" && hints.filterBlobSize !== undefined)) { throw new Error(`git fetch did not provide the requested object${ From 19765037157378912ee7d4ea1b74623aea3c6823 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:37:43 -0500 Subject: [PATCH 15/24] Give a pack's objects their own size cap, and raise the pack's MAX_GIT_PACK_BYTES did two jobs: it capped the pack, and through maxObjectSize it capped each object inflated from one. The 64 MiB came from gatekeeper-context's artifact sync, where it bounds a repository loaded into memory, and consumePack matched it while it still read the whole pack into memory before decoding. Since #656 the pack streams through, so its cap no longer bounds memory, only how much one pull may download, decode and store. The per-object cap still is a memory bound. They are now separate: MAX_GIT_PACK_OBJECT_SIZE stays at 64 MiB, and MAX_GIT_PACK_BYTES goes to 256 MiB, the smallest round figure that admits a depth-1 mount of torvalds/linux (190.8 MiB, 98,591 objects). That pack through consumePack, local workerd on Durable Object SQLite: every object stored (222.8 MiB of rows), the tree listing all 102,334 paths, in 44 s fed from a local stream and 70 s through an RPC hop in GitHub-sized pieces. In production it would run past the 30 s CPU limit, which this does not touch. Co-Authored-By: Claude Code --- .../__tests__/git-cache.test.ts | 17 +++++++++++++++++ packages/workshop-backend/src/git-cache.ts | 15 +++++++++++---- 2 files changed, 28 insertions(+), 4 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index 660d8b6a77..24ca252357 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -9,6 +9,7 @@ import { GitCacheImpl, GitObjectTooLargeError, MAX_GIT_OBJECT_SIZE, + MAX_GIT_PACK_OBJECT_SIZE, MAX_OVERSIZED_BASE_BYTES, WorkspaceGitCache, } from "../src/git-cache"; @@ -905,6 +906,22 @@ describe("consumePack", () => { expect(meta.onRemote).toStrictEqual([G1]); }); + it("caps one object at MAX_GIT_PACK_OBJECT_SIZE, however large a pack may be", async () => { + // An entry that only declares a size over the cap: it is refused before anything inflates. + let size = MAX_GIT_PACK_OBJECT_SIZE + 1; + let entry = [0x80 | (3 << 4) | (size & 0x0f)]; + for (let rest = Math.floor(size / 16); rest > 0; rest >>>= 7) { + entry.push(rest > 0x7f ? rest & 0x7f | 0x80 : rest); + } + let body = concatBytes([ + concatBytes(await buildPackBytes([])).subarray(0, 12), new Uint8Array(entry), + deflate(new Uint8Array(1))]); + new DataView(body.buffer).setUint32(8, 1); + let pack = concatBytes([body, new Uint8Array(await crypto.subtle.digest("SHA-1", body))]); + await expect(new GitCacheImpl(makeCache().cache, G1).consumePack(byteStream(pack))) + .rejects.toThrow(`entry of ${size} bytes exceeds the ${size - 1}-byte limit`); + }); + it("measures an oversized entry even when the pack then fails", async () => { // The recorded size is what lets a later read fail fast instead of pulling the object again. let t = makeCache(); diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index e4ffb65e86..ee2487ad59 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -81,10 +81,17 @@ export const MAX_GIT_OBJECT_SIZE = 1 << 20; /** * Maximum byte size of a packfile accepted by `consumePack()` (matching the transfer-size - * limiter gatekeepers are expected to apply to fetch bodies), and the hard per-object - * inflation bound while decoding one. + * limiter gatekeepers are expected to apply to fetch bodies). A pack streams through, so this + * bounds what one pull may download, decode and store, not memory. */ -export const MAX_GIT_PACK_BYTES = 64 << 20; +export const MAX_GIT_PACK_BYTES = 256 << 20; + +/** + * Hard cap on any one object inflated while decoding a pack. This one is a memory bound, so it + * does not grow with MAX_GIT_PACK_BYTES. It sits well above MAX_GIT_OBJECT_SIZE so that an + * object too large to store can still be measured, and serve as the base of a delta. + */ +export const MAX_GIT_PACK_OBJECT_SIZE = 64 << 20; /** * How many bytes of oversized objects `consumePack()` keeps at once as bases a later delta may @@ -292,7 +299,7 @@ export class WorkspaceGitCache { let stored: GitOid[] = []; let objects = decodePackStream(pack, { maxPackSize: MAX_GIT_PACK_BYTES, - maxObjectSize: MAX_GIT_PACK_BYTES, + maxObjectSize: MAX_GIT_PACK_OBJECT_SIZE, resolveBase: oid => held.get(oid) ?? oversized.get(oid) ?? this.readLocalObject(oid), }); for await (let { oid, ...object } of objects) { From 60c5292e6626f908c83c715cfc17bc7126c5506d Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:37:45 -0500 Subject: [PATCH 16/24] Raise the git fetch transfer limit to match the pack cap MAX_GIT_FETCH_BYTES bounds the raw body of an upload-pack fetch and is kept equal to the overseer's cap on the pack inside it, which is now 256 MiB. The gatekeeper holds none of the body: it strips the framing and streams the pack on, so the limit bounds the work of one pull, not memory here. A depth-1 mount of torvalds/linux is a 190.9 MiB response. Under the 64 MiB limit it failed after 67,067,903 pack bytes. Co-Authored-By: Claude Code --- packages/gatekeeper-github/src/git-transport.ts | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/packages/gatekeeper-github/src/git-transport.ts b/packages/gatekeeper-github/src/git-transport.ts index 6d63fad98d..860e792513 100644 --- a/packages/gatekeeper-github/src/git-transport.ts +++ b/packages/gatekeeper-github/src/git-transport.ts @@ -31,10 +31,11 @@ import type { GitOid, GitPullHints } from "@gadgets/workshop-shared/gatekeeper"; /** * Maximum raw HTTP body size accepted from one upload-pack fetch, enforced while streaming (the - * transfer-size limiter pattern from gatekeeper-context's artifact-sync, same 64MB budget -- - * also matching the cap the overseer's `consumePack()` applies to the pack itself). + * transfer-size limiter pattern from gatekeeper-context's artifact-sync), matching the cap the + * overseer's `consumePack()` applies to the pack itself. Nothing here holds the body in memory: + * it streams through to the overseer, so the limit bounds the work of one pull. */ -export const MAX_GIT_FETCH_BYTES = 64 << 20; +export const MAX_GIT_FETCH_BYTES = 256 << 20; /** The `agent` capability sent with every request, mirroring the REST layer's User-Agent. */ const GIT_AGENT = "cloudflare-gadgets"; From 3f5e2a6fdd161e0416482d380ada5892049fd387 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:40:38 -0500 Subject: [PATCH 17/24] Time a git fetch out on a quiet server, not on its total length fetchGitUploadPack gave the whole fetch 120 s, body included. That suited a body read straight into memory. Since #656 the pack streams on to the overseer, which reads it only as fast as it stores it, so the budget now covers decoding and storing the pack too. A depth-1 mount of torvalds/linux took 70 s that way in local workerd, and a pack at the 256 MiB cap would be past 90 s. The 120 s now covers the wait for the response and an error body. Once the pack is streaming, git-transport.ts fails the fetch if the server sends nothing for GIT_FETCH_STALL_MS (60 s) while the pull is waiting on it. Time the reader takes between reads does not count. git's upload-pack sends a keepalive every few seconds when it has nothing else to send, so a live fetch stays well inside that. Two tests with fake timers: a body that goes quiet fails after 60 s, and a reader that pauses for ten times that between reads does not. The change to fetchGitUploadPack itself has no test; nothing covered its timeout before either. Co-Authored-By: Claude Code --- .../__tests__/git-transport.test.ts | 41 ++++++++++++- .../gatekeeper-github/src/git-transport.ts | 25 +++++++- packages/gatekeeper-github/src/github-api.ts | 57 +++++++++++-------- 3 files changed, 98 insertions(+), 25 deletions(-) diff --git a/packages/gatekeeper-github/__tests__/git-transport.test.ts b/packages/gatekeeper-github/__tests__/git-transport.test.ts index 7267e9c350..1c534ae4ff 100644 --- a/packages/gatekeeper-github/__tests__/git-transport.test.ts +++ b/packages/gatekeeper-github/__tests__/git-transport.test.ts @@ -7,10 +7,11 @@ // canonical empty pack, report-status parsing, and the push driver's body composition. import type { GitOid, GitPullHints } from "@gadgets/workshop-shared/gatekeeper"; -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; import { DELIM_PKT, FLUSH_PKT, + GIT_FETCH_STALL_MS, GitRefUpdateRejectedError, MAX_GIT_FETCH_BYTES, PktLineParser, @@ -310,6 +311,44 @@ describe("demuxGitFetchResponse", () => { await expect(collect(demuxGitFetchResponse(streamOf(response), limit))) .rejects.toThrow(/transfer limit/); }); + + it("gives the fetch up when the server goes quiet", async () => { + vi.useFakeTimers(); + try { + // A body that opens the packfile section and then neither sends nor ends. + let opened = false; + const quiet = new ReadableStream({ + pull(controller) { + if (opened) return new Promise(() => {}); + opened = true; + controller.enqueue(encodePktLine("packfile")); + }, + }); + const outcome = expect(collect(demuxGitFetchResponse(quiet, MAX_GIT_FETCH_BYTES))) + .rejects.toThrow("git fetch stalled: the server sent nothing for 60 seconds"); + await vi.advanceTimersByTimeAsync(GIT_FETCH_STALL_MS); + await outcome; + } finally { + vi.useRealTimers(); + } + }); + + it("does not count the time its reader takes between reads", async () => { + // The overseer reads the pack only as fast as it stores it. + vi.useFakeTimers(); + try { + const reader = + demuxGitFetchResponse(streamOf(packfileResponse()), MAX_GIT_FETCH_BYTES).getReader(); + const chunks = [(await reader.read()).value!]; + await vi.advanceTimersByTimeAsync(10 * GIT_FETCH_STALL_MS); + for (let next = await reader.read(); !next.done; next = await reader.read()) { + chunks.push(next.value); + } + expect(concatBytes(chunks)).toEqual(PACK_BYTES); + } finally { + vi.useRealTimers(); + } + }); }); describe("pullGitObjectsIntoCache", () => { diff --git a/packages/gatekeeper-github/src/git-transport.ts b/packages/gatekeeper-github/src/git-transport.ts index 860e792513..e0d03ae595 100644 --- a/packages/gatekeeper-github/src/git-transport.ts +++ b/packages/gatekeeper-github/src/git-transport.ts @@ -37,6 +37,13 @@ import type { GitOid, GitPullHints } from "@gadgets/workshop-shared/gatekeeper"; */ export const MAX_GIT_FETCH_BYTES = 256 << 20; +/** + * How long the server may send nothing, while a pull is waiting on it, before the fetch is given + * up as stalled. Time between reads does not count: the overseer stores the pack as it streams + * and reads at that pace, so a fetch has no fixed duration to hold it to. + */ +export const GIT_FETCH_STALL_MS = 60_000; + /** The `agent` capability sent with every request, mirroring the REST layer's User-Agent. */ const GIT_AGENT = "cloudflare-gadgets"; @@ -304,7 +311,7 @@ async function* demuxPackData( let received = 0; let inPackfile = false; while (true) { - let result = await reader.read(); + let result = await readOrStall(reader); if (result.done) { parser.finish(); throw new Error(inPackfile @@ -351,6 +358,22 @@ async function* demuxPackData( } } +// Reads the next chunk of a fetch body, failing if none arrives within GIT_FETCH_STALL_MS. +async function readOrStall(reader: ReadableStreamDefaultReader) + : Promise> { + let timer: ReturnType | undefined; + let stalled = new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error( + `git fetch stalled: the server sent nothing for ${GIT_FETCH_STALL_MS / 1000} seconds`)), + GIT_FETCH_STALL_MS); + }); + try { + return await Promise.race([reader.read(), stalled]); + } finally { + clearTimeout(timer); + } +} + // ======================================================================================= // Pull driver diff --git a/packages/gatekeeper-github/src/github-api.ts b/packages/gatekeeper-github/src/github-api.ts index da64df7037..39e73e8a83 100644 --- a/packages/gatekeeper-github/src/github-api.ts +++ b/packages/gatekeeper-github/src/github-api.ts @@ -1369,30 +1369,41 @@ export class GitHubApi { */ async fetchGitUploadPack(owner: string, repo: string, requestBody: Uint8Array): Promise { const url = `${LOGIN_BASE_URL}/${encodeURIComponent(owner)}/${encodeURIComponent(repo)}.git/git-upload-pack`; - const response = await fetch(url, { - method: "POST", - headers: { - "Content-Type": "application/x-git-upload-pack-request", - Accept: "application/x-git-upload-pack-result", - "Git-Protocol": "version=2", - "User-Agent": USER_AGENT, - Authorization: `Basic ${encodeBasicAuth("x-access-token", await this.#getToken())}`, - }, - body: requestBody, - // Longer than REQUEST_TIMEOUT_MS: the signal also covers streaming the response body, - // which may be a pack of tens of megabytes. - signal: AbortSignal.timeout(GIT_UPLOAD_PACK_TIMEOUT_MS), - }); - if (!response.ok) { - // The error body is short prose (e.g. "Repository not found"); a truncated copy makes the - // failure actionable without trusting its size. - const detail = (await response.text().catch(() => "")).trim().slice(0, 200); - throw new GitHubApiError( - response.status, - `git fetch failed: ${response.status} ${response.statusText}${detail ? `: ${detail}` : ""}`, - ); + // The timeout covers the wait for the response, and an error body, but not the pack. That + // streams on to the overseer, which reads it only as fast as it stores it, so no fixed + // budget fits; git-transport.ts gives the fetch up if the server goes quiet instead. + const token = await this.#getToken(); + const waiting = new AbortController(); + const timer = setTimeout( + () => waiting.abort(new DOMException("The operation timed out.", "TimeoutError")), + GIT_UPLOAD_PACK_TIMEOUT_MS, + ); + try { + const response = await fetch(url, { + method: "POST", + headers: { + "Content-Type": "application/x-git-upload-pack-request", + Accept: "application/x-git-upload-pack-result", + "Git-Protocol": "version=2", + "User-Agent": USER_AGENT, + Authorization: `Basic ${encodeBasicAuth("x-access-token", token)}`, + }, + body: requestBody, + signal: waiting.signal, + }); + if (!response.ok) { + // The error body is short prose (e.g. "Repository not found"); a truncated copy makes + // the failure actionable without trusting its size. + const detail = (await response.text().catch(() => "")).trim().slice(0, 200); + throw new GitHubApiError( + response.status, + `git fetch failed: ${response.status} ${response.statusText}${detail ? `: ${detail}` : ""}`, + ); + } + return response; + } finally { + clearTimeout(timer); } - return response; } /** From cef22545738ad7ea4ed46eacb14c7be8f0671a1e Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 15:57:00 -0500 Subject: [PATCH 18/24] Bound how much of one pack entry the reader buffers Native inflate needs an entry's whole zlib stream in one piece, so the reader buffers it. When the first attempt ran out of input it doubled the amount and tried again, for as long as the stream stayed unfinished. A stream can be padded without limit (empty stored blocks produce no output), so the only bound on that buffer was the pack cap, which the pako reader never leaned on: it took input a chunk at a time. At a 256 MiB pack cap that is more than the Overseer's memory. Three megabytes of padding around ten bytes decoded without complaint. inflate() now makes two attempts: with as much input as zlib itself could produce for the declared size, then with a deflater's worst case (an eighth over, plus a read's slack for small objects). A stream still unfinished there is rejected. The buffer is now bounded by the object cap, not the pack's. The vscode, TypeScript and Linux mount packs from GitHub (23,280, 65,044 and 98,591 objects) all decode within it. The test that held two 200 kB arrays to toStrictEqual now compares oids: under load it took nine seconds and timed out. Co-Authored-By: Claude Code --- .../__tests__/git-codec.test.ts | 36 +++++++++++++------ packages/workshop-backend/src/git-codec.ts | 12 ++++--- 2 files changed, 34 insertions(+), 14 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-codec.test.ts b/packages/workshop-backend/__tests__/git-codec.test.ts index 4d0220e4a8..ae7ac7efa0 100644 --- a/packages/workshop-backend/__tests__/git-codec.test.ts +++ b/packages/workshop-backend/__tests__/git-codec.test.ts @@ -249,32 +249,48 @@ describe("pack decoding", () => { let after = new TextEncoder().encode("the entry after it"); let pack = concatBytes(await buildPackBytes( [{ type: "blob", payload }, { type: "blob", payload: after }])); + let oids = [await gitObjectOid("blob", payload), await gitObjectOid("blob", after)]; for (let step of [undefined, 1000]) { - let objects = await decodePack(pack, { step }); - expect(objects.map(o => o.payload)).toStrictEqual([payload, after]); + expect((await decodePack(pack, { step })).map(o => o.oid)).toStrictEqual(oids); } }); - it("decodes an entry whose stream is longer than zlib would make it", async () => { - // Fifty empty stored blocks, then the payload in a final stored block and its Adler-32: a - // valid stream several times the length the first attempt allows for. - let payload = new TextEncoder().encode("padded out"); + // A zlib stream holding `payload` in one stored block, after `padding` empty stored blocks: + // valid, and as long as the padding makes it. + function storedStream(payload: Uint8Array, padding: number): Uint8Array { let a = 1, b = 0; for (let byte of payload) { a = (a + byte) % 65521; b = (b + a) % 65521; } - let pack = await packOfEntry(payload, entry => concatBytes([ - entry.subarray(0, 1), + let empty = new Uint8Array(5 * padding); + for (let i = 0; i < empty.length; i += 5) empty.set([0, 0, 0, 0xff, 0xff], i); + return concatBytes([ new Uint8Array([0x78, 0x01]), - ...Array.from({ length: 50 }, () => new Uint8Array([0, 0, 0, 0xff, 0xff])), + empty, new Uint8Array([1, payload.length, 0, ~payload.length & 0xff, 0xff]), payload, new Uint8Array([b >> 8, b & 0xff, a >> 8, a & 0xff]), - ])); + ]); + } + + it("decodes an entry whose stream is longer than zlib would make it", async () => { + let payload = new TextEncoder().encode("padded out"); + let pack = await packOfEntry(payload, + entry => concatBytes([entry.subarray(0, 1), storedStream(payload, 50)])); expect((await decodePack(pack))[0].payload).toStrictEqual(payload); }); + it("rejects an entry padded past any deflate of its size", async () => { + // Three megabytes of padding around ten bytes. A reader that buffers an entry until its + // stream ends is otherwise bounded only by the size of the pack. + let payload = new TextEncoder().encode("padded out"); + let pack = await packOfEntry(payload, + entry => concatBytes([entry.subarray(0, 1), storedStream(payload, 600_000)])); + await expect(decodePack(pack)) + .rejects.toThrow(/longer than any deflate of its declared size/); + }); + it("enforces the pack size cap", async () => { let pack = b64Bytes(PACK_NO_DELTA); await expect(decodePack(pack, { maxPackSize: pack.length - 1 })) diff --git a/packages/workshop-backend/src/git-codec.ts b/packages/workshop-backend/src/git-codec.ts index bd6618bad3..c7271bcb5d 100644 --- a/packages/workshop-backend/src/git-codec.ts +++ b/packages/workshop-backend/src/git-codec.ts @@ -473,11 +473,14 @@ class PackReader { // header; it was pre-checked against the object-size cap, and bounds the output here, so a // lying header cannot cause a larger allocation than it claimed. // - // inflateSync needs the whole stream in one piece and nothing records its length, so it is - // first given as much input as zlib itself could turn `size` bytes into. A stream that runs - // past that (another deflater's, or a padded one) is inflated again from twice as much. + // inflateSync needs the whole stream in one piece, and nothing records its length. It is first + // given as much input as zlib itself could turn `size` bytes into, then as much as any + // deflater's worst case: an eighth over, with a read's slack for a small object. A stream + // still unfinished there is padding, and the window that would hold more of it is bounded by + // nothing but the pack. async inflate(size: number): Promise { - for (let want = size + (size >>> 12) + (size >>> 14) + 13;; want *= 2) { + let longest = size + (size >>> 3) + PACK_READ_SIZE; + for (let want of [size + (size >>> 12) + (size >>> 14) + 13, longest]) { let buffered = await this.#fill(want); let input = this.#window.subarray(this.#pos, this.#pos + Math.min(buffered, want)); let inflated: InflatedEntry; @@ -509,6 +512,7 @@ class PackReader { return buffer.byteLength < buffer.buffer.byteLength ? new Uint8Array(buffer) : new Uint8Array(buffer.buffer); } + throw new Error("invalid packfile: object data longer than any deflate of its declared size"); } /** Ends the hashed body at the read position, returning its SHA-1 (hex). */ From d918113f9ec0f61fac4739220c28d43c0f8181c1 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:15:50 -0500 Subject: [PATCH 19/24] Leave the size of a pack to the overseer demuxGitFetchResponse counted the raw response body against MAX_GIT_FETCH_BYTES, a second copy of the overseer's cap on the pack that had to be kept equal to it by hand. The gatekeeper holds none of the body, so the limit protected nothing here that the overseer's does not: consumePack() rejects a pack over its cap and cancels the stream it was reading, which ends the fetch. The limiter is gone, along with the constant. A reader's cancel does reach the source across RPC: in workerd, a default stream handed to another Worker had its cancel() called one pull after that Worker's reader cancelled. One thing is given up. The limiter also counted the bytes around the pack (the sections before it, progress messages), which the overseer never sees. Those are a few bytes per packet from a server that was asked for no progress, but nothing bounds them now. Co-Authored-By: Claude Code --- .../__tests__/git-transport.test.ts | 38 +++++++++------ .../gatekeeper-github/src/git-transport.ts | 47 +++++++------------ 2 files changed, 40 insertions(+), 45 deletions(-) diff --git a/packages/gatekeeper-github/__tests__/git-transport.test.ts b/packages/gatekeeper-github/__tests__/git-transport.test.ts index 1c534ae4ff..4f6889ae6f 100644 --- a/packages/gatekeeper-github/__tests__/git-transport.test.ts +++ b/packages/gatekeeper-github/__tests__/git-transport.test.ts @@ -13,7 +13,6 @@ import { FLUSH_PKT, GIT_FETCH_STALL_MS, GitRefUpdateRejectedError, - MAX_GIT_FETCH_BYTES, PktLineParser, ZERO_OID, buildGitFetchRequest, @@ -260,26 +259,26 @@ function packfileResponse(options: { withSections?: boolean; progress?: boolean describe("demuxGitFetchResponse", () => { it("yields exactly the band-1 payload bytes", async () => { const pack = await collect( - demuxGitFetchResponse(streamOf(packfileResponse()), MAX_GIT_FETCH_BYTES)); + demuxGitFetchResponse(streamOf(packfileResponse()))); expect(pack).toEqual(PACK_BYTES); }); it("skips leading sections and progress frames", async () => { const pack = await collect(demuxGitFetchResponse( - streamOf(packfileResponse({ withSections: true, progress: true })), MAX_GIT_FETCH_BYTES)); + streamOf(packfileResponse({ withSections: true, progress: true })))); expect(pack).toEqual(PACK_BYTES); }); it("parses across arbitrary chunk boundaries", async () => { const whole = concatBytes(packfileResponse({ withSections: true, progress: true })); const rechunked = [whole.subarray(0, 3), whole.subarray(3, 27), whole.subarray(27)]; - const pack = await collect(demuxGitFetchResponse(streamOf(rechunked), MAX_GIT_FETCH_BYTES)); + const pack = await collect(demuxGitFetchResponse(streamOf(rechunked))); expect(pack).toEqual(PACK_BYTES); }); it("fails the stream on an ERR pkt with the server's message", async () => { const response = [encodePktLine(`ERR upload-pack: not our ref ${oid(3)}`), FLUSH_PKT]; - await expect(collect(demuxGitFetchResponse(streamOf(response), MAX_GIT_FETCH_BYTES))) + await expect(collect(demuxGitFetchResponse(streamOf(response)))) .rejects.toThrow(`git fetch failed: upload-pack: not our ref ${oid(3)}`); }); @@ -289,27 +288,36 @@ describe("demuxGitFetchResponse", () => { sidebandPkt(1, PACK_BYTES.subarray(0, 4)), sidebandPkt(3, new TextEncoder().encode("fatal: the remote end hung up")), ]; - await expect(collect(demuxGitFetchResponse(streamOf(response), MAX_GIT_FETCH_BYTES))) + await expect(collect(demuxGitFetchResponse(streamOf(response)))) .rejects.toThrow("git fetch failed: fatal: the remote end hung up"); }); it("rejects a response that ends without a flush", async () => { const truncated = packfileResponse().slice(0, -1); - await expect(collect(demuxGitFetchResponse(streamOf(truncated), MAX_GIT_FETCH_BYTES))) + await expect(collect(demuxGitFetchResponse(streamOf(truncated)))) .rejects.toThrow(/missing final flush/); }); it("rejects a response with no packfile section", async () => { const response = [encodePktLine("acknowledgments"), encodePktLine("NAK"), FLUSH_PKT]; - await expect(collect(demuxGitFetchResponse(streamOf(response), MAX_GIT_FETCH_BYTES))) + await expect(collect(demuxGitFetchResponse(streamOf(response)))) .rejects.toThrow(/no packfile section/); }); - it("enforces the transfer-size limit on the raw body", async () => { - const response = packfileResponse(); - const limit = concatBytes(response).byteLength - 1; - await expect(collect(demuxGitFetchResponse(streamOf(response), limit))) - .rejects.toThrow(/transfer limit/); + it("ends the fetch when its reader cancels", async () => { + // How a pack over the overseer's cap stops downloading: consumePack() rejects it and + // cancels the stream it was reading. + let cancelled = false; + const pieces = packfileResponse(); + let index = 0; + const body = new ReadableStream({ + pull(controller) { controller.enqueue(pieces[index++]); }, + cancel() { cancelled = true; }, + }, { highWaterMark: 0 }); + const reader = demuxGitFetchResponse(body).getReader(); + await reader.read(); + await reader.cancel(); + expect(cancelled).toBe(true); }); it("gives the fetch up when the server goes quiet", async () => { @@ -324,7 +332,7 @@ describe("demuxGitFetchResponse", () => { controller.enqueue(encodePktLine("packfile")); }, }); - const outcome = expect(collect(demuxGitFetchResponse(quiet, MAX_GIT_FETCH_BYTES))) + const outcome = expect(collect(demuxGitFetchResponse(quiet))) .rejects.toThrow("git fetch stalled: the server sent nothing for 60 seconds"); await vi.advanceTimersByTimeAsync(GIT_FETCH_STALL_MS); await outcome; @@ -338,7 +346,7 @@ describe("demuxGitFetchResponse", () => { vi.useFakeTimers(); try { const reader = - demuxGitFetchResponse(streamOf(packfileResponse()), MAX_GIT_FETCH_BYTES).getReader(); + demuxGitFetchResponse(streamOf(packfileResponse())).getReader(); const chunks = [(await reader.read()).value!]; await vi.advanceTimersByTimeAsync(10 * GIT_FETCH_STALL_MS); for (let next = await reader.read(); !next.done; next = await reader.read()) { diff --git a/packages/gatekeeper-github/src/git-transport.ts b/packages/gatekeeper-github/src/git-transport.ts index e0d03ae595..f85b42398a 100644 --- a/packages/gatekeeper-github/src/git-transport.ts +++ b/packages/gatekeeper-github/src/git-transport.ts @@ -29,14 +29,6 @@ import type { GitOid, GitPullHints } from "@gadgets/workshop-shared/gatekeeper"; -/** - * Maximum raw HTTP body size accepted from one upload-pack fetch, enforced while streaming (the - * transfer-size limiter pattern from gatekeeper-context's artifact-sync), matching the cap the - * overseer's `consumePack()` applies to the pack itself. Nothing here holds the body in memory: - * it streams through to the overseer, so the limit bounds the work of one pull. - */ -export const MAX_GIT_FETCH_BYTES = 256 << 20; - /** * How long the server may send nothing, while a pull is waiting on it, before the fetch is given * up as stalled. Time between reads does not count: the overseer stores the pack as it streams @@ -199,10 +191,10 @@ const OID_PATTERN = /^[0-9a-f]{40}$/; * strictly stronger than any `blob:limit`. * - A fetch whose wants are themselves blobs (`hints.type === "blob"`) sends **no filter**: a * blob want has no traversal for a filter to prune, and the filter would not suppress the - * wanted blob anyway -- an oversized blob arrives huge, the transfer limiter bounds the - * download, and the overseer's `put()`-equivalent size rejection measures and records its - * exact size (so later reads fail fast). This is the second of spike 1's two possible worlds; - * nothing is ever inferred from an absence. + * wanted blob anyway -- an oversized blob arrives huge, the overseer's caps on a pack and on + * any one object in it bound the download, and its `put()`-equivalent size rejection measures + * and records the blob's exact size (so later reads fail fast). This is the second of spike + * 1's two possible worlds; nothing is ever inferred from an absence. * - Otherwise `filterBlobSize` maps directly: 0 (fetch no blobs) → `blob:none`, N → `blob:limit=N` * (git's semantics -- omit blobs of size at least N -- match the hint's). */ @@ -272,17 +264,20 @@ const BAND_ERROR = 3; * irrelevant here, since every fetch is independent and nothing tracks shallow boundaries), * then demultiplex the packfile section's sideband (band 1 = pack data, band 2 = progress, * discarded, band 3 = fatal server error). An `ERR` pkt or band-3 message fails the stream with - * the server's message; `maxBytes` bounds the raw body (see MAX_GIT_FETCH_BYTES); a response - * that ends without a flush-pkt, or without ever reaching a packfile section, is an error -- - * a truncated pack must never look like a short success. `onFailure` is told what the stream - * failed with, which a reader on the far side of an RPC hop does not learn. + * the server's message; a response that ends without a flush-pkt, or without ever reaching a + * packfile section, is an error -- a truncated pack must never look like a short success. + * `onFailure` is told what the stream failed with, which a reader on the far side of an RPC hop + * does not learn. + * + * Nothing here limits the size of the body. The pack is the overseer's to limit: `consumePack()` + * rejects one over its cap and cancels this stream, which ends the fetch. The gatekeeper holds + * none of it either way. */ export function demuxGitFetchResponse( body: ReadableStream, - maxBytes: number, onFailure?: (error: unknown) => void, ): ReadableStream { - let iterator = demuxPackData(body, maxBytes); + let iterator = demuxPackData(body); return new ReadableStream({ async pull(controller) { let next: IteratorResult; @@ -303,12 +298,10 @@ export function demuxGitFetchResponse( async function* demuxPackData( body: ReadableStream, - maxBytes: number, ): AsyncGenerator { let reader = body.getReader(); try { let parser = new PktLineParser(); - let received = 0; let inPackfile = false; while (true) { let result = await readOrStall(reader); @@ -318,12 +311,7 @@ async function* demuxPackData( ? "truncated git fetch response: missing final flush" : "git fetch response contained no packfile section"); } - let value = result.value; - received += value.byteLength; - if (received > maxBytes) { - throw new Error(`git fetch response exceeded the ${maxBytes}-byte transfer limit`); - } - for (let item of parser.push(value)) { + for (let item of parser.push(result.value)) { if (item.kind === "delim" || item.kind === "response-end") continue; if (item.kind === "flush") { if (!inPackfile) { @@ -414,11 +402,10 @@ export async function pullGitObjectsIntoCache( throw new Error("git fetch failed: response had no body"); } // When the pack stream itself fails, the overseer's reader sees only that it ended early, and - // consumePack() rejects with that. What the stream failed with says why: the transfer limit, - // the server's own error, a truncated response. + // consumePack() rejects with that. What the stream failed with says why: the server's own + // error, a truncated response, a fetch that stalled. let failure: unknown; - let pack = demuxGitFetchResponse( - response.body, MAX_GIT_FETCH_BYTES, error => { failure = error; }); + let pack = demuxGitFetchResponse(response.body, error => { failure = error; }); let stored: Set; try { stored = new Set(await cache.consumePack(pack)); From 7f38a7758302003c4db856e60ee755f1b7c9b21f Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:15:52 -0500 Subject: [PATCH 20/24] Say that the pack cap is the one limit on a pull's size MAX_GIT_PACK_BYTES was documented as matching a transfer limit that gatekeepers were expected to apply to their fetch. The GitHub gatekeeper no longer applies one, and none needs to: a pack over the cap fails in consumePack(), which cancels the stream it was read from. Co-Authored-By: Claude Code --- packages/workshop-backend/src/git-cache.ts | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index ee2487ad59..1396bb6901 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -80,9 +80,10 @@ const logger = createWorkshopLogger("workshop.git-cache"); export const MAX_GIT_OBJECT_SIZE = 1 << 20; /** - * Maximum byte size of a packfile accepted by `consumePack()` (matching the transfer-size - * limiter gatekeepers are expected to apply to fetch bodies). A pack streams through, so this - * bounds what one pull may download, decode and store, not memory. + * Maximum byte size of a packfile accepted by `consumePack()`. A pack streams through, so this + * bounds what one pull may download, decode and store, not memory. It is the one limit on a + * pull's size: a gatekeeper need not bound its own fetch, because a pack over this fails here + * and the stream it was read from is cancelled. */ export const MAX_GIT_PACK_BYTES = 256 << 20; From 90e30af1246896082f2d1212b425b3ca9f852af0 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:24:02 -0500 Subject: [PATCH 21/24] Give the stall guard a timer that is never undefined readOrStall declared its timer as possibly undefined and assigned it inside the Promise executor. Under the Workers types clearTimeout takes `number | null`, so `tsc` in this package failed with TS2345 on the clearTimeout call, and the workspace build with it. The timer is now created outside the executor, so it always has a value. Checked with `tsc` run directly in the package: it fails on the previous commit and passes on this one. The earlier commit was pushed on the strength of `vp run build` reporting success, which was a cached result replayed for content it had not checked. Co-Authored-By: Claude Code --- packages/gatekeeper-github/src/git-transport.ts | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/packages/gatekeeper-github/src/git-transport.ts b/packages/gatekeeper-github/src/git-transport.ts index f85b42398a..2d1081db4c 100644 --- a/packages/gatekeeper-github/src/git-transport.ts +++ b/packages/gatekeeper-github/src/git-transport.ts @@ -349,12 +349,11 @@ async function* demuxPackData( // Reads the next chunk of a fetch body, failing if none arrives within GIT_FETCH_STALL_MS. async function readOrStall(reader: ReadableStreamDefaultReader) : Promise> { - let timer: ReturnType | undefined; - let stalled = new Promise((_, reject) => { - timer = setTimeout(() => reject(new Error( - `git fetch stalled: the server sent nothing for ${GIT_FETCH_STALL_MS / 1000} seconds`)), - GIT_FETCH_STALL_MS); - }); + let stall!: (error: Error) => void; + let stalled = new Promise((_, reject) => { stall = reject; }); + let timer = setTimeout(() => stall(new Error( + `git fetch stalled: the server sent nothing for ${GIT_FETCH_STALL_MS / 1000} seconds`)), + GIT_FETCH_STALL_MS); try { return await Promise.race([reader.read(), stalled]); } finally { From 84353e994a7a9e067a2143c1cd645b5827ec6192 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:24:02 -0500 Subject: [PATCH 22/24] Decode a pack entry from what has arrived, and name its input limit Two review comments on PackReader.inflate, both right. It asked for as much input as zlib could turn the declared size into before its first attempt. For an object that compresses well that is far more than its stream: two megabytes of zeros are a couple of kilobytes, yet the reader waited on two more megabytes of the pack behind them. If the stream failed in that wait, an entry that had fully arrived was never decoded, and an oversized one never measured. inflate() now starts from what is already buffered, grows to zlib's bound or one read's worth, and doubles from there. The cap on that input was described as the longest any deflate could be, and a stream past it as padding. Neither is true: a stream may be any length for its output, twenty thousand one-byte stored blocks being six times theirs. The cap is a limit on memory, since the window holds an entry's whole stream. The comment and the error now say so, and that a valid stream past it is refused. The other open comment, that the window could grow to the size of the pack, was answered by the commit that added the cap. New test: an oversized entry inside the first read is measured though the stream fails on the second. It fails on the previous commit. Co-Authored-By: Claude Code --- .../__tests__/git-cache.test.ts | 26 +++++++++++++++++++ .../__tests__/git-codec.test.ts | 10 +++---- packages/workshop-backend/src/git-codec.ts | 24 ++++++++++++----- 3 files changed, 49 insertions(+), 11 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index 24ca252357..7b686848b2 100644 --- a/packages/workshop-backend/__tests__/git-cache.test.ts +++ b/packages/workshop-backend/__tests__/git-cache.test.ts @@ -934,6 +934,32 @@ describe("consumePack", () => { .toBe(big.byteLength); }); + it("measures an oversized entry that has arrived, though the stream then fails", async () => { + // Two megabytes of zeros deflate to a couple of kilobytes, so the whole entry is in the + // first read. Decoding it must not wait on as much of the pack again as it inflates to. + let t = makeCache(); + let big = new Uint8Array(2 * MAX_GIT_OBJECT_SIZE); + let next = new Uint8Array(200_000); + for (let pos = 0; pos < next.length; pos += 50_000) { + crypto.getRandomValues(next.subarray(pos, pos + 50_000)); + } + let pack = concatBytes(await buildPackBytes( + [{ type: "blob", payload: big }, { type: "blob", payload: next }])); + let sent = false; + let failing = new ReadableStream({ + type: "bytes", + pull(controller) { + if (sent) controller.error(new Error("connection lost")); + else controller.enqueue(pack.slice(0, 64 << 10)); + sent = true; + }, + }); + await expect(new GitCacheImpl(t.cache, G1).consumePack(failing)) + .rejects.toThrow("connection lost"); + expect(t.storage.gitObjectMetadata.get(await gitObjectOid("blob", big))?.size) + .toBe(big.byteLength); + }); + it("resolves a delta against an oversized base it declines to store", async () => { // How git packs a file similar to a large one (e.g. a second lockfile in a whole-tree blob // pull), hand-built: the large blob, then a ref-delta copying its first 16 bytes. diff --git a/packages/workshop-backend/__tests__/git-codec.test.ts b/packages/workshop-backend/__tests__/git-codec.test.ts index ae7ac7efa0..de9713d4c6 100644 --- a/packages/workshop-backend/__tests__/git-codec.test.ts +++ b/packages/workshop-backend/__tests__/git-codec.test.ts @@ -281,14 +281,14 @@ describe("pack decoding", () => { expect((await decodePack(pack))[0].payload).toStrictEqual(payload); }); - it("rejects an entry padded past any deflate of its size", async () => { - // Three megabytes of padding around ten bytes. A reader that buffers an entry until its - // stream ends is otherwise bounded only by the size of the pack. + it("refuses an entry with more compressed data than one entry may buffer", async () => { + // Three megabytes of empty blocks around ten bytes: a valid stream, but the reader holds an + // entry's whole stream, and nothing else would bound that short of the pack's own size. let payload = new TextEncoder().encode("padded out"); let pack = await packOfEntry(payload, entry => concatBytes([entry.subarray(0, 1), storedStream(payload, 600_000)])); - await expect(decodePack(pack)) - .rejects.toThrow(/longer than any deflate of its declared size/); + await expect(decodePack(pack)).rejects.toThrow( + "packfile entry of 10 bytes has more than 65547 bytes of compressed data"); }); it("enforces the pack size cap", async () => { diff --git a/packages/workshop-backend/src/git-codec.ts b/packages/workshop-backend/src/git-codec.ts index c7271bcb5d..3dc508393d 100644 --- a/packages/workshop-backend/src/git-codec.ts +++ b/packages/workshop-backend/src/git-codec.ts @@ -474,13 +474,20 @@ class PackReader { // lying header cannot cause a larger allocation than it claimed. // // inflateSync needs the whole stream in one piece, and nothing records its length. It is first - // given as much input as zlib itself could turn `size` bytes into, then as much as any - // deflater's worst case: an eighth over, with a read's slack for a small object. A stream - // still unfinished there is padding, and the window that would hold more of it is bounded by - // nothing but the pack. + // given what is already buffered, up to what zlib itself could turn `size` bytes into or one + // read's worth if that is less: an entry that has arrived is then decoded without waiting on + // the pack behind it. While the stream is unfinished the input grows, to that first amount and + // then doubling, up to `longest`: an eighth over `size`, with a read's slack for a small + // object. + // + // `longest` is a limit on memory, not a rule of the format. The window holds an entry's whole + // stream, and a stream may be any length for its output (short stored blocks, empty ones), so + // without it one entry could buffer as much as the pack cap allows. No deflater git servers + // use goes past it, but a valid stream that does is refused. async inflate(size: number): Promise { let longest = size + (size >>> 3) + PACK_READ_SIZE; - for (let want of [size + (size >>> 12) + (size >>> 14) + 13, longest]) { + let first = Math.min(size + (size >>> 12) + (size >>> 14) + 13, PACK_READ_SIZE); + for (let want = Math.min(first, this.#end - this.#pos || first);;) { let buffered = await this.#fill(want); let input = this.#window.subarray(this.#pos, this.#pos + Math.min(buffered, want)); let inflated: InflatedEntry; @@ -491,6 +498,12 @@ class PackReader { let code = err instanceof Error && "code" in err ? err.code : undefined; if (code === "Z_BUF_ERROR") { if (buffered < want) throw new Error("invalid packfile: truncated", { cause: err }); + if (want === longest) { + throw new Error( + `packfile entry of ${size} bytes has more than ${longest} bytes of compressed ` + + `data, over the limit for one entry`, { cause: err }); + } + want = want < first ? first : Math.min(2 * want, longest); continue; } let detail = err instanceof Error ? err.message : String(err); @@ -512,7 +525,6 @@ class PackReader { return buffer.byteLength < buffer.buffer.byteLength ? new Uint8Array(buffer) : new Uint8Array(buffer.buffer); } - throw new Error("invalid packfile: object data longer than any deflate of its declared size"); } /** Ends the hashed body at the read position, returning its SHA-1 (hex). */ From 78d9be0f601731eb3fef8f68b53a389aae85f560 Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:36:12 -0500 Subject: [PATCH 23/24] Bound the part of a fetch response that is not pack data Removing the transfer limit left one thing unbounded, which review caught: the bytes of a response that are not pack data. The overseer caps the pack it is sent, but it never sees packet framing, the sections before the pack, progress or keepalives, and each of those that arrives also holds off the stall timeout. A server that kept sending progress would be read without end. demuxPackData now counts what it receives against what it delivers, and gives the fetch up once the difference passes 1 MiB plus a sixteenth of the pack data delivered so far. That scales with the pack the overseer is already limiting, so there is still no second copy of its cap here. A depth-1 mount of torvalds/linux through this transport: 200,205,634 bytes of pack and 106,875 of everything else, 0.053%. Co-Authored-By: Claude Code --- .../__tests__/git-transport.test.ts | 14 +++++++++++ .../gatekeeper-github/src/git-transport.ts | 25 ++++++++++++++++--- 2 files changed, 36 insertions(+), 3 deletions(-) diff --git a/packages/gatekeeper-github/__tests__/git-transport.test.ts b/packages/gatekeeper-github/__tests__/git-transport.test.ts index 4f6889ae6f..66e925cf46 100644 --- a/packages/gatekeeper-github/__tests__/git-transport.test.ts +++ b/packages/gatekeeper-github/__tests__/git-transport.test.ts @@ -304,6 +304,20 @@ describe("demuxGitFetchResponse", () => { .rejects.toThrow(/no packfile section/); }); + it("gives up on a response that is mostly not pack data", async () => { + // The overseer limits the pack it is sent, but it never sees progress, framing or + // keepalives, and each of those that arrives also holds off the stall timeout. + const progress = new Uint8Array(60_000); + const response = [ + encodePktLine("packfile"), + sidebandPkt(1, PACK_BYTES), + ...Array.from({ length: 20 }, () => sidebandPkt(2, progress)), + FLUSH_PKT, + ]; + await expect(collect(demuxGitFetchResponse(streamOf(response)))) + .rejects.toThrow(/bytes that are not pack data/); + }); + it("ends the fetch when its reader cancels", async () => { // How a pack over the overseer's cap stops downloading: consumePack() rejects it and // cancels the stream it was reading. diff --git a/packages/gatekeeper-github/src/git-transport.ts b/packages/gatekeeper-github/src/git-transport.ts index 2d1081db4c..6964d53523 100644 --- a/packages/gatekeeper-github/src/git-transport.ts +++ b/packages/gatekeeper-github/src/git-transport.ts @@ -29,6 +29,15 @@ import type { GitOid, GitPullHints } from "@gadgets/workshop-shared/gatekeeper"; +/** + * How much of a fetch response may be something other than pack data (packet framing, the + * sections before the pack, progress, keepalives) before the fetch is given up: this much, plus + * a sixteenth of the pack data delivered so far. The overseer bounds the pack, but it never sees + * these bytes, and each one that arrives also holds off the stall timeout. git's own framing is + * five bytes on a packet of eight thousand or more. + */ +export const MAX_GIT_FETCH_OVERHEAD_BYTES = 1 << 20; + /** * How long the server may send nothing, while a pull is waiting on it, before the fetch is given * up as stalled. Time between reads does not count: the overseer stores the pack as it streams @@ -269,9 +278,10 @@ const BAND_ERROR = 3; * `onFailure` is told what the stream failed with, which a reader on the far side of an RPC hop * does not learn. * - * Nothing here limits the size of the body. The pack is the overseer's to limit: `consumePack()` - * rejects one over its cap and cancels this stream, which ends the fetch. The gatekeeper holds - * none of it either way. + * The pack's size is the overseer's to limit: `consumePack()` rejects one over its cap and + * cancels this stream, which ends the fetch. What is limited here is the rest of the response, + * which the overseer never sees (see MAX_GIT_FETCH_OVERHEAD_BYTES). The gatekeeper holds none + * of the body either way. */ export function demuxGitFetchResponse( body: ReadableStream, @@ -303,6 +313,8 @@ async function* demuxPackData( try { let parser = new PktLineParser(); let inPackfile = false; + let received = 0; + let delivered = 0; while (true) { let result = await readOrStall(reader); if (result.done) { @@ -311,6 +323,7 @@ async function* demuxPackData( ? "truncated git fetch response: missing final flush" : "git fetch response contained no packfile section"); } + received += result.value.byteLength; for (let item of parser.push(result.value)) { if (item.kind === "delim" || item.kind === "response-end") continue; if (item.kind === "flush") { @@ -333,6 +346,7 @@ async function* demuxPackData( let band = item.data[0]; let payload = item.data.subarray(1); if (band === BAND_PACK) { + delivered += payload.byteLength; if (payload.byteLength > 0) yield payload; } else if (band === BAND_ERROR) { throw new Error(`git fetch failed: ${pktText(payload)}`); @@ -340,6 +354,11 @@ async function* demuxPackData( throw new Error(`malformed sideband frame: unknown band ${band}`); } } + if (received - delivered > MAX_GIT_FETCH_OVERHEAD_BYTES + delivered / 16) { + throw new Error( + `git fetch response carried ${received - delivered} bytes that are not pack data, ` + + `with ${delivered} bytes that are`); + } } } finally { await reader.cancel().catch(() => {}); From ec0ceddb8eeabf71d7ea0d2e57def7c0b61b2d0f Mon Sep 17 00:00:00 2001 From: Maximo Guk <62088388+Maximo-Guk@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:42:30 -0500 Subject: [PATCH 24/24] Inflate a large pack entry a read at a time Native inflate needs an entry's whole compressed stream at once, and joins its output from 16 KiB chunks. For one large object that is the stream, the chunks and the joined copy all live together: review worked it through for a 48 MiB blob deflated at level 0, about 150 MiB against the Overseer's 128 MB. The pako reader in #656 held one read of input, so the same entry fit there. An entry declaring more than 1 MiB now goes back to that loop: pako, a read at a time, with no input kept beyond the read in hand. Its output is written into one buffer of the declared size, so the chunk list and the copy that joined it are gone too, and inflating that 48 MiB entry takes about 48 MiB where #656 took 96. Entries up to 1 MiB, which is all of a mount pack but the odd large tree, stay on native inflate, where the most one can buffer is now a little over a megabyte. pako's Inflate import, its internals interface and that cast are back for this path. It is given windowBits 15, so it reads zlib only, not the gzip it would otherwise also accept. Run for real: the 48 MiB level-0 entry goes through consumePack and is measured; the vscode, TypeScript and Linux mount packs decode to the same totals, TypeScript's two trees over 1 MiB on the new path. Memory is reasoned from what each path holds, not measured. Co-Authored-By: Claude Code --- .../__tests__/git-codec.test.ts | 38 ++++++++++ packages/workshop-backend/src/git-codec.ts | 69 ++++++++++++++++++- 2 files changed, 104 insertions(+), 3 deletions(-) diff --git a/packages/workshop-backend/__tests__/git-codec.test.ts b/packages/workshop-backend/__tests__/git-codec.test.ts index de9713d4c6..9ba6e19bd2 100644 --- a/packages/workshop-backend/__tests__/git-codec.test.ts +++ b/packages/workshop-backend/__tests__/git-codec.test.ts @@ -255,6 +255,44 @@ describe("pack decoding", () => { } }); + it("decodes an entry large enough to be inflated a read at a time", async () => { + // Over a megabyte the reader stops holding an entry's whole stream. Random bytes make that + // stream as long as the object, some fifty reads. + let payload = new Uint8Array(3 << 20); + for (let pos = 0; pos < payload.length; pos += 65536) { + crypto.getRandomValues(payload.subarray(pos, pos + 65536)); + } + let after = new TextEncoder().encode("the entry after it"); + let pack = concatBytes(await buildPackBytes( + [{ type: "blob", payload }, { type: "blob", payload: after }])); + let oids = [await gitObjectOid("blob", payload), await gitObjectOid("blob", after)]; + for (let step of [undefined, 1000]) { + expect((await decodePack(pack, { step })).map(o => o.oid)).toStrictEqual(oids); + } + }); + + it("rejects a large entry whose data is not the size it declares", async () => { + // The same check as for a small entry, on the path that inflates a read at a time. + let payload = new Uint8Array((1 << 20) + 10).fill(7); + let header = (size: number) => { + let bytes = []; + let first = (3 << 4) | (size & 0x0f); + for (let rest = size >>> 4; rest > 0; rest >>>= 7) { + bytes.push(first | 0x80); + first = rest & 0x7f; + } + return new Uint8Array([...bytes, first]); + }; + let declaring = (size: number) => packOfEntry(payload, entry => concatBytes( + [header(size), entry.subarray(header(payload.length).length)])); + await expect(decodePack(await declaring(payload.length - 1))) + .rejects.toThrow(/larger than its declared/); + await expect(decodePack(await declaring(payload.length + 1))) + .rejects.toThrow(/smaller than its declared/); + expect((await decodePack(await declaring(payload.length)))[0].oid) + .toBe(await gitObjectOid("blob", payload)); + }); + // A zlib stream holding `payload` in one stored block, after `padding` empty stored blocks: // valid, and as long as the padding makes it. function storedStream(payload: Uint8Array, padding: number): Uint8Array { diff --git a/packages/workshop-backend/src/git-codec.ts b/packages/workshop-backend/src/git-codec.ts index 3dc508393d..fe8ab52760 100644 --- a/packages/workshop-backend/src/git-codec.ts +++ b/packages/workshop-backend/src/git-codec.ts @@ -17,11 +17,13 @@ // use workerd's native node:zlib: pako's deflate takes about twice the CPU of the native one, and // its inflate allocates some 100 KiB per stream. Pack entries are concatenated zlib streams with // no recorded lengths, so the decoder has to learn where each one ends -- inflateSync reports the -// input it consumed when asked for `info`, which DecompressionStream cannot do. Only writing a -// pack (buildPackBytes, the push path) still uses pako, the library isomorphic-git bundles. +// input it consumed when asked for `info`, which DecompressionStream cannot do. It needs the +// whole stream at once, though, so an entry too large to hold that way is inflated a read at a +// time with pako (the library isomorphic-git bundles), which also writes packs (buildPackBytes, +// the push path). import { constants, deflateSync, inflateSync } from "node:zlib"; -import { deflate } from "pako"; +import { Inflate, deflate } from "pako"; import type { GitObjectType, GitOid } from "@gadgets/workshop-shared/gatekeeper"; const ENCODER = new TextEncoder(); @@ -421,8 +423,27 @@ interface InflatedEntry { engine: { bytesWritten: number }; } +// The pako Inflate internals the streamed path relies on beyond @types/pako's declarations, all +// stable pako API in practice (isomorphic-git's own pack parser relies on `strm.avail_in` the +// same way): `ended` flips when the zlib stream completes mid-input, and `strm.avail_in` is how +// many bytes of the last push() the stream did not consume -- together they locate the entry +// boundary. +interface InflateWithInternals { + ended: boolean; + err: number; + msg: string; + strm: { avail_in: number }; + onData: (chunk: Uint8Array) => void; + push(data: Uint8Array, flush: boolean): void; +} + const PACK_READ_SIZE = 64 << 10; +// Entries declaring more than this are inflated a read at a time. Native inflate holds an entry's +// whole compressed stream beside its output and a second copy of that output while it joins it: +// three times the object, which for one of tens of megabytes is more than the Overseer has. +const PACK_STREAMED_ENTRY_SIZE = 1 << 20; + // decodePackStream's reads, in order, from a window over the bytes that have arrived and are not // yet consumed, hashing every byte before `endBody()` (the trailer's SHA-1 input) as it leaves // the window. Reads are BYOB, which a gatekeeper facet's pack stream supports once Workers RPC @@ -485,6 +506,7 @@ class PackReader { // without it one entry could buffer as much as the pack cap allows. No deflater git servers // use goes past it, but a valid stream that does is refused. async inflate(size: number): Promise { + if (size > PACK_STREAMED_ENTRY_SIZE) return this.#inflateStreamed(size); let longest = size + (size >>> 3) + PACK_READ_SIZE; let first = Math.min(size + (size >>> 12) + (size >>> 14) + 13, PACK_READ_SIZE); for (let want = Math.min(first, this.#end - this.#pos || first);;) { @@ -527,6 +549,47 @@ class PackReader { } } + // Inflates a large entry a read at a time, into a buffer of the size it declares. No input is + // kept beyond the read in hand, so the stream may be any length for its output, and the output + // is written once. + async #inflateStreamed(size: number): Promise { + let inflator = new Inflate({ windowBits: 15 }) as unknown as InflateWithInternals; + let out = new Uint8Array(size); + let total = 0; + let overflow = false; + inflator.onData = (chunk: Uint8Array) => { + if (total + chunk.byteLength > size) { + overflow = true; + // pako offers no abort; raising here unwinds through push() below. + throw new Error("pack entry exceeds declared size"); + } + out.set(chunk, total); + total += chunk.byteLength; + }; + + try { + while (!inflator.ended) { + if (await this.#fill(1) === 0) throw new Error("invalid packfile: truncated"); + inflator.push(this.#window.subarray(this.#pos, this.#end), false); + if (inflator.err) { + throw new Error( + `invalid packfile: corrupt object data (${inflator.msg || inflator.err})`); + } + this.#pos = this.#end - inflator.strm.avail_in; + } + } catch (err) { + if (overflow) { + throw new Error("invalid packfile: object larger than its declared size", { cause: err }); + } + throw err; + } + + if (total !== size) { + throw new Error("invalid packfile: object smaller than its declared size"); + } + return out; + } + /** Ends the hashed body at the read position, returning its SHA-1 (hex). */ async endBody(): Promise { let hash = this.#hash!;