Skip to content
Merged
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
20 changes: 9 additions & 11 deletions src/jsc/bindings/webcore/streams/JSDirectStreamController.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -435,17 +435,6 @@ static void callUnderlyingSourceClose(JSC::VM& vm, JSGlobalObject* globalObject,
}
}

// Drop the direct controller's retained user-source state once no further pull/close callbacks
// can run (m_closed set, or the stream has left Readable). Idempotent.
static void directStreamControllerClearSource(JSDirectStreamController* controller)
{
controller->m_underlyingSource.clear();
controller->m_pull.clear();
controller->m_deferCloseReason.clear();
if (auto* stream = controller->m_stream.get())
readableStreamClearSourceBarriers(stream);
}

void JSDirectStreamController::handleError(JSGlobalObject* globalObject, JSValue error)
{
auto& vm = getVM(globalObject);
Expand Down Expand Up @@ -992,6 +981,15 @@ using namespace JSC;
using WebCore::JSDirectStreamController;
using WebCore::JSStreamsRuntime;

void directStreamControllerClearSource(JSDirectStreamController* controller)
{
controller->m_underlyingSource.clear();
controller->m_pull.clear();
controller->m_deferCloseReason.clear();
if (auto* stream = controller->m_stream.get())
readableStreamClearSourceBarriers(stream);
}

void setUpDirectStreamController(JSC::JSGlobalObject* globalObject, JSReadableStream* stream, DirectSinkKind sinkKind, double highWaterMark)
{
auto& vm = getVM(globalObject);
Expand Down
7 changes: 3 additions & 4 deletions src/jsc/bindings/webcore/streams/ReadableStreamOperations.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -447,9 +447,8 @@ JSPromise* readableStreamCancel(JSGlobalObject* globalObject, JSReadableStream*
pendingRead->fulfill(vm, doneResult);
RETURN_IF_EXCEPTION(scope, nullptr);
}
controller->m_underlyingSource.clear();
controller->m_pull.clear();
controller->m_deferCloseReason.clear();
controller->m_closed = true;
directStreamControllerClearSource(controller);
sourceCancelPromise = promiseFulfilledWith(globalObject, JSC::jsUndefined());
break;
}
Expand Down Expand Up @@ -550,7 +549,7 @@ void readableStreamReaderGenericRelease(JSGlobalObject* globalObject, JSReadable
auto* controller = defaultControllerOf(stream);
controller->releaseSteps();
// Bun: drop the native handle's event-loop ref when its consumer releases the lock.
if (stream->m_nativePtr && controller->m_algorithms.kind == SourceKind::Native) {
if (controller->m_algorithms.kind == SourceKind::Native) {
const auto* adapter = uncheckedDowncast<WebCore::JSNativeStreamSourceAdapter>(controller->m_algorithms.algorithmContext.get());
if (auto* handle = adapter->handle()) {
JSValue updateRef = handle->getIfPropertyExists(globalObject, builtinNames(vm).updateRefPublicName());
Expand Down
3 changes: 3 additions & 0 deletions src/jsc/bindings/webcore/streams/WebStreamsInternals.h
Original file line number Diff line number Diff line change
Expand Up @@ -540,6 +540,9 @@ JSC::JSPromise* readStreamIntoSink(JSC::JSGlobalObject*, JSReadableStream*, JSC:
// Installs a JSDirectStreamController of the given flavor on the stream, nulls the stream's
// m_directUnderlyingSource, and sets m_bunMode = Default.
void setUpDirectStreamController(JSC::JSGlobalObject*, JSReadableStream*, DirectSinkKind, double highWaterMark); // userJS: yes — JSDirectStreamController.cpp
// Drop the direct controller's retained user-source state once no further pull/close callbacks
// can run (m_closed set, or the stream has left Readable). Idempotent.
Comment thread
robobun marked this conversation as resolved.
void directStreamControllerClearSource(JSDirectStreamController*); // userJS: no — JSDirectStreamController.cpp

// BunStreamConsumers.cpp — Bun.readableStreamTo*, the buffered fast path, the direct
// consumers, and the generic accumulators. These are the native entry points; their
Expand Down
243 changes: 93 additions & 150 deletions test/js/web/streams/readable-stream-terminal-barrier-release.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { describe, expect, test } from "bun:test";
import { expect, test } from "bun:test";
import { bunEnv, bunExe } from "harness";

// A ReadableStream captures the ambient AsyncLocalStorage context at construction
Expand All @@ -12,165 +12,108 @@ import { bunEnv, bunExe } from "harness";
// WriteBarrier. Before this change the probe survived every GC; with eager
// clearing it is collectable once the terminal transition completes.

const timeout = 60_000;

async function run(src: string) {
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", src],
env: bunEnv,
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect(stderr).toBe("");
expect(exitCode).toBe(0);
return JSON.parse(stdout.trim()) as { alive: number; total: number };
}

// Shared child-process scaffolding: runs `body` N times under an AsyncLocalStorage
// store keyed by a fresh probe object tracked by a FinalizationRegistry, awaits
// microtasks, then GC-storms and reports how many probes survived.
function fixture(body: string) {
return `
test("ReadableStream releases source-only WriteBarriers once terminal", async () => {
const src = `
const { AsyncLocalStorage } = require("node:async_hooks");
const als = new AsyncLocalStorage();
const N = 200;
const N = 20;
const retained = [];
let alive = 0;
const registry = new FinalizationRegistry(() => { alive--; });
const results = {};
const registry = new FinalizationRegistry(name => { results[name]--; });

async function alsCase(name, body) {
results[name] = 0;
for (let i = 0; i < N; i++) {
const probe = { i };
results[name]++;
registry.register(probe, name);
await als.run({ probe }, () => body(retained));
}
}

await alsCase("default-cancel", async retained => {
const rs = new ReadableStream({ pull() {} });
retained.push(rs);
await rs.cancel();
});
await alsCase("default-error", async retained => {
let ctrl;
const rs = new ReadableStream({ start(c) { ctrl = c; } });
retained.push(rs, ctrl);
ctrl.error(new Error("boom"));
});
await alsCase("default-close", async retained => {
let ctrl;
const rs = new ReadableStream({ start(c) { ctrl = c; } });
retained.push(rs, ctrl);
ctrl.close();
});
await alsCase("byte-cancel", async retained => {
const rs = new ReadableStream({ type: "bytes", pull() {} });
retained.push(rs);
await rs.cancel();
});
await alsCase("byte-error", async retained => {
let ctrl;
const rs = new ReadableStream({ type: "bytes", start(c) { ctrl = c; } });
retained.push(rs, ctrl);
ctrl.error(new Error("boom"));
});
await alsCase("direct-end", async retained => {
const rs = new ReadableStream({ type: "direct", pull(c) { c.write("x"); c.end(); } });
retained.push(rs);
for await (const _ of rs) {}
});

// DirectPending underlyingSource (no ALS: the probe is the source object itself).
results["direct-pending-cancel"] = 0;
for (let i = 0; i < N; i++) {
const probe = { i };
alive++;
registry.register(probe, undefined);
await als.run({ probe }, async () => { ${body} });
const source = { type: "direct", pull() {} };
results["direct-pending-cancel"]++;
registry.register(source, "direct-pending-cancel");
const rs = new ReadableStream(source);
retained.push(rs);
await rs.cancel();
}
Comment thread
robobun marked this conversation as resolved.

for (let r = 0; r < 20; r++) { Bun.gc(true); await new Promise(f => setImmediate(f)); }
process.stdout.write(JSON.stringify({ alive, total: N }));
for (let r = 0; r < 8; r++) { Bun.gc(true); await new Promise(f => setImmediate(f)); }
process.stdout.write(JSON.stringify({ results, N }));
if (retained.length === 0) throw new Error("retained was cleared");
`;
}

describe.concurrent("ReadableStream releases its captured async context once terminal", () => {
test(
"default controller: cancel()",
async () => {
const { alive, total } = await run(
fixture(`
const rs = new ReadableStream({ pull() {} });
retained.push(rs);
await rs.cancel();
`),
);
expect(alive).toBeLessThan(total * 0.2);
},
timeout,
);

test(
"default controller: controller.error()",
async () => {
const { alive, total } = await run(
fixture(`
let ctrl;
const rs = new ReadableStream({ start(c) { ctrl = c; } });
retained.push(rs, ctrl);
ctrl.error(new Error("boom"));
`),
);
expect(alive).toBeLessThan(total * 0.2);
},
timeout,
);

test(
"default controller: controller.close()",
async () => {
const { alive, total } = await run(
fixture(`
let ctrl;
const rs = new ReadableStream({ start(c) { ctrl = c; } });
retained.push(rs, ctrl);
ctrl.close();
`),
);
expect(alive).toBeLessThan(total * 0.2);
},
timeout,
);

test(
"byte controller: cancel()",
async () => {
const { alive, total } = await run(
fixture(`
const rs = new ReadableStream({ type: "bytes", pull() {} });
retained.push(rs);
await rs.cancel();
`),
);
expect(alive).toBeLessThan(total * 0.2);
},
timeout,
);

test(
"byte controller: controller.error()",
async () => {
const { alive, total } = await run(
fixture(`
let ctrl;
const rs = new ReadableStream({ type: "bytes", start(c) { ctrl = c; } });
retained.push(rs, ctrl);
ctrl.error(new Error("boom"));
`),
);
expect(alive).toBeLessThan(total * 0.2);
},
timeout,
);

test(
"direct stream: cancel() before materialize drops the pending underlyingSource",
async () => {
// No ALS here: the probe is the underlyingSource object itself, which the
// DirectPending stream stores in m_directUnderlyingSource until first use.
const src = `
const N = 200;
const retained = [];
let alive = 0;
const registry = new FinalizationRegistry(() => { alive--; });
for (let i = 0; i < N; i++) {
const src = { type: "direct", pull() {} };
alive++;
registry.register(src, undefined);
const rs = new ReadableStream(src);
retained.push(rs);
await rs.cancel();
}
for (let r = 0; r < 20; r++) { Bun.gc(true); await new Promise(f => setImmediate(f)); }
process.stdout.write(JSON.stringify({ alive, total: N }));
if (retained.length === 0) throw new Error("retained was cleared");
`;
const { alive, total } = await run(src);
expect(alive).toBeLessThan(total * 0.2);
},
timeout,
);
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", src],
env: bunEnv,
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect(stderr).toBe("");
const { results, N } = JSON.parse(stdout.trim()) as { results: Record<string, number>; N: number };
// Before this change every case reported N live probes; allow generous GC slack.
const threshold = Math.floor(N * 0.2);
const leaked = Object.entries(results).filter(([, alive]) => alive >= threshold);
expect({ leaked, results, threshold }).toEqual({ leaked: [], results, threshold });
expect(exitCode).toBe(0);
});

test(
"direct controller: end() drops underlyingSource and async context",
async () => {
const { alive, total } = await run(
fixture(`
const src = { type: "direct", pull(c) { c.write("x"); c.end(); } };
const rs = new ReadableStream(src);
retained.push(rs);
for await (const _ of rs) {}
`),
);
expect(alive).toBeLessThan(total * 0.2);
// readableStreamCancel's Direct arm used to leave m_closed false (onClose early-returned
// because the stream was already Closed), so a captured controller could still reach
// writeToDirectSink after cancel. With m_closed set, the bound methods throw.
test("direct controller write() after reader.cancel() throws closed", async () => {
let ctrl: any;
let pullStarted = Promise.withResolvers<void>();
const rs = new ReadableStream({
type: "direct",
async pull(c: any) {
ctrl = c;
c.write("first");
pullStarted.resolve();
await new Promise(() => {});
},
timeout,
);
});
const reader = rs.getReader();
reader.read().catch(() => {});
await pullStarted.promise;
await reader.cancel();
expect(() => ctrl.write("after-cancel")).toThrow(/closed/);
});
Loading