diff --git a/src/core/anonymize-utils.ts b/src/core/anonymize-utils.ts index ad644f8..379f481 100644 --- a/src/core/anonymize-utils.ts +++ b/src/core/anonymize-utils.ts @@ -442,12 +442,31 @@ function replaceTerm(content: string, term: CompiledTermVariant): string { if (term.pattern instanceof RegExp) { 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[] = []; let cursor = 0; - while (matcher.find()) { - const start = matcher.start(); - const end = matcher.end(); + const matches = function* () { + if (candidates) { + 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 // 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; diff --git a/src/core/zipStream.ts b/src/core/zipStream.ts index 51e3baf..c4c5983 100644 --- a/src/core/zipStream.ts +++ b/src/core/zipStream.ts @@ -1,4 +1,5 @@ import got from "got"; +import { Readable, Transform } from "stream"; import { Parse } from "unzip-stream"; import archiver = require("archiver"); @@ -112,33 +113,24 @@ export async function streamAnonymizedZip( } 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 parser = Parse() as Transform; + const activeStreams = new Set(); const compiledTerms = compileTerms(opt.anonymizerOptions.terms || []); - - // Track whether the upstream zipball finished cleanly. If it didn't, - // we must NOT finalize the archive — finalizing while bytes are still - // flowing to the response produces a valid-looking ZIP that's missing - // 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)); + let stopped = false; + const stop = () => { + if (stopped) return; + stopped = true; downloadStream.destroy(); - downloadStream.unpipe(parser); + parser.destroy(); + for (const stream of activeStreams) stream.destroy(); + activeStreams.clear(); archive.abort(); + }; + const fail = (error: Error) => { + if (stopped) return; + logger.error("zip stream failed", serializeError(error)); + stop(); const destroyable = res as unknown as { destroy?: (err?: Error) => void; end?: () => void; @@ -149,14 +141,23 @@ export async function streamAnonymizedZip( 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.pipe(res); downloadStream .on("error", fail) .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") { try { const fileName = anonymizePathCompiled( @@ -175,12 +176,12 @@ export async function streamAnonymizedZip( ...opt.anonymizerOptions, filePath: entry.path, }); - entry.on("error", fail); anonymizer.on("error", fail); + activeStreams.add(anonymizer); + anonymizer.once("close", () => activeStreams.delete(anonymizer)); const st = entry.pipe(anonymizer); archive.append(st, { name: fileName }); } catch (error) { - entry.autodrain(); fail(error as Error); } } else { @@ -189,13 +190,7 @@ export async function streamAnonymizedZip( }) .on("error", fail) .on("finish", () => { - if (failed) return; - try { - archive.finalize().catch(fail); - } catch (error) { - fail(error as Error); - } + if (stopped) return; + void archive.finalize().catch(fail); }); - - archive.pipe(res).on("error", fail); } diff --git a/test/large-notebook-download.test.js b/test/large-notebook-download.test.js new file mode 100644 index 0000000..73b0a18 --- /dev/null +++ b/test/large-notebook-download.test.js @@ -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); + }); + }); +});