diff --git a/src/js/internal/fs/streams.ts b/src/js/internal/fs/streams.ts index 0d1ed409d4f2..535e35e8c5f0 100644 --- a/src/js/internal/fs/streams.ts +++ b/src/js/internal/fs/streams.ts @@ -593,7 +593,6 @@ function underscoreWriteFast(this: FSStream, data: any, encoding: any, cb: any) this._write = _write; return this._write(data, encoding, cb); } - const hasCallback = typeof cb === "function"; try { if (fileSink === true) { fileSink = this[kWriteStreamFastPath] = Bun.file(this.path).writer(); @@ -610,12 +609,7 @@ function underscoreWriteFast(this: FSStream, data: any, encoding: any, cb: any) }, err => { if (cb) cb(err); - // If no callback was provided, emit the error on the stream - // This matches Node.js behavior where unhandled write errors - // are emitted as 'error' events on the stream - if (!hasCallback) { - this.destroy(err); - } + require("internal/streams/destroy").errorOrDestroy(this, err); }, ); return false; @@ -625,10 +619,7 @@ function underscoreWriteFast(this: FSStream, data: any, encoding: any, cb: any) } } catch (e) { if (cb) process.nextTick(cb, e); - // If no callback was provided, emit the error on the stream - if (!hasCallback) { - this.destroy(e); - } + require("internal/streams/destroy").errorOrDestroy(this, e, true); return false; } } @@ -639,10 +630,10 @@ const kWriteMonkeyPatchDefense = Symbol("!"); function writeFast(this: FSStream, data: any, encoding: any, cb: any) { if (this[kWriteMonkeyPatchDefense]) return writablePrototypeWrite.$call(this, data, encoding, cb); - // After end() the Writable contract requires write() to fail with - // ERR_STREAM_WRITE_AFTER_END and not reach the sink. + // After end()/destroy() the Writable contract requires write() to fail with + // ERR_STREAM_WRITE_AFTER_END / ERR_STREAM_DESTROYED and not reach the sink. const state = this._writableState; - if (state !== undefined && state.ending) { + if (state !== undefined && (state.ending || state.destroyed)) { return writablePrototypeWrite.$call(this, data, encoding, cb); } @@ -650,8 +641,7 @@ function writeFast(this: FSStream, data: any, encoding: any, cb: any) { cb = encoding; encoding = undefined; } - const hasCallback = typeof cb === "function"; - if (!hasCallback) { + if (typeof cb !== "function") { cb = streamNoop; } @@ -659,21 +649,20 @@ function writeFast(this: FSStream, data: any, encoding: any, cb: any) { if (fileSink && fileSink !== true) { const maybePromise = fileSink.write(data); if ($isPromise(maybePromise)) { - maybePromise - .then(() => { + // Two-arg then(): a throw from the fulfillment handler must not be + // mistaken for a write failure. + maybePromise.then( + () => { this.emit("drain"); // Emit drain event cb(null); - }) - .catch(err => { - // Always call the callback with the error + }, + err => { cb(err); - // If no callback was provided, emit the error on the stream - // This matches Node.js behavior where unhandled write errors - // are emitted as 'error' events on the stream - if (!hasCallback) { - this.destroy(err); - } - }); + // Node.js onwriteError: callback AND destroy are both invoked; the + // callback is additive, not a replacement for the 'error' event. + require("internal/streams/destroy").errorOrDestroy(this, err); + }, + ); return false; // Indicate backpressure } else { cb(null); diff --git a/test/js/node/child_process/child_process.test.ts b/test/js/node/child_process/child_process.test.ts index 291ec5ae656a..5d0d91053455 100644 --- a/test/js/node/child_process/child_process.test.ts +++ b/test/js/node/child_process/child_process.test.ts @@ -405,6 +405,54 @@ describe("spawn()", () => { expect(child.stderr).not.toBeNull(); }); }); + + it.skipIf(isWindows)( + "stdin write failure (EPIPE) emits 'error' and destroys even with a write callback", + async () => { + // Child closes its own stdin fd, signals ready on stdout, then stays alive. + const child = spawn( + bunExe(), + ["-e", `require("fs").closeSync(0); process.stdout.write("ready\\n"); setInterval(() => {}, 1e5);`], + { env: bunEnv, stdio: ["pipe", "pipe", "ignore"] }, + ); + try { + await new Promise((resolve, reject) => { + child.on("error", reject); + child.on("exit", () => reject(new Error("child exited before ready"))); + child.stdout!.once("data", () => resolve()); + }); + child.removeAllListeners("error"); + child.removeAllListeners("exit"); + + const errEv = Promise.withResolvers(); + const cb1 = Promise.withResolvers(); + child.stdin!.on("error", e => errEv.resolve(e)); + child.stdin!.write(Buffer.alloc(65536, 0x41), e => cb1.resolve(e)); + + const [cb1Err, errEvErr] = await Promise.all([cb1.promise, errEv.promise]); + + expect({ + cb1: cb1Err?.code, + errEv: errEvErr?.code, + destroyed: child.stdin!.destroyed, + writable: child.stdin!.writable, + }).toEqual({ + cb1: "EPIPE", + errEv: "EPIPE", + destroyed: true, + writable: false, + }); + + const cb2 = Promise.withResolvers(); + const r2 = child.stdin!.write("more-bytes", e => cb2.resolve(e)); + const cb2Err = await cb2.promise; + + expect({ r2, cb2: cb2Err?.code }).toEqual({ r2: false, cb2: "ERR_STREAM_DESTROYED" }); + } finally { + child.kill("SIGKILL"); + } + }, + ); }); describe("execFile()", () => {