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
18 changes: 15 additions & 3 deletions src/runtime/webcore/FileSink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -401,6 +401,13 @@ impl FileSink {
}

let was_pending = (*this).pending.get().state == streams::PendingState::Pending;
let resumes = was_pending || (status == WriteStatus::Drained && !has_pending_data);
// Opened before `run_pending` so the settle and the resume share one checkpoint.
let _entered = if resumes && (*this).source_pending_pull.get() {
(*this).completion_scope()
} else {
None
};
Comment thread
robobun marked this conversation as resolved.
if was_pending {
// `consumed` was credited when the pending operation accepted its
// bytes; `amount` is only what this drain pushed to the fd.
Expand All @@ -419,9 +426,7 @@ impl FileSink {
FileSink::run_pending(this);
}

if (was_pending || (status == WriteStatus::Drained && !has_pending_data))
&& (*this).source_pending_pull.replace(false)
{
if resumes && (*this).source_pending_pull.replace(false) {
let mut src = *(*this).source.get();
src.ready(None, None);
}
Expand Down Expand Up @@ -494,6 +499,12 @@ impl FileSink {
}
}

/// Microtask checkpoint owed for the JS a writer callback enters while a stream is piped in.
fn completion_scope(&self) -> Option<bun_jsc::event_loop_handle::EnteredEventLoop> {
self.pipe.get().cell()?;
Some(self.event_loop().entered())
}

/// This sink opened its file itself, for the script that is running: if that is a
/// `Bun.ModuleGraph`'s, the file is closed with the graph. (The host's sinks are left to
/// flush at exit as they always have.)
Expand All @@ -515,6 +526,7 @@ impl FileSink {
unsafe {
// `source.close()` may drop the last ref (a Subprocess whose `.stdin` was never read).
let _guard = RefPtr::init_ref(this);
let _entered = (*this).completion_scope();

(*this).abort_handle.leave();
if (*this).js_global().is_some() {
Expand Down
267 changes: 266 additions & 1 deletion test/js/bun/util/filesink.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { createSocketPair, fileSinkInternals } from "bun:internal-for-testing";
import { describe, expect, it } from "bun:test";
import { bunEnv, bunExe, fileDescriptorLeakChecker, isLinux, isPosix, isWindows, tmpdirSync } from "harness";
import { bunEnv, bunExe, fileDescriptorLeakChecker, isLinux, isPosix, isWindows, tempDir, tmpdirSync } from "harness";
import { mkfifo } from "mkfifo";
import { join } from "node:path";

Expand Down Expand Up @@ -1067,6 +1067,271 @@ describe("FileSink on a pipe stays alive until end() has drained the buffer", ()
});
});

// `Bun.write(dest, new Response(stream))` pumps a JS stream into the sink, and each step of
// the pump arrives as a microtask: the next chunk, the stream's close, its error. A short
// write parks the pump. When the reader drains the pipe, the sink's writable callback
// resumes it or closes the sink, and that callback comes from the event loop, not from
// script: the microtasks it queues are its own to run. It did not run them, and with
// nothing else alive the process exited with them queued. The rest of the stream was
// never written and the promise never settled.
//
// The fixtures have no timer, and only one has a 'beforeExit' listener. Either gives the
// loop one more turn, which runs the queued step and hides the bug.
describe("a stream piped into a FileSink on a pipe is pumped to its end", () => {
// Larger than a pipe, so the first write is short and the pump parks on it.
const first = 2 * 1024 * 1024;
const chunk = 64 * 1024;
const chunks = 16;
// The whole stream as it must arrive. A stream that ends after the first chunk is a prefix.
const expected = Buffer.concat([Buffer.alloc(first, 97), Buffer.alloc(chunks * chunk, 98)]);

// "parked" goes out on the second pull(), which runs once the pump has taken the first
// chunk. The test reads only after that, so the first drain finds the pump parked.
const sources = {
chunks: `{
pull(c) {
if (++pulls === 1) return c.enqueue(first);
if (pulls === 2) console.error("parked");
if (pulls <= 1 + ${chunks}) return c.enqueue(chunk);
c.close();
},
}`,
close: `{
pull(c) {
if (++pulls === 1) return c.enqueue(first);
console.error("parked");
c.close();
},
}`,
error: `{
pull(c) {
if (++pulls === 1) return c.enqueue(first);
console.error("parked");
c.error(new Error("boom"));
},
}`,
// Nothing parks here: the sink is ended with the tail of the chunk still buffered, so
// its close comes later, from the callback that drains the tail.
directClose: `{
type: "direct",
pull(c) {
c.write(first);
console.error("parked");
c.close();
},
}`,
};

function fixture(source: string, dest: string, countBeforeExit = false) {
return `
const first = new Uint8Array(${first}).fill(97);
const chunk = new Uint8Array(${chunk}).fill(98);
let pulls = 0;
let settled = "pending";
let beforeExit = ${countBeforeExit ? 0 : undefined};
if (beforeExit === 0) process.on("beforeExit", () => { beforeExit++; });
process.on("exit", code => console.error(JSON.stringify({ settled, code, beforeExit })));
Bun.write(${dest}, new Response(new ReadableStream(${source}))).then(
n => { settled = "resolved " + n; },
e => { settled = "rejected " + e.message; },
);
`;
}

// Returns once the child has printed "parked"; `rest()` collects stderr to its end.
async function parked(stderr: ReadableStream<Uint8Array>) {
const decoder = new TextDecoder();
const reader = stderr.getReader();
let text = "";
while (!text.startsWith("parked\n")) {
const { value, done } = await reader.read();
if (done) break;
text += decoder.decode(value, { stream: true });
}
return async function rest() {
while (true) {
const { value, done } = await reader.read();
if (done) return text;
text += decoder.decode(value, { stream: true });
}
};
}

// read(2) in a poll on the non-blocking read end the test holds: it returns 0 once the
// writer is gone, on every POSIX, and needs no readiness notification.
async function drain(fd: number) {
const buffer = Buffer.alloc(64 * 1024);
const parts: Buffer[] = [];
let total = 0;
let progressAt = performance.now();
while (true) {
let n: number;
try {
n = fs.readSync(fd, buffer);
} catch (e: any) {
if (e.code !== "EAGAIN") throw e;
if (performance.now() - progressAt > 10_000) {
throw new Error(`the writer holds the FIFO open and wrote nothing for 10 s, after ${total} bytes`);
}
await Bun.sleep(1);
Comment thread
robobun marked this conversation as resolved.
continue;
}
if (n === 0) return Buffer.concat(parts);
parts.push(Buffer.from(buffer.subarray(0, n)));
total += n;
progressAt = performance.now();
}
}

async function run(
source: keyof typeof sources,
{ fifo, countBeforeExit }: { fifo?: { path: string; fd: number }; countBeforeExit?: boolean } = {},
) {
const dest = fifo ? `Bun.file(${JSON.stringify(fifo.path)})` : "Bun.stdout";
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", fixture(sources[source], dest, countBeforeExit)],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});
const rest = await parked(proc.stderr);
const [received, stderr, exitCode] = await Promise.all([
fifo ? drain(fifo.fd) : proc.stdout.bytes(),
rest(),
proc.exited,
]);
// `intact` is about the bytes that arrived: each one is the byte the stream has there.
const intact = expected.subarray(0, received.length).equals(received);
return { received: received.length, intact, stderr, exitCode };
}

const total = first + chunks * chunk;
Comment thread
robobun marked this conversation as resolved.

it.concurrent("chunks that arrive after the pump parked are written", async () => {
expect(await run("chunks")).toEqual({
received: total,
intact: true,
stderr: "parked\n" + JSON.stringify({ settled: `resolved ${total}`, code: 0 }) + "\n",
exitCode: 0,
});
});

// The test holds the read end from before the child opens the FIFO to write.
it.concurrent.skipIf(!isPosix)("the same into a FIFO the sink opened itself", async () => {
using dir = tempDir("filesink-piped-fifo", {});
const path = join(String(dir), "piped.fifo");
mkfifo(path, 0o666);
const fd = fs.openSync(path, fs.constants.O_RDONLY | fs.constants.O_NONBLOCK);
try {
expect(await run("chunks", { fifo: { path, fd } })).toEqual({
received: total,
intact: true,
stderr: "parked\n" + JSON.stringify({ settled: `resolved ${total}`, code: 0 }) + "\n",
exitCode: 0,
});
} finally {
fs.closeSync(fd);
}
});

// With a listener the stream arrives on unfixed code too: the loop goes back to work each
// time the listener has run. It comes to the edge of its exit after each drain, though,
// and the listener hears of each one.
it.concurrent("'beforeExit' is emitted once while the stream is pumped", async () => {
expect(await run("chunks", { countBeforeExit: true })).toEqual({
received: total,
intact: true,
stderr: "parked\n" + JSON.stringify({ settled: `resolved ${total}`, code: 0, beforeExit: 1 }) + "\n",
exitCode: 0,
});
});

it.concurrent("a stream that closes while the pump is parked settles the promise", async () => {
expect(await run("close")).toEqual({
received: first,
intact: true,
stderr: "parked\n" + JSON.stringify({ settled: `resolved ${first}`, code: 0 }) + "\n",
exitCode: 0,
});
});

it.concurrent("a stream that fails while the pump is parked rejects the promise", async () => {
expect(await run("error")).toEqual({
received: first,
intact: true,
stderr: "parked\n" + JSON.stringify({ settled: "rejected boom", code: 0 }) + "\n",
exitCode: 0,
});
});

it.concurrent("a sink that is closed once its buffer drains settles the promise", async () => {
expect(await run("directClose")).toEqual({
received: first,
intact: true,
stderr: "parked\n" + JSON.stringify({ settled: `resolved ${first}`, code: 0 }) + "\n",
exitCode: 0,
});
});

// No write is short here and nothing parks: the chunk is small and the sink buffers it.
// The child's stdout is a socket whose send buffer is already full, so the end() that
// follows the stream's close cannot flush. The sink closes from the callback that
// drains it, which resumes no source and has only the close to run.
it.concurrent.skipIf(!isPosix)("a sink that is closed with no source parked settles the promise", async () => {
const tail = "the only chunk\n";
const [readFd, writeFd] = createSocketPair();
let writeFdOpen = true;
try {
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;
}

await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
fixture(
`{
start(c) {
c.enqueue(new TextEncoder().encode(${JSON.stringify(tail)}));
c.close();
},
}`,
"Bun.stdout",
) + `console.error("parked");`,
],
env: bunEnv,
stdout: writeFd,
stderr: "pipe",
});
// The child holds the only write end now: its exit is the end of the data.
fs.closeSync(writeFd);
writeFdOpen = false;

const rest = await parked(proc.stderr);
const [received, stderr, exitCode] = await Promise.all([Bun.file(readFd).bytes(), rest(), proc.exited]);
expect({
filler: received.length - tail.length,
tail: new TextDecoder().decode(received.subarray(filled)),
stderr,
exitCode,
}).toEqual({
filler: filled,
tail,
stderr: "parked\n" + JSON.stringify({ settled: `resolved ${tail.length}`, code: 0 }) + "\n",
exitCode: 0,
});
} finally {
if (writeFdOpen) fs.closeSync(writeFd);
fs.closeSync(readFd);
}
});
});

// FileSink::on_close tells the owner of the sink that it closed, and a Subprocess then drops its ref
// on its stdin sink. That is the only ref when script never read `proc.stdin`, and on_close used the
// sink after it (ASAN: heap-use-after-free in settle_stream_done). Outside tests only the stop phase
Expand Down
Loading