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
40 changes: 31 additions & 9 deletions src/js/internal/streams/native-readable.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,14 @@ let dynamicallyAdjustChunkSize = (_?) => (
type NodeReadable = import("node:stream").Readable;

interface NativeReadable extends NodeReadable {
_readableState: {
flowing: boolean | null;
ended: boolean;
sync: boolean;
buffer: unknown[];
bufferIndex: number;
length: number;
};
$bunNativePtr: NativePtr | undefined;
$start?: typeof ensureConstructed;
ref: typeof ref;
Expand Down Expand Up @@ -193,9 +201,7 @@ function handleResult(stream: NativeReadable, result: any, chunk: Buffer | undef
return handleNumberResult(stream, result, chunk, isClosed);
} else if (typeof result === "boolean") {
$debug(`[${stream.debugId}] handleResult(${result})`, chunk, isClosed);
process.nextTick(() => {
stream.push(null);
});
process.nextTick(pushEof, stream);
return (chunk?.byteLength ?? 0) > 0 ? chunk : undefined;
} else if ($isTypedArrayView(result)) {
if (result.byteLength >= stream[kHighWaterMark] && !stream[kHasResized] && !isClosed) {
Expand All @@ -207,6 +213,11 @@ function handleResult(stream: NativeReadable, result: any, chunk: Buffer | undef
}
}

// EOF is pushed a tick after the last chunk. After a destroy() in between, Node emits 'close' without 'end'.
function pushEof(stream: NativeReadable) {
if (!stream.destroyed) stream.push(null);
}

// `push()` returning false means the Readable's buffer is at/above hwm (or
// the consumer paused); stop the native reader so kernel backpressure reaches
// the writer (readStop, like net.Socket). The next `_read()` re-enables it.
Expand All @@ -227,9 +238,7 @@ function handleNumberResult(stream: NativeReadable, result: number, chunk: any,
}

if (isClosed) {
process.nextTick(() => {
stream.push(null);
});
process.nextTick(pushEof, stream);
}

return chunk;
Expand All @@ -241,9 +250,7 @@ function handleArrayBufferViewResult(stream: NativeReadable, result: any, chunk:
}

if (isClosed) {
process.nextTick(() => {
stream.push(null);
});
process.nextTick(pushEof, stream);
}

return chunk;
Expand All @@ -259,12 +266,27 @@ function destroy(this: NativeReadable, error: any, cb: () => void) {
if (ptr) {
ptr.cancel(error);
}
dropReadAhead(this);
if (cb) {
// `_destroy` reports its error through the callback.
process.nextTick(cb, error);
}
}

// `_read()` pushes synchronously, so flow() stays one chunk ahead of the 'data' listener. Node's async sources do not.
function dropReadAhead(stream: NativeReadable) {
const state = stream._readableState;
// Paused: Node has this buffered too, and a later read() returns it.
if (!state.flowing) return;
// Ended: the buffer is all that is left, and 'end' must not follow dropped data.
if (state.ended) return;
// Inside `_read()` the source failed: the bytes it read before the error are still delivered.
if (state.sync) return;
state.buffer.length = 0;
state.bufferIndex = 0;
state.length = 0;
}

function ref(this: NativeReadable) {
const ptr = this.$bunNativePtr;
if (ptr === undefined) return;
Expand Down
53 changes: 50 additions & 3 deletions test/js/bun/spawn/spawn-stdio-syscall-error.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,13 @@ const cc = Bun.which("cc") || Bun.which("gcc") || Bun.which("clang");

// SPAWN_FAULT_RECV_AT=N the Nth recv() on each AF_UNIX socket fails with EIO (1-based).
// SPAWN_FAULT_SEND_AT=N the Nth send() on each AF_UNIX socket fails with ENOBUFS.
// SPAWN_FAULT_REPORT=path the failing recv() writes how many bytes that socket received before it to this file.
const SHIM_C = /* c */ `
#define _GNU_SOURCE
#include <dlfcn.h>
#include <errno.h>
#include <stdarg.h>
#include <stdio.h>
#include <stdlib.h>
#include <sys/socket.h>
#include <sys/stat.h>
Expand All @@ -37,6 +39,7 @@ static int fail_recv_at = -1;
static int fail_send_at = -1;
static unsigned int recv_count[MAX_FD];
static unsigned int send_count[MAX_FD];
static unsigned long recv_bytes[MAX_FD];

static void init_modes(void) {
const char *s;
Expand Down Expand Up @@ -64,9 +67,20 @@ ssize_t recv(int fd, void *buf, size_t len, int flags) {
real_recv = (ssize_t (*)(int, void *, size_t, int))dlsym(RTLD_NEXT, "recv");
init_modes();
}
if (fail_recv_at > 0 && is_unix_sock(fd) && ++recv_count[fd] == (unsigned)fail_recv_at) {
errno = EIO;
return -1;
if (fail_recv_at > 0 && is_unix_sock(fd)) {
if (++recv_count[fd] == (unsigned)fail_recv_at) {
const char *report = getenv("SPAWN_FAULT_REPORT");
FILE *f = report ? fopen(report, "w") : NULL;
if (f) {
fprintf(f, "%lu", recv_bytes[fd]);
fclose(f);
}
errno = EIO;
return -1;
}
ssize_t n = real_recv(fd, buf, len, flags);
if (n > 0) recv_bytes[fd] += (unsigned long)n;
return n;
}
return real_recv(fd, buf, len, flags);
}
Expand All @@ -90,6 +104,7 @@ static void reset_fd(int fd) {
if (fd >= 0 && fd < MAX_FD) {
recv_count[fd] = 0;
send_count[fd] = 0;
recv_bytes[fd] = 0;
}
}

Expand Down Expand Up @@ -204,6 +219,23 @@ child.on("close", () => {
});
`;

// node:child_process: every byte the parent read before the error reaches 'data'.
const CHILD_PROCESS_BYTES_FIXTURE = /* js */ `
import { spawn } from "node:child_process";
import { readFileSync } from "node:fs";
const events = [];
let got = 0;
const child = spawn(${JSON.stringify(WRITER_CMD[0])}, ${JSON.stringify(WRITER_CMD.slice(1))}, { stdio: ["ignore", "pipe", "ignore"] });
child.stdout.on("data", chunk => (got += chunk.length));
child.stdout.on("error", e => events.push("stdout.error:" + e.code));
child.stdout.on("close", () => events.push("stdout.close"));
child.on("close", () => {
events.push("close");
const received = Number(readFileSync(process.env.SPAWN_FAULT_REPORT, "utf8"));
console.log(JSON.stringify({ receivedSome: received > 0, lost: received - got, events }));
});
`;

// Bun.spawnSync / child_process.spawnSync / execFileSync: the lost output is an error, not a success.
const SPAWN_SYNC_FIXTURE = /* js */ `
import { spawnSync, execFileSync } from "node:child_process";
Expand Down Expand Up @@ -241,6 +273,7 @@ beforeAll(async () => {
"stdout-text.mjs": STDOUT_TEXT_FIXTURE,
"stdout-write.mjs": STDOUT_WRITE_FIXTURE,
"child-process.mjs": CHILD_PROCESS_FIXTURE,
"child-process-bytes.mjs": CHILD_PROCESS_BYTES_FIXTURE,
"spawn-sync.mjs": SPAWN_SYNC_FIXTURE,
});
shimPath = join(String(dir), "shim.so");
Expand All @@ -266,6 +299,7 @@ async function runWithFault(fixture: string, fault: Record<string, string>) {
LD_PRELOAD: bunEnv.LD_PRELOAD ? `${shimPath}:${bunEnv.LD_PRELOAD}` : shimPath,
SPAWN_FAULT_RECV_AT: undefined,
SPAWN_FAULT_SEND_AT: undefined,
SPAWN_FAULT_REPORT: undefined,
...fault,
};
await using proc = Bun.spawn({
Expand Down Expand Up @@ -348,6 +382,19 @@ describe.skipIf(!isLinux || !cc)("subprocess stdio syscall errors", () => {
});
});

// child.stdout reads one chunk ahead of its 'data' listener, so once the stream flows a chunk is still buffered
// when a later read fails. The stream is destroyed with the error, and that chunk has to reach 'data' first.
test.concurrent("node:child_process: stdout delivers every byte read before the error", async () => {
const report = join(String(dir), "recv-report.txt");
expect(
await runWithFault("child-process-bytes.mjs", { SPAWN_FAULT_RECV_AT: "6", SPAWN_FAULT_REPORT: report }),
).toEqual({
parsed: { receivedSome: true, lost: 0, events: ["stdout.error:EIO", "stdout.close", "close"] },
stderr: "",
exitCode: 0,
});
});

describe.each([
["the first write", "1"],
["a later write", "3"],
Expand Down
76 changes: 76 additions & 0 deletions test/js/node/child_process/child_process.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1406,6 +1406,82 @@ it("child.stdout.pause() after flowing stops native reads and blocks the child",
}
});

// child.stdout and child.stderr read ahead: `_read()` pushes each native pull
// result synchronously, so while the stream flows Readable holds the next chunk
// in its buffer when a 'data' listener runs. destroy() left that chunk there and
// flow() emitted it after destroy() returned, with `destroyed === true`. Node's
// child stdio is a net.Socket that pushes asynchronously, so no 'data' follows
// destroy() there.
it.concurrent.each(["stdout", "stderr"] as const)(
"child.%s.destroy() inside a 'data' listener stops 'data' and 'end'",
async name => {
const SIZE = 8 * 1024 * 1024;
const writer = `const s=process.${name};s.on('error',()=>process.exit(0));const c=Buffer.alloc(1<<16,97);let w=0;(function f(){while(w<${SIZE}){w+=c.length;if(!s.write(c)){s.once('drain',f);return}}})()`;
const c = spawn(bunExe(), ["-e", writer], {
stdio: ["ignore", name === "stdout" ? "pipe" : "ignore", name === "stderr" ? "pipe" : "ignore"],
env: bunEnv,
});
try {
const stream = c[name]!;
let bytes = 0;
let destroyed = false;
const afterDestroy: string[] = [];
stream.on("data", (d: Buffer) => {
if (destroyed) {
afterDestroy.push(`data(${d.length}) destroyed=${stream.destroyed}`);
return;
}
bytes += d.length;
// Destroy as soon as another chunk is already buffered behind this
// one: that is the chunk that used to follow destroy(). If the reader
// never gets ahead of this listener, destroy half way instead.
if (stream.readableLength > 0 || bytes >= SIZE / 2) {
destroyed = true;
stream.destroy();
c.kill();
}
});
stream.on("end", () => afterDestroy.push("end"));
await once(c, "close");
expect(afterDestroy).toEqual([]);
expect(destroyed).toBe(true);
} finally {
c.kill();
}
},
);

// The exec()/execFile() 'data' listener is a port of Node's: at the chunk that
// crosses maxBuffer it destroys the stream, and it relies on no 'data' following
// destroy(). The chunk that used to follow made `maxBuffer - (totalLen - length)`
// negative, and slice(0, negative) appended most of it, so the callback got up
// to a chunk more than maxBuffer. Whether a chunk is buffered at the crossing is
// a race (about 1 run in 3 without the fix), so several run at once.
describe.concurrent("execFile() maxBuffer against a fast writer", () => {
const maxBuffer = 1024 * 1024;
// 3 MiB in 64 KiB blocks, each block filled with its own letter. ASCII on
// purpose: the handler, like Node's, counts bytes but slices a string chunk by
// code units, so multi-byte output is over maxBuffer in bytes in Node too
// (Node v26.3.0, maxBuffer 1000000, 2-byte characters: 1016960 bytes).
const writer = `const s=process.stdout;s.on('error',()=>process.exit(0));let k=0;(function f(){while(k<48){if(!s.write(Buffer.alloc(1<<16,65+(k++%26)))){s.once('drain',f);return}}})()`;
const expected = Buffer.concat(
Array.from({ length: maxBuffer >> 16 }, (_, k) => Buffer.alloc(1 << 16, 65 + (k % 26))),
);

it.each(["buffer", "utf8"] as const)("truncates at exactly maxBuffer (encoding: %s)", async encoding => {
const runs = Array.from({ length: 4 }, () => {
const { promise, resolve } = Promise.withResolvers();
execFile(bunExe(), ["-e", writer], { maxBuffer, encoding, env: bunEnv }, (err, stdout) => {
resolve({ code: err?.code, length: stdout.length, isPrefix: expected.equals(Buffer.from(stdout)) });
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});
return promise;
});
expect(await Promise.all(runs)).toEqual(
Array(4).fill({ code: "ERR_CHILD_PROCESS_STDIO_MAXBUFFER", length: maxBuffer, isPrefix: true }),
);
});
});

// When spawn fails (ENOENT, bad cwd, etc.) the ChildProcess emits 'error' and
// 'close' but never 'exit'. The abort listener on options.signal was only
// removed on 'exit', so every failed spawn against a shared AbortSignal leaked
Expand Down
58 changes: 58 additions & 0 deletions test/js/node/stream/node-stream.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -628,6 +628,64 @@ it("Readable.fromWeb: destroy(err) after consuming a chunk cancels the web sourc
});
});

// A native-backed Readable pushes its pull results synchronously, so while it flows Readable has the next chunk
// buffered when a 'data' listener runs, and it pushes EOF a tick after the last chunk. destroy() stopped neither: flow()
// emitted the buffered chunk with `destroyed === true`, and 'end' followed. Node's fromWeb pushes asynchronously, so
// nothing is buffered or due at that point, and it emits only 'close'.
describe.each([
["Blob.stream()", size => new Blob([Buffer.alloc(size, "x")]).stream()],
["Response.body", size => new Response(Buffer.alloc(size, "x")).body],
])("Readable.fromWeb(%s): destroy() inside a 'data' listener", (_, makeWeb) => {
// With today's chunking, 100 bytes are the only chunk and EOF is due a tick later, 16484 bytes are two chunks with
// the last one buffered behind the first, and 1 MiB has more buffered behind every chunk.
it.each([100, 16384 + 100, 1024 * 1024])("of a %d byte body stops 'data' and 'end'", async size => {
const r = Readable.fromWeb(makeWeb(size));
const events = [];
const { promise: closed, resolve } = Promise.withResolvers();
r.on("data", () => {
events.push(`data destroyed=${r.destroyed}`);
r.destroy();
});
r.on("end", () => events.push("end"));
r.on("close", () => {
events.push("close");
resolve();
});
await closed;
expect(events).toEqual(["data destroyed=false", "close"]);
});
});

it("Readable.fromWeb: destroy() on a paused stream keeps the buffered chunk for read(), as in Node", async () => {
const r = Readable.fromWeb(new Blob([Buffer.alloc(1024, "x")]).stream());
const { promise, resolve } = Promise.withResolvers();
r.once("readable", () => {
r.destroy();
resolve(r.read());
});
const chunk = await promise;
expect(chunk?.length).toBe(1024);
});

// Once the source has ended, what is buffered is all that is left. Node delivers it after destroy(), and to drop it
// would let 'end' follow data that never arrived.
it("Readable.fromWeb: a stream that ended while paused still delivers every byte after destroy(), as in Node", async () => {
const size = 16384 + 100;
const r = Readable.fromWeb(new Blob([Buffer.alloc(size, "x")]).stream());
let bytes = 0;
r.on("data", chunk => {
bytes += chunk.length;
r.destroy();
});
r.pause();
const closed = new Promise(resolve => r.once("close", resolve));
r.read(0);
while (!r._readableState.ended) await new Promise(resolve => setImmediate(resolve));
r.resume();
await closed;
expect(bytes).toBe(size);
});

it("Readable.toWeb(Readable.fromWeb(rs)).cancel(reason) propagates to the web source", async () => {
let cancelReason;
const web = new ReadableStream({
Expand Down
Loading