Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
e85dac8
Bound a pack entry's size varint
Maximo-Guk Oct 4, 2026
953cf79
Read a pack in full 64 KiB buffers
Maximo-Guk Oct 4, 2026
25fa791
Skip the pull-routing hint a gatekeeper has already proven
Maximo-Guk Oct 4, 2026
bba045a
Leave an object the gatekeeper already stored untouched
Maximo-Guk Oct 4, 2026
6355337
Measure an oversized pack entry as it arrives
Maximo-Guk Oct 4, 2026
4b4e0e9
Store a streamed blob without a transaction of its own
Maximo-Guk Oct 4, 2026
e1a748a
Deflate loose git objects at zlib's fastest level
Maximo-Guk Oct 4, 2026
579aa3c
Inflate pack entries with native zlib
Maximo-Guk Oct 4, 2026
08146e0
Test a pack arriving over RPC, and pin the delta order rule
Maximo-Guk Oct 4, 2026
bb3aff8
Say what consumePack returns and leaves behind on failure
Maximo-Guk Oct 4, 2026
2163ddb
Bound the oversized objects kept as delta bases
Maximo-Guk Oct 4, 2026
3f75573
Drop oversized bases only from a pack of blobs alone
Maximo-Guk Oct 4, 2026
0f45e11
Drop only the oversized bases a pull asked for by name
Maximo-Guk Oct 4, 2026
b3dd592
Report why a git fetch failed, not the disconnect the cache saw
Maximo-Guk Oct 4, 2026
1976503
Give a pack's objects their own size cap, and raise the pack's
Maximo-Guk Oct 4, 2026
60c5292
Raise the git fetch transfer limit to match the pack cap
Maximo-Guk Oct 4, 2026
3f5e2a6
Time a git fetch out on a quiet server, not on its total length
Maximo-Guk Oct 4, 2026
cef2254
Bound how much of one pack entry the reader buffers
Maximo-Guk Oct 4, 2026
d918113
Leave the size of a pack to the overseer
Maximo-Guk Oct 4, 2026
7f38a77
Say that the pack cap is the one limit on a pull's size
Maximo-Guk Oct 4, 2026
90e30af
Give the stall guard a timer that is never undefined
Maximo-Guk Oct 4, 2026
84353e9
Decode a pack entry from what has arrived, and name its input limit
Maximo-Guk Oct 4, 2026
78d9be0
Bound the part of a fetch response that is not pack data
Maximo-Guk Oct 4, 2026
ec0cedd
Inflate a large pack entry a read at a time
Maximo-Guk Oct 4, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
114 changes: 100 additions & 14 deletions packages/gatekeeper-github/__tests__/git-transport.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)}`);
});

Expand All @@ -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<Uint8Array>({
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<Uint8Array>({
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();
}
});
});

Expand Down Expand Up @@ -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<Uint8Array>): Promise<GitOid[]> {
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<Uint8Array>): Promise<GitOid[]> {
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(
Expand Down
97 changes: 72 additions & 25 deletions packages/gatekeeper-github/src/git-transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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).
*/
Expand Down Expand Up @@ -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<Uint8Array>,
maxBytes: number,
onFailure?: (error: unknown) => void,
): ReadableStream<Uint8Array> {
let iterator = demuxPackData(body, maxBytes);
let iterator = demuxPackData(body);
return new ReadableStream({
async pull(controller) {
let next = await iterator.next();
let next: IteratorResult<Uint8Array, void>;
try {
next = await iterator.next();
} catch (error) {
onFailure?.(error);
throw error;
}
if (next.done) controller.close();
else controller.enqueue(next.value);
},
Expand All @@ -287,27 +308,23 @@ export function demuxGitFetchResponse(

async function* demuxPackData(
body: ReadableStream<Uint8Array>,
maxBytes: number,
): AsyncGenerator<Uint8Array, void, unknown> {
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)) {

@devin-ai-integration devin-ai-integration Bot Oct 4, 2026 •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Final flush bypasses fetch overhead limit

When a response chunk ends with a flush, demuxPackData returns before checking accumulated overhead. A single chunk can carry excessive progress data and still complete successfully.

Learn more

The demultiplexer counts body bytes in received and pack bytes in delivered, but checks their difference only after processing every parsed packet in a chunk. Processing a flush exits the generator from inside that loop, so a chunk containing a flush never reaches this check. That is common for a small response received as one chunk, and it can also happen after a large progress packet. The intended response overhead limit is then unenforced for the final chunk.

Example: A fetch body arrives as one chunk containing the packfile header, a small band-1 pack, twenty 60,000-byte band-2 progress packets, and a flush. The generator yields the pack and returns on the flush; it never checks the roughly 1.2 MB of progress against the 1 MB allowance.

Recommended fix: Check the accumulated overhead before returning on a flush, and ensure the check also runs for terminal packets in a chunk. Keep the ordinary between-chunk check for responses that continue without a flush.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

if (item.kind === "delim" || item.kind === "response-end") continue;
if (item.kind === "flush") {
if (!inPackfile) {
Expand All @@ -329,19 +346,40 @@ 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)}`);
} else if (band !== BAND_PROGRESS) {
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<Uint8Array>)
: Promise<ReadableStreamReadResult<Uint8Array>> {
let stall!: (error: Error) => void;
let stalled = new Promise<never>((_, 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

Expand Down Expand Up @@ -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<GitOid>;
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${
Expand Down
Loading
Loading