Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
76 changes: 66 additions & 10 deletions src/jsc/bindings/webcore/streams/BunStreamConsumers.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Comment thread
robobun marked this conversation as resolved.
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<typename CharacterType>
static WTF::String tryJoinStringChunks(JSGlobalObject* globalObject, const MarkedArgumentBuffer& chunks, unsigned length, unsigned skip)
{
auto scope = DECLARE_THROW_SCOPE(getVM(globalObject));
std::span<CharacterType> 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)
Expand Down Expand Up @@ -585,15 +639,18 @@ 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);
RETURN_IF_EXCEPTION(scope, {});
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);
Expand All @@ -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<Latin1Character>(globalObject, values, codeUnits.value(), skip)
: tryJoinStringChunks<char16_t>(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
Expand Down
23 changes: 23 additions & 0 deletions test/harness.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
129 changes: 126 additions & 3 deletions test/js/web/streams/streams-string-limit.test.ts
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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<string, string | undefined> = bunEnv,
): Promise<{ stdout: string; stderr: string; exitCode: number }> {
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", script],
env: bunEnv,
env,
stdout: "pipe",
stderr: "pipe",
});
Expand Down Expand Up @@ -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);
Comment thread
robobun marked this conversation as resolved.
});
});

// 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
Expand Down
29 changes: 29 additions & 0 deletions test/js/web/streams/streams.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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] });
Expand Down
Loading