Merge pull request #835 from tdurieux/codex/fix-large-notebook-downloads

fix: prevent large notebook download timeouts
This commit is contained in:
Thomas Durieux
2026-09-23 06:00:09 +02:00
committed by GitHub
3 changed files with 183 additions and 40 deletions
+23 -4
View File
@@ -442,12 +442,31 @@ function replaceTerm(content: string, term: CompiledTermVariant): string {
if (term.pattern instanceof RegExp) { if (term.pattern instanceof RegExp) {
return content.replace(term.pattern, () => term.mask); return content.replace(term.pattern, () => term.mask);
} }
const matcher = term.pattern.matcher(content); // Fixed-width literal searches run natively. Validate each candidate with
// RE2 on a small window to preserve its case folding and boundary semantics
// without scanning megabytes of notebook outputs in the JS regex engine.
const candidates = term.literalPrefilter
? content.matchAll(new RegExp(term.literalPrefilter.source, "giu"))
: null;
const matcher = candidates ? null : term.pattern.matcher(content);
const pieces: string[] = []; const pieces: string[] = [];
let cursor = 0; let cursor = 0;
while (matcher.find()) { const matches = function* () {
const start = matcher.start(); if (candidates) {
const end = matcher.end(); for (const candidate of candidates) {
const start = candidate.index!;
const end = start + candidate[0].length;
const offset = Math.max(0, start - 2);
const check = (term.pattern as RE2JS).matcher(content.slice(offset, end + 2));
if (check.find(start - offset) && check.start() === start - offset && check.end() === end - offset) {
yield { start, end };
}
}
} else {
while (matcher!.find()) yield { start: matcher!.start(), end: matcher!.end() };
}
};
for (const { start, end } of matches()) {
// RE2 has no lookahead. Check the generated Unicode word boundaries // RE2 has no lookahead. Check the generated Unicode word boundaries
// outside the engine, without executing any user-supplied native regex. // outside the engine, without executing any user-supplied native regex.
if (term.before && /[\p{L}\p{N}_]$/u.test(content.slice(Math.max(0, start - 2), start))) continue; if (term.before && /[\p{L}\p{N}_]$/u.test(content.slice(Math.max(0, start - 2), start))) continue;
+31 -36
View File
@@ -1,4 +1,5 @@
import got from "got"; import got from "got";
import { Readable, Transform } from "stream";
import { Parse } from "unzip-stream"; import { Parse } from "unzip-stream";
import archiver = require("archiver"); import archiver = require("archiver");
@@ -112,33 +113,24 @@ export async function streamAnonymizedZip(
} }
const downloadStream = got.stream(response.url); const downloadStream = got.stream(response.url);
res.on("error", (error) => {
logger.error("response stream error", serializeError(error));
downloadStream.destroy();
});
res.on("close", () => {
downloadStream.destroy();
});
const archive = archiver("zip", {}); const archive = archiver("zip", {});
const parser = Parse() as Transform;
const activeStreams = new Set<Readable>();
const compiledTerms = compileTerms(opt.anonymizerOptions.terms || []); const compiledTerms = compileTerms(opt.anonymizerOptions.terms || []);
let stopped = false;
// Track whether the upstream zipball finished cleanly. If it didn't, const stop = () => {
// we must NOT finalize the archive — finalizing while bytes are still if (stopped) return;
// flowing to the response produces a valid-looking ZIP that's missing stopped = true;
// entries, which the client has no way to detect (status 200, archive
// opens). Destroy the response instead so the client sees a connection
// drop and knows the download failed. Same class of silent-truncation
// bug as #694.
let failed = false;
const parser = Parse();
const fail = (error: Error) => {
if (failed) return;
failed = true;
logger.error("upstream zipball failed", serializeError(error));
downloadStream.destroy(); downloadStream.destroy();
downloadStream.unpipe(parser); parser.destroy();
for (const stream of activeStreams) stream.destroy();
activeStreams.clear();
archive.abort(); archive.abort();
};
const fail = (error: Error) => {
if (stopped) return;
logger.error("zip stream failed", serializeError(error));
stop();
const destroyable = res as unknown as { const destroyable = res as unknown as {
destroy?: (err?: Error) => void; destroy?: (err?: Error) => void;
end?: () => void; end?: () => void;
@@ -149,14 +141,23 @@ export async function streamAnonymizedZip(
destroyable.end(); destroyable.end();
} }
}; };
// pipe() returns the destination. Listen on the archive itself as well,
// including errors emitted after the upstream ZIP has finished downloading. res.on("error", (error) => fail(error as Error));
res.on("close", stop);
archive.on("error", fail); archive.on("error", fail);
archive.pipe(res);
downloadStream downloadStream
.on("error", fail) .on("error", fail)
.pipe(parser) .pipe(parser)
.on("entry", (entry: NodeJS.ReadableStream & { type: string; path: string; autodrain: () => void }) => { .on("entry", (entry: Readable & { type: string; path: string; autodrain: () => void }) => {
if (stopped) {
entry.destroy();
return;
}
entry.on("error", fail);
activeStreams.add(entry);
entry.once("close", () => activeStreams.delete(entry));
if (entry.type === "File") { if (entry.type === "File") {
try { try {
const fileName = anonymizePathCompiled( const fileName = anonymizePathCompiled(
@@ -175,12 +176,12 @@ export async function streamAnonymizedZip(
...opt.anonymizerOptions, ...opt.anonymizerOptions,
filePath: entry.path, filePath: entry.path,
}); });
entry.on("error", fail);
anonymizer.on("error", fail); anonymizer.on("error", fail);
activeStreams.add(anonymizer);
anonymizer.once("close", () => activeStreams.delete(anonymizer));
const st = entry.pipe(anonymizer); const st = entry.pipe(anonymizer);
archive.append(st, { name: fileName }); archive.append(st, { name: fileName });
} catch (error) { } catch (error) {
entry.autodrain();
fail(error as Error); fail(error as Error);
} }
} else { } else {
@@ -189,13 +190,7 @@ export async function streamAnonymizedZip(
}) })
.on("error", fail) .on("error", fail)
.on("finish", () => { .on("finish", () => {
if (failed) return; if (stopped) return;
try { void archive.finalize().catch(fail);
archive.finalize().catch(fail);
} catch (error) {
fail(error as Error);
}
}); });
archive.pipe(res).on("error", fail);
} }
+129
View File
@@ -0,0 +1,129 @@
const { expect } = require("chai");
const { Readable, PassThrough } = require("stream");
const { once } = require("events");
const archiver = require("archiver");
const { Parse } = require("unzip-stream");
require("ts-node/register/transpile-only");
const got = require("got");
const GitHubDownload = require("../src/core/source/GitHubDownload").default;
const config = require("../src/config").default;
const { AnonymizeTransformer, ContentAnonimizer } = require("../src/core/anonymize-utils");
const { streamAnonymizedZip } = require("../src/core/zipStream");
function notebook() {
return JSON.stringify({ nbformat: 4, nbformat_minor: 0, metadata: {}, cells: [{
cell_type: "code", metadata: {}, execution_count: 1,
source: ["# Alice's experiment"],
outputs: [{ output_type: "display_data", metadata: {}, data: {
"image/png": "AbCd".repeat(12 * 1024 * 1024),
"text/plain": ["Alice completed the experiment"],
} }],
}] });
}
async function collect(stream) {
const chunks = [];
for await (const chunk of stream) chunks.push(chunk);
return Buffer.concat(chunks);
}
async function zip(entries) {
const archive = archiver("zip");
const result = collect(archive);
for (const [name, content] of Object.entries(entries)) archive.append(content, { name });
await archive.finalize();
return result;
}
async function unzip(buffer) {
const parser = Parse();
const entries = {};
const pending = [];
parser.on("entry", entry => pending.push(collect(entry).then(data => { entries[entry.path] = data; })));
const done = once(parser, "finish");
Readable.from([buffer]).pipe(parser);
await done;
await Promise.all(pending);
return entries;
}
describe("large notebook downloads", function () {
this.timeout(15000);
it("redacts names before and after large outputs within the anonymization deadline", async function () {
const input = notebook();
const transformer = new AnonymizeTransformer({ filePath: "example.ipynb", terms: ["Alice"] });
const output = collect(transformer);
Readable.from([Buffer.from(input)]).pipe(transformer);
const result = JSON.parse((await output).toString());
expect(result.cells[0].source[0]).to.equal("# XXXX-1's experiment");
expect(result.cells[0].outputs[0].data["text/plain"][0]).to.equal("XXXX-1 completed the experiment");
expect(result.cells[0].outputs[0].data["image/png"]).to.equal("AbCd".repeat(12 * 1024 * 1024));
});
it("matches RE2 boundaries and case folding for nearby literal candidates", function () {
for (const term of ["k", "s", "Alice", "a-a", "Σ", "研究", "@Alice", "😀", "---", "@@", " "]) {
const input = ["", " ", "é", "_", "😀", "𐐀"].flatMap(left =>
["", " ", "é", "_", "😀", "𐐀"].map(right => `${left}${term} ${term.toUpperCase()} K ſ${right}`)
).join(" ") + term.repeat(10);
const reference = new ContentAnonimizer({ terms: [term] });
for (const compiled of reference.compiledTerms) delete compiled.literalPrefilter;
expect(new ContentAnonimizer({ terms: [term] }).anonymize(input)).to.equal(reference.anonymize(input));
}
});
describe("ZIP streaming", function () {
let originalStream, originalUrl, originalLimit, upstream;
beforeEach(function () {
originalStream = got.stream;
originalUrl = GitHubDownload.prototype.getZipUrl;
originalLimit = config.MAX_FILE_SIZE;
GitHubDownload.prototype.getZipUrl = async () => ({ url: "https://example.invalid/source.zip" });
got.stream = () => upstream;
});
afterEach(function () {
got.stream = originalStream;
GitHubDownload.prototype.getZipUrl = originalUrl;
config.MAX_FILE_SIZE = originalLimit;
upstream?.destroy();
});
const options = {
repoId: "fixture", organization: "owner", repoName: "repo", commit: "HEAD",
getToken: () => "test", anonymizerOptions: { filePath: "", terms: ["Alice"] },
};
it("completes an archive containing a large notebook and subsequent entries", async function () {
upstream = Readable.from([await zip({ "root/example.ipynb": notebook(), "root/after.txt": "Alice" })]);
const response = new PassThrough();
const output = collect(response);
await streamAnonymizedZip(options, response);
const entries = await unzip(await output);
expect(Object.keys(entries)).to.have.members(["example.ipynb", "after.txt"]);
expect(entries["after.txt"].toString()).to.equal("XXXX-1");
expect(JSON.parse(entries["example.ipynb"].toString()).cells[0].source[0]).to.equal("# XXXX-1's experiment");
});
it("aborts the response when an entry cannot be anonymized", async function () {
upstream = Readable.from([await zip({ "root/large.txt": "Alice".repeat(1000) })]);
config.MAX_FILE_SIZE = 100;
const response = new PassThrough();
response.resume();
const error = once(response, "error");
await streamAnonymizedZip(options, response);
expect((await error)[0].message).to.contain("Text file exceeded");
expect(response.destroyed).to.equal(true);
expect(response.writableFinished).to.equal(false);
expect(upstream.destroyed).to.equal(true);
});
it("stops downloading when the client disconnects", async function () {
upstream = new PassThrough();
const response = new PassThrough();
await streamAnonymizedZip(options, response);
const closed = once(response, "close");
response.destroy();
await closed;
expect(upstream.destroyed).to.equal(true);
});
});
});