diff --git a/packages/gatekeeper-github/__tests__/git-transport.test.ts b/packages/gatekeeper-github/__tests__/git-transport.test.ts index 604b6566f0..66e925cf46 100644 --- a/packages/gatekeeper-github/__tests__/git-transport.test.ts +++ b/packages/gatekeeper-github/__tests__/git-transport.test.ts @@ -7,12 +7,12 @@ // 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, ZERO_OID, buildGitFetchRequest, @@ -259,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)}`); }); @@ -288,27 +288,88 @@ 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("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. + 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 () => { + 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))) + .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())).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(); + } }); }); @@ -344,6 +405,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..6964d53523 100644 --- a/packages/gatekeeper-github/src/git-transport.ts +++ b/packages/gatekeeper-github/src/git-transport.ts @@ -30,11 +30,20 @@ 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). + * 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_BYTES = 64 << 20; +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 + * 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"; @@ -191,10 +200,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). */ @@ -264,18 +273,30 @@ 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. + * 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. + * + * 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, - maxBytes: number, + onFailure?: (error: unknown) => void, ): ReadableStream { - let iterator = demuxPackData(body, maxBytes); + let iterator = demuxPackData(body); 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); }, @@ -287,27 +308,23 @@ 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; + let received = 0; + let delivered = 0; while (true) { - let result = await reader.read(); + let result = await readOrStall(reader); if (result.done) { parser.finish(); throw new Error(inPackfile ? "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)) { + 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") { if (!inPackfile) { @@ -329,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)}`); @@ -336,12 +354,32 @@ 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(() => {}); } } +// Reads the next chunk of a fetch body, failing if none arrives within GIT_FETCH_STALL_MS. +async function readOrStall(reader: ReadableStreamDefaultReader) + : Promise> { + 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 { + clearTimeout(timer); + } +} + // ======================================================================================= // Pull driver @@ -381,8 +419,17 @@ 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 server's own + // error, a truncated response, a fetch that stalled. + let failure: unknown; + let pack = demuxGitFetchResponse(response.body, 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${ 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; } /** diff --git a/packages/workshop-backend/__tests__/git-cache.test.ts b/packages/workshop-backend/__tests__/git-cache.test.ts index 7eea12662b..7b686848b2 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 } from "vitest"; +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"; @@ -8,6 +9,8 @@ import { GitCacheImpl, GitObjectTooLargeError, MAX_GIT_OBJECT_SIZE, + MAX_GIT_PACK_OBJECT_SIZE, + MAX_OVERSIZED_BASE_BYTES, WorkspaceGitCache, } from "../src/git-cache"; import { GitStore, blobOid } from "../src/git-store"; @@ -802,6 +805,65 @@ 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(); + 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("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(); @@ -844,6 +906,60 @@ 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(); + 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("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. @@ -865,6 +981,92 @@ describe("consumePack", () => { expect(t.cache.readLocalObject(stored[0])!.payload).toStrictEqual(target); expect(t.storage.gitObjectMetadata.get(bigOid)!.size).toBe(big.byteLength); }); + + // 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 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. + 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 after = concatBytes(await buildPackBytes(trailing)); + let body = concatBytes([ + 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, 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 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); + 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(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 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, undefined, [commitOid]) + .consumePack(byteStream(pack)); + expect(stored).toStrictEqual([targetOid, commitOid]); + }); + + 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), asked); + }); + 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/__tests__/git-codec.test.ts b/packages/workshop-backend/__tests__/git-codec.test.ts index e2ccc59f17..9ba6e19bd2 100644 --- a/packages/workshop-backend/__tests__/git-codec.test.ts +++ b/packages/workshop-backend/__tests__/git-codec.test.ts @@ -194,6 +194,141 @@ 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/); + }); + + // 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))]); + } + + 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 }])); + 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("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 { + let a = 1, b = 0; + for (let byte of payload) { + a = (a + byte) % 65521; + b = (b + a) % 65521; + } + 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]), + 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("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( + "packfile entry of 10 bytes has more than 65547 bytes of compressed data"); + }); + it("enforces the pack size cap", async () => { let pack = b64Bytes(PACK_NO_DELTA); await expect(decodePack(pack, { maxPackSize: pack.length - 1 })) @@ -213,6 +348,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"); diff --git a/packages/workshop-backend/src/git-cache.ts b/packages/workshop-backend/src/git-cache.ts index a4f872967c..1396bb6901 100644 --- a/packages/workshop-backend/src/git-cache.ts +++ b/packages/workshop-backend/src/git-cache.ts @@ -80,11 +80,27 @@ 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), and the hard per-object - * inflation bound while decoding one. + * 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 = 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 + * 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; /** * Blob size fetched eagerly when pulling a worktree base: the pull requests the base commit with @@ -255,40 +271,60 @@ 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, + * 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. 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 droppable = new Set(asked); 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), + maxObjectSize: MAX_GIT_PACK_OBJECT_SIZE, + 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)); - stored.push(oid); - } else { + if (object.type !== "blob" && object.payload.byteLength <= MAX_GIT_OBJECT_SIZE) { held.set(oid, object); + } else if (this.#storeVerifiedObject(gatekeeperId, oid, object)) { + stored.push(oid); + } else if (!oversized.has(oid)) { + oversized.set(oid, object); + oversizedBytes += object.payload.byteLength; + // Oldest out first: never the one just added, nor one the pull did not ask for. + for (let [oldest, { payload }] of oversized) { + if (oversizedBytes <= MAX_OVERSIZED_BASE_BYTES || oldest === oid) break; + if (!droppable.has(oldest)) continue; + oversized.delete(oldest); + oversizedBytes -= payload.byteLength; + } } } 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; } @@ -335,23 +371,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) { @@ -1158,10 +1193,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); } } @@ -1176,18 +1215,25 @@ 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. + // 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); @@ -1214,12 +1260,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(); } @@ -1256,7 +1304,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/git-codec.ts b/packages/workshop-backend/src/git-codec.ts index 58f553ffee..fe8ab52760 100644 --- a/packages/workshop-backend/src/git-codec.ts +++ b/packages/workshop-backend/src/git-codec.ts @@ -13,13 +13,16 @@ // 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. - -import { deflateSync, inflateSync } from "node:zlib"; +// 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. 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 { Inflate, deflate } from "pako"; import type { GitObjectType, GitOid } from "@gadgets/workshop-shared/gatekeeper"; @@ -77,10 +80,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. */ @@ -351,6 +358,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; } @@ -402,10 +415,19 @@ 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. +// 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 }; +} + +// 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; @@ -417,21 +439,31 @@ interface InflateWithInternals { 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). +// 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 +// 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" }); @@ -440,68 +472,110 @@ 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.read(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. It is first + // 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 inflator = new Inflate() as unknown as InflateWithInternals; - let chunks: Uint8Array[] = []; + 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);;) { + 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 }); + 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); + throw new Error( + code === "ERR_BUFFER_TOO_LARGE" + ? "invalid packfile: object larger than its declared size" + : `invalid packfile: corrupt object data (${detail})`, + { 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`); + } + 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); + } + } + + // 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) => { - total += chunk.byteLength; - if (total > size) { + 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"); } - chunks.push(chunk); + out.set(chunk, total); + total += chunk.byteLength; }; try { while (!inflator.ended) { - await this.#fill(); - inflator.push(this.#chunk.subarray(this.#pos), false); + 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})`); + throw new Error( + `invalid packfile: corrupt object data (${inflator.msg || inflator.err})`); } - this.#pos = this.#chunk.byteLength - inflator.strm.avail_in; + this.#pos = this.#end - inflator.strm.avail_in; } } catch (err) { if (overflow) { @@ -513,14 +587,14 @@ class PackReader { if (total !== size) { throw new Error("invalid packfile: object smaller than its declared size"); } - return concatBytes(chunks); + return out; } /** 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)); } @@ -531,8 +605,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; } } 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 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;