diff --git a/src/jsc/bindings/webcore/streams/BunStreamConsumers.cpp b/src/jsc/bindings/webcore/streams/BunStreamConsumers.cpp index f510e1a76142..2d6081fc688a 100644 --- a/src/jsc/bindings/webcore/streams/BunStreamConsumers.cpp +++ b/src/jsc/bindings/webcore/streams/BunStreamConsumers.cpp @@ -540,6 +540,60 @@ static JSValue convertChunksToBytes(JSGlobalObject* globalObject, JSValue chunks static JSValue textAccumulatorWrite(JSC::VM& vm, JSGlobalObject*, JSC::JSObject* owner, BunTextAccumulator&, JSValue chunk); static WTF::String finishTextAccumulator(JSC::VM& vm, JSGlobalObject*, JSC::JSObject* owner, BunTextAccumulator&); +// How many U+FEFF code units, at most two, start the text of the string chunks. The text result leaves them out. +static unsigned leadingBOMCount(JSGlobalObject* globalObject, const MarkedArgumentBuffer& chunks) +{ + auto scope = DECLARE_THROW_SCOPE(getVM(globalObject)); + unsigned count = 0; + for (unsigned i = 0; i < chunks.size() && count < 2; i++) { + JSString* chunk = asString(chunks.at(i)); + if (!chunk->length()) + continue; + if (chunk->is8Bit()) + break; + auto view = chunk->view(globalObject); + RETURN_IF_EXCEPTION(scope, 0); + unsigned inChunk = 0; + while (inChunk < view->length() && count < 2 && view[inChunk] == 0xFEFF) { + inChunk++; + count++; + } + if (inChunk < view->length()) + break; + } + return count; +} + +// The string chunks as one string without its first `skip` code units, from one allocation. Null when that fails. +template +static WTF::String tryJoinStringChunks(JSGlobalObject* globalObject, const MarkedArgumentBuffer& chunks, unsigned length, unsigned skip) +{ + auto scope = DECLARE_THROW_SCOPE(getVM(globalObject)); + std::span characters; + WTF::String joined = WTF::String::tryCreateUninitialized(length - skip, characters); + if (joined.isNull()) [[unlikely]] + return joined; + for (unsigned i = 0; i < chunks.size(); i++) { + JSString* chunk = asString(chunks.at(i)); + unsigned chunkLength = chunk->length(); + if (skip && chunkLength) [[unlikely]] { + unsigned skipped = std::min(skip, chunkLength); + skip -= skipped; + if (skipped < chunkLength) { + auto view = chunk->view(globalObject); + RETURN_IF_EXCEPTION(scope, {}); + view->substring(skipped).getCharacters(characters); + characters = characters.subspan(chunkLength - skipped); + } + continue; + } + // This copies a rope from its fibers. It does not make the string of the rope. + chunk->resolveToBuffer(characters.first(chunkLength)); + characters = characters.subspan(chunkLength); + } + return joined; +} + // The chunk-array -> text conversion: pure-string arrays join once (no UTF-8 round trip); // mixed/binary chunk arrays run through the shared text accumulator. static JSValue convertChunksToText(JSGlobalObject* globalObject, JSValue chunksValue) @@ -585,6 +639,7 @@ static JSValue convertChunksToText(JSGlobalObject* globalObject, JSValue chunksV // MarkedArgumentBuffer for the conversion below. MarkedArgumentBuffer values; bool allStrings = true; + bool all8Bit = true; WTF::CheckedUint32 codeUnits = 0; for (unsigned i = 0; i < length; i++) { JSValue chunk = chunks->getIndex(globalObject, i); @@ -592,8 +647,10 @@ static JSValue convertChunksToText(JSGlobalObject* globalObject, JSValue chunksV values.append(chunk); if (!chunk.isString()) allStrings = false; - else if (allStrings) + else if (allStrings) { codeUnits += asString(chunk)->length(); + all8Bit = all8Bit && asString(chunk)->is8Bit(); + } } if (values.hasOverflowed()) [[unlikely]] { throwOutOfMemoryError(globalObject, scope); @@ -604,18 +661,17 @@ static JSValue convertChunksToText(JSGlobalObject* globalObject, JSValue chunksV throwOutOfMemoryError(globalObject, scope); return {}; } - WTF::StringBuilder rope; - rope.reserveCapacity(codeUnits.value()); - for (unsigned i = 0; i < length; i++) { - WTF::String string = asString(values.at(i))->value(globalObject); - RETURN_IF_EXCEPTION(scope, {}); - rope.append(string); - } - if (rope.hasOverflowed()) [[unlikely]] { + unsigned skip = all8Bit ? 0 : leadingBOMCount(globalObject, values); + RETURN_IF_EXCEPTION(scope, {}); + WTF::String joined = all8Bit + ? tryJoinStringChunks(globalObject, values, codeUnits.value(), skip) + : tryJoinStringChunks(globalObject, values, codeUnits.value(), skip); + RETURN_IF_EXCEPTION(scope, {}); + if (joined.isNull()) [[unlikely]] { throwOutOfMemoryError(globalObject, scope); return {}; } - RELEASE_AND_RETURN(scope, jsString(vm, stripTextResultBOM(rope.toString()))); + RELEASE_AND_RETURN(scope, jsString(vm, WTF::move(joined))); } // Mixed string/binary chunks: drive the shared accumulator so adjacent-string rope diff --git a/test/harness.ts b/test/harness.ts index 60fffcbc865c..3e9a7f21c043 100644 --- a/test/harness.ts +++ b/test/harness.ts @@ -338,6 +338,29 @@ export async function runFixtureMaxRSS(fixture: string, expected: unknown) { return maxRSS; } +/** + * The env of a child whose allocator refuses one allocation of more than + * `megabytes` MiB. It needs an ASAN build: use it under `skipIf(!isASAN)`. + * `Malloc=1` is not optional. Without it WebKit takes the buffers of its + * strings from bmalloc, which the cap does not reach, and a test that expects + * a refusal passes or fails for another reason. ASAN logs every allocation + * that it refuses to stderr. + */ +export function allocationCapEnv(megabytes: number): typeof bunEnv { + return { + ...bunEnv, + Malloc: "1", + ASAN_OPTIONS: [ + bunEnv.ASAN_OPTIONS, + "allocator_may_return_null=1", + `max_allocation_size_mb=${megabytes}`, + "detect_leaks=0", + ] + .filter(Boolean) + .join(":"), + }; +} + /** * Runs `cmd` (a script that prints `{"deltaMiB": number}` as its last stdout * line) under bun with ASAN quarantine disabled, and asserts the delta is below diff --git a/test/js/web/streams/streams-string-limit.test.ts b/test/js/web/streams/streams-string-limit.test.ts index 956b81189b62..ae7ce6146458 100644 --- a/test/js/web/streams/streams-string-limit.test.ts +++ b/test/js/web/streams/streams-string-limit.test.ts @@ -1,5 +1,5 @@ import { describe, expect, test } from "bun:test"; -import { bunEnv, bunExe } from "harness"; +import { allocationCapEnv, bunEnv, bunExe, emptyProcessMaxRSS, isASAN, isDebug, runFixtureMaxRSS } from "harness"; import { totalmem } from "node:os"; // Consuming a stream as text must reject with a catchable error when the accumulated @@ -21,10 +21,13 @@ function consumeToText(streamSource: string): string { `; } -async function run(script: string): Promise<{ stdout: string; stderr: string; exitCode: number }> { +async function run( + script: string, + env: Record = bunEnv, +): Promise<{ stdout: string; stderr: string; exitCode: number }> { await using proc = Bun.spawn({ cmd: [bunExe(), "-e", script], - env: bunEnv, + env, stdout: "pipe", stderr: "pipe", }); @@ -184,6 +187,126 @@ test.skipIf(!enoughMemory)("arrayBuffer() and bytes() reject mixed chunks summin }); }); +// The text of a stream whose chunks are all strings is one string of the sum of their lengths. +// The consumer joined them in a WTF::StringBuilder that aborts the process when it cannot grow: +// when the allocator refuses its buffer, and when it doubles the buffer for the first 16-bit +// chunk and the double is longer than a 16-bit string can be. +describe("the text of string chunks is one allocation of its length", () => { + const MIB = 1024 * 1024; + const outOfMemory = "RangeError: Out of memory"; + // "aaabbc" is "a3 b2 c1". + const runs = `text => text.replace(/(.)\\1*/gs, (run, character) => character + run.length + " ").trim()`; + const consume = (chunks: string, text: string, describeText = runs) => ` + const megabyte = letter => Buffer.alloc(1024 * 1024, letter).toString("latin1"); + const chunks = ${chunks}; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + const describeText = ${describeText}; + const settled = await ${text}.then(text => ({ text: describeText(text) }), e => ({ rejected: e.name + ": " + e.message })); + console.log(JSON.stringify(settled)); + `; + + describe.skipIf(!isASAN)("under a cap of 4 MiB for one allocation", () => { + const env = allocationCapEnv(4); + const megabytes = (letters: string) => `[...${JSON.stringify(letters)}].map(megabyte)`; + + test.concurrent.each([ + ["three megabytes", megabytes("abc"), { text: `a${MIB} b${MIB} c${MIB}` }], + ["four megabytes", megabytes("abcd"), { rejected: outOfMemory }], + // A megabyte of 16-bit characters is 2 MiB. + ["one megabyte, then a 16-bit character", `[megabyte("a"), "\\u20AC"]`, { text: `a${MIB} \u20AC1` }], + ["three megabytes, then a 16-bit character", `[...${megabytes("abc")}, "\\u20AC"]`, { rejected: outOfMemory }], + ["a 16-bit character, then three megabytes", `["\\u20AC", ...${megabytes("abc")}]`, { rejected: outOfMemory }], + ])("%s", async (_name, chunks, expected) => { + const { stdout, exitCode } = await run(consume(chunks, "Bun.readableStreamToText(stream)"), env); + expect({ stdout: JSON.parse(stdout || "null"), exitCode }).toEqual({ stdout: expected, exitCode: 0 }); + }); + + test.concurrent.each([ + "stream.text()", + "new Response(stream).text()", + "stream.json()", + "new Response(stream).json()", + ])("%s of four megabytes", async text => { + const { stdout, exitCode } = await run(consume(megabytes("abcd"), text), env); + expect({ stdout: JSON.parse(stdout || "null"), exitCode }).toEqual({ + stdout: { rejected: outOfMemory }, + exitCode: 0, + }); + }); + }); + + // A 16-bit string holds at most 2,147,483,635 characters. The consumer refuses a longer text + // before it allocates, so the child stays small. + test("a 16-bit text of 2,147,483,636 characters", async () => { + const chunks = `["\\u20AC", ...Array(2047).fill(megabyte("x")), megabyte("x").slice(0, 2 ** 20 - 13)]`; + const { stdout, stderr, exitCode } = await run(consume(chunks, "Bun.readableStreamToText(stream)")); + expect({ stdout: JSON.parse(stdout || "null"), stderr, exitCode }).toEqual({ + stdout: { rejected: outOfMemory }, + stderr: "", + exitCode: 0, + }); + }); + + // This text fits in a string: it is 2 GiB in 16 bits. A debug build takes 6 s to copy it. + test.skipIf(!enoughMemory || isDebug)("a 16-bit character after a gigabyte of Latin-1", async () => { + const chunks = `[...Array(16).fill(Buffer.alloc(2 ** 26, "x").toString("latin1")), "\\u20AC"]`; + const ends = `text => ({ length: text.length, start: text.slice(0, 3), end: text.slice(-3) })`; + const { stdout, stderr, exitCode } = await run(consume(chunks, "Bun.readableStreamToText(stream)", ends)); + expect({ stdout: JSON.parse(stdout || "null"), stderr, exitCode }).toEqual({ + stdout: { text: { length: 2 ** 30 + 1, start: "xxx", end: "xx\u20AC" } }, + stderr: "", + exitCode: 0, + }); + }); + + // The peak RSS of a child that reads a stream of string chunks as text, above the peak of an + // empty child, in MiB. + const peakOfText = async (chunks: string, report: string, expected: unknown) => { + const fixture = ` + const latin1 = (length, letter) => Buffer.alloc(length, letter).toString("latin1"); + ${chunks} + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + const text = await stream.text(); + console.log(JSON.stringify(${report})); + `; + const [peak, emptyPeak] = await Promise.all([runFixtureMaxRSS(fixture, expected), emptyProcessMaxRSS()]); + return (peak - emptyPeak) / MIB; + }; + + // The join leaves the BOM out. No copy removes it, and the text is not a part of another + // string, so structuredClone() shares it. The text is 256 MiB in 16 bits. A copy makes 512 MiB. + test("a BOM before 128 MiB of Latin-1", async () => { + const peak = await peakOfText( + `const chunks = ["\\uFEFF", ...Array(128).fill(latin1(${MIB}, "x"))];`, + `{ length: text.length, start: text.slice(0, 3), clone: structuredClone(text).length }`, + { length: 128 * MIB, start: "xxx", clone: 128 * MIB }, + ); + expect(peak).toBeLessThan(384); + }); + + // The join copies a chunk that is a rope from the two strings of the rope. It does not make + // the string of each rope first, which is 128 MiB more. + test("128 MiB of chunks that are ropes", async () => { + const peak = await peakOfText( + `const a = latin1(${MIB / 2}, "a"), b = latin1(${MIB / 2}, "b"); + const chunks = Array.from({ length: 128 }, () => a + b);`, + `{ length: text.length, start: text.slice(0, 3), end: text.slice(-3) }`, + { length: 128 * MIB, start: "aaa", end: "bbb" }, + ); + expect(peak).toBeLessThan(192); + }); +}); + // TextDecoderStream joins a chunk with the bytes it carried over from an incomplete UTF-8 // sequence. The join sized a WTF::Vector that aborts past 2^31-1 bytes, so a 2^31-byte chunk // after a carried byte killed the process before it read a byte. The child reserves the diff --git a/test/js/web/streams/streams.test.js b/test/js/web/streams/streams.test.js index ed0b09ce8e81..c6859c306f9b 100644 --- a/test/js/web/streams/streams.test.js +++ b/test/js/web/streams/streams.test.js @@ -1895,6 +1895,35 @@ describe("multi-chunk consumers produce exactly the concatenated bytes", () => { await expect(new Response(source([42])).text()).rejects.toThrow(expect.objectContaining({ name: "TypeError" })); }); + // The text of string chunks is one string of the sum of their lengths, 16-bit when a chunk is. + const unresolvedRopes = () => { + const part = Buffer.alloc(40, "p").toString(); + return [1, 2, 3].map(i => part + i + (i === 2 ? "\u20AC" : "\u00E9") + part); + }; + it.each([ + ["Latin-1 chunks", () => ["caf\u00E9", "", " au lait"], "caf\u00E9 au lait"], + ["a 16-bit chunk after Latin-1 chunks", () => ["abc", "d\u00E9", "\u20AC"], "abcd\u00E9\u20AC"], + ["a Latin-1 chunk after a 16-bit chunk", () => ["\u4F60\u597D", "abc"], "\u4F60\u597Dabc"], + ["only empty chunks", () => ["", ""], ""], + ["chunks that are ropes", unresolvedRopes, unresolvedRopes().join("")], + ["a BOM before the text", () => ["\uFEFF", "abc"], "abc"], + ["two BOMs before the text", () => ["\uFEFF", "\uFEFFabc"], "abc"], + ["three BOMs before the text", () => ["\uFEFF\uFEFF", "\uFEFFabc"], "\uFEFFabc"], + ["only a BOM", () => ["\uFEFF", ""], ""], + ["two BOMs and text in one chunk", () => ["\uFEFF\uFEFFab", "c"], "abc"], + ["BOMs with empty chunks between them", () => ["", "\uFEFF", "", "\uFEFF", "abc"], "abc"], + ["a BOM after the first character", () => ["a", "\uFEFFbc"], "a\uFEFFbc"], + ["a BOM at the start of a rope", () => ["\uFEFF" + unresolvedRopes()[1], "abc"], unresolvedRopes()[1] + "abc"], + [ + "a BOM at the start of a part of another string", + () => [("x\uFEFF" + unresolvedRopes()[0]).slice(1), "abc"], + unresolvedRopes()[0] + "abc", + ], + ])("text: string chunks join in order: %s", async (_name, chunks, text) => { + expect(await Bun.readableStreamToText(source(chunks()))).toBe(text); + expect(await new Response(source(chunks())).text()).toBe(text); + }); + it("a detached chunk throws", () => { const chunk = new Uint8Array([1, 2, 3]); structuredClone(chunk.buffer, { transfer: [chunk.buffer] });