diff --git a/src/js/internal/streams/native-readable.ts b/src/js/internal/streams/native-readable.ts index 2b8a8c4fabe4..3d62e724416f 100644 --- a/src/js/internal/streams/native-readable.ts +++ b/src/js/internal/streams/native-readable.ts @@ -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; @@ -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) { @@ -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. @@ -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; @@ -241,9 +250,7 @@ function handleArrayBufferViewResult(stream: NativeReadable, result: any, chunk: } if (isClosed) { - process.nextTick(() => { - stream.push(null); - }); + process.nextTick(pushEof, stream); } return chunk; @@ -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; diff --git a/test/js/bun/spawn/spawn-stdio-syscall-error.test.ts b/test/js/bun/spawn/spawn-stdio-syscall-error.test.ts index 70c6fdd2c53b..a6f323796272 100644 --- a/test/js/bun/spawn/spawn-stdio-syscall-error.test.ts +++ b/test/js/bun/spawn/spawn-stdio-syscall-error.test.ts @@ -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 #include #include +#include #include #include #include @@ -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; @@ -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); } @@ -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; } } @@ -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"; @@ -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"); @@ -266,6 +299,7 @@ async function runWithFault(fixture: string, fault: Record) { 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({ @@ -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"], diff --git a/test/js/node/child_process/child_process.test.ts b/test/js/node/child_process/child_process.test.ts index 4285422d73d7..212a12c9ba2f 100644 --- a/test/js/node/child_process/child_process.test.ts +++ b/test/js/node/child_process/child_process.test.ts @@ -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)) }); + }); + 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 diff --git a/test/js/node/stream/node-stream.test.js b/test/js/node/stream/node-stream.test.js index 02e1d6c06ea4..782ed9039c62 100644 --- a/test/js/node/stream/node-stream.test.js +++ b/test/js/node/stream/node-stream.test.js @@ -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({