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
45 changes: 17 additions & 28 deletions src/js/internal/fs/streams.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand All @@ -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;
Expand All @@ -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;
}
}
Expand All @@ -639,41 +630,39 @@ 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);
}

if (typeof encoding === "function") {
cb = encoding;
encoding = undefined;
}
const hasCallback = typeof cb === "function";
if (!hasCallback) {
if (typeof cb !== "function") {
cb = streamNoop;
}

const fileSink = this[kWriteStreamFastPath];
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);
Expand Down
48 changes: 48 additions & 0 deletions test/js/node/child_process/child_process.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>((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<any>();
const cb1 = Promise.withResolvers<any>();
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<any>();
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()", () => {
Expand Down
Loading