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
9 changes: 9 additions & 0 deletions src/runtime/webcore/FileSink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -891,9 +891,18 @@ impl FileSink {
return sys::Result::Ok(JSValue::UNDEFINED);
}

let had_buffered_data = self.writer.get().has_pending_data();
// SAFETY(JsCell): `IOWriter::flush` is pure I/O; no JS re-entry while
// the `&mut IOWriter` is held.
let rc = self.writer.with_mut(|w| w.flush());
// `on_write` keeps the event loop alive while bytes sit in the buffer
// and `on_auto_flush` releases that once it drains them. A flush from
// JS that drained them has to release it too, or the loop still counts
// as alive until the next deferred-task drain (a write()+flush() from a
// 'beforeExit' listener then re-emits 'beforeExit').
if had_buffered_data && !self.writer.get().has_pending_data() {
self.update_ref(false);
}
let flushed = match rc {
WriteResult::Done(written)
| WriteResult::Pending(written)
Expand Down
35 changes: 35 additions & 0 deletions test/js/bun/console/console-write.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,3 +36,38 @@ console.write("ok");
signalCode: null,
});
});

// console.write() is a FileSink write() + flush(). With stdout on a pipe the
// write() parks the bytes in the sink's buffer and marks the event loop alive
// until they are flushed; the flush() drained them right away but used to leave
// that mark in place, so a console.write() from a 'beforeExit' listener looked
// like newly scheduled work and 'beforeExit' was emitted again. A listener that
// writes on every emit kept the process alive forever; this one writes once, so
// the broken behavior shows up as a count of 2.
test("console.write inside a 'beforeExit' listener does not re-emit 'beforeExit'", async () => {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`
let count = 0;
process.on("beforeExit", () => {
count++;
if (count === 1) console.write("from beforeExit\\n");
});
process.on("exit", () => console.log("beforeExit emitted " + count + " time(s)"));
`,
],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});

const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);

expect({ stdout, stderr, exitCode }).toEqual({
stdout: "from beforeExit\nbeforeExit emitted 1 time(s)\n",
stderr: "",
exitCode: 0,
});
});
136 changes: 136 additions & 0 deletions test/js/bun/util/filesink.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -676,6 +676,142 @@ describe.skipIf(isWindows)("FileSink buffered data is flushed on process exit",
);
});

// On a pipe or socket, a small write() parks the bytes in the sink's buffer and
// marks the event loop alive until they are flushed; the deferred auto-flush
// clears that once it has drained them. An explicit flush() that drained them
// itself (console.write is write()+flush()) used to leave the mark in place
// until that deferred task ran. 'beforeExit' is where this shows: it is
// re-emitted whenever a listener left the loop alive, so a listener doing
// write()+flush() got it a second time, and one writing on every emit kept the
// process alive forever. The fixtures below write on the first emit only and
// report on stderr, from 'exit', how many emits they saw.
describe("FileSink flush() from a 'beforeExit' listener", () => {
const marker = "from beforeExit\n";

function fixture(stream: "stdout" | "stderr") {
return `
let count = 0;
let code = null;
process.on("beforeExit", () => {
count++;
if (count !== 1) return;
const sink = Bun.${stream}.writer();
sink.write(${JSON.stringify(marker)});
try {
sink.flush();
} catch (e) {
code = e.code;
}
});
process.on("exit", () => console.error(JSON.stringify({ count, code })));
`;
}

// `stdio` replaces the child's stdin/stdout with raw fds; stdout is then not
// captured and comes back as null.
async function run(script: string, stdio: { stdin?: number; stdout?: number } = {}) {
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", script],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
...stdio,
});
const [stdout, stderr, exitCode] = await Promise.all([
stdio.stdout === undefined ? (proc.stdout as ReadableStream).text() : null,
(proc.stderr as ReadableStream).text(),
proc.exited,
]);
return { stdout, stderr, exitCode };
}

const once = JSON.stringify({ count: 1, code: null }) + "\n";

it.concurrent("Bun.stdout.writer() write()+flush() on a pipe emits 'beforeExit' once", async () => {
expect(await run(fixture("stdout"))).toEqual({ stdout: marker, stderr: once, exitCode: 0 });
});

it.concurrent("Bun.stderr.writer() write()+flush() on a pipe emits 'beforeExit' once", async () => {
expect(await run(fixture("stderr"))).toEqual({ stdout: "", stderr: marker + once, exitCode: 0 });
});

// flush() can also fail outright: the writer drops the buffered bytes and
// flush() throws. Nothing is pending after that either, so this path must not
// leave the loop marked alive any more than the success path does. The
// child's stdout is a socket whose peer is already closed.
it.concurrent.skipIf(!isPosix)("a flush() that fails with EPIPE emits 'beforeExit' once", async () => {
const [readFd, writeFd] = createSocketPair();
fs.closeSync(readFd);
try {
expect(await run(fixture("stdout"), { stdout: writeFd })).toEqual({
stdout: null,
stderr: JSON.stringify({ count: 1, code: "EPIPE" }) + "\n",
exitCode: 0,
});
} finally {
fs.closeSync(writeFd);
}
});

// The other direction has to keep working: when flush() cannot drain the
// buffer, the bytes are still pending and the process has to stay alive until
// they go out. The child's stdout is a socket whose send buffer is already
// full; its stdin is the other end. The unref'd timer that drains it does not
// hold the process open by itself, it only gets to run because the pending
// bytes do, so releasing the loop on this path would exit the child before
// the timer fires and before flush()'s promise settles.
it.concurrent.skipIf(!isPosix)("a flush() that could not drain keeps the process alive until it does", async () => {
const [readFd, writeFd] = createSocketPair();
try {
// createSocketPair() hands out non-blocking fds: write until the kernel
// refuses more.
const filler = Buffer.alloc(64 * 1024, 0x61);
let filled = 0;
try {
while (true) filled += fs.writeSync(writeFd, filler);
} catch (e: any) {
if (e.code !== "EAGAIN") throw e;
}

const result = await run(
`
const fs = require("node:fs");
const sink = Bun.stdout.writer();
const wrote = sink.write(${JSON.stringify(marker)});
const flushed = sink.flush();
let settled = "pending";
Promise.resolve(flushed).then(
() => { settled = "resolved"; },
e => { settled = "rejected: " + e.code; },
);

setTimeout(() => {
const buf = Buffer.alloc(64 * 1024);
let drained = 0;
while (drained < ${filled}) drained += fs.readSync(0, buf);
}, 0).unref();

let count = 0;
process.on("beforeExit", () => { count++; });
process.on("exit", () =>
console.error(JSON.stringify({ wrote, flushReturnedPromise: flushed instanceof Promise, settled, count })),
);
`,
{ stdin: readFd, stdout: writeFd },
);
expect(result).toEqual({
stdout: null,
stderr:
JSON.stringify({ wrote: marker.length, flushReturnedPromise: true, settled: "resolved", count: 1 }) + "\n",
exitCode: 0,
});
} finally {
fs.closeSync(readFd);
fs.closeSync(writeFd);
}
});
});

it("fs.promises.writeFile with iterables under GC pressure does not crash", async () => {
const dir = tmpdirSync();
await using proc = Bun.spawn({
Expand Down