From 4fa0b57f405816c2379687ffcf4eeb7456cf9950 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 12 Aug 2026 12:25:33 +0000 Subject: [PATCH 1/7] Bun.file(fifo).bytes(): wait for a named pipe with select(2) on macOS XNU attaches a FIFO's EVFILT_READ filter to the vnode, where it only activates while bytes are buffered; closing the last writer posts nothing, and poll(2) is implemented on top of kqueue. So once ReadFile had drained a named pipe and handed itself to the io thread it never learned about EOF, and Bun.file(fifo).bytes()/text() (or Bun.stdin with a FIFO as stdin) never resolved. select(2) reaches the FIFO's socket, whose readability includes the state set by the last writer's close, so on macOS a ReadFile on a FIFO vnode (S_IFIFO with a non-zero st_dev; pipe(2) pipes report 0 and keep using kqueue) now waits in select on the pool thread instead, including before its first read so a writer that has not connected yet is waited for as on Linux. --- src/runtime/webcore/blob/read_file.rs | 49 +++++++++ src/sys/lib.rs | 62 ++++++++++- test/js/bun/util/bun-file-read.test.ts | 136 ++++++++++++++++++++++++- 3 files changed, 245 insertions(+), 2 deletions(-) diff --git a/src/runtime/webcore/blob/read_file.rs b/src/runtime/webcore/blob/read_file.rs index 8bc36a6b5805..02d76a17bcd1 100644 --- a/src/runtime/webcore/blob/read_file.rs +++ b/src/runtime/webcore/blob/read_file.rs @@ -285,6 +285,11 @@ pub struct ReadFile { pub(crate) io_request: io::Request, #[cfg(not(windows))] pub(crate) could_block: bool, + /// A FIFO vnode as opposed to a `pipe(2)` pipe. kqueue never reports EOF + /// for these, so their waits happen in `block_until_readable` instead of + /// going through the io thread; see `bun_sys::block_until_readable`. + #[cfg(target_os = "macos")] + pub(crate) is_named_pipe: bool, pub(crate) close_after_io: bool, pub(crate) state: AtomicU8, // ClosingState } @@ -380,6 +385,8 @@ impl ReadFile { scheduled: false, }, could_block: false, + #[cfg(target_os = "macos")] + is_named_pipe: false, close_after_io: false, state: AtomicU8::new(ClosingState::Running as u8), }; @@ -465,6 +472,22 @@ impl ReadFile { } } + /// The named-pipe counterpart of `wait_for_readable`: waits right here on + /// the pool thread, like `fs.readFile` does for every FIFO. Returns `false` + /// when the wait itself failed, with the error recorded for `then()`. + #[cfg(target_os = "macos")] + fn block_until_readable(&mut self) -> bool { + bloblog!("ReadFile.blockUntilReadable"); + match bun_sys::block_until_readable(self.opened_fd) { + Ok(()) => true, + Err(err) => { + self.errno = Some(bun_errno::from_errno(err.errno as i32).into()); + self.system_error = Some(err.to_system_error().into()); + false + } + } + } + /// Pick the read target: `buffer`'s spare capacity if it is at least as /// large as `stack_buffer`, otherwise `stack_buffer`; capped by /// `max_length - read_off`. Returns `(use_stack, target)` so the caller @@ -682,6 +705,13 @@ impl ReadFile { } self.could_block = !bun_sys::is_regular_file(stat.st_mode as _); + #[cfg(target_os = "macos")] + { + // `pipe(2)` pipes are S_IFIFO too; XNU's pipe_stat leaves their + // st_dev 0, whereas a FIFO on a filesystem carries its volume's + // device number. + self.is_named_pipe = bun_sys::S::ISFIFO(stat.st_mode as _) && stat.st_dev != 0; + } self.total_size = SizeType::try_from((stat.st_size as i64).max(0).min(MAX_SIZE as i64)).unwrap(); @@ -754,6 +784,18 @@ impl ReadFile { // If we immediately call read(), it will block until stdin is // readable. if self.could_block { + #[cfg(target_os = "macos")] + if self.is_named_pipe { + // Waiting before the first read is also what keeps a FIFO + // whose writer has not connected yet from reading as empty. + if self.block_until_readable() { + self.do_read_loop(); + } else { + self.on_finish(); + } + return; + } + if bun_core::is_readable(fd) == bun_core::Pollable::NotReady { self.wait_for_readable(); return; @@ -865,6 +907,13 @@ impl ReadFile { bun_core::Pollable::Ready | bun_core::Pollable::Hup => continue, } } + #[cfg(target_os = "macos")] + if self.is_named_pipe { + if self.block_until_readable() { + continue; + } + break; + } self.read_eof = false; self.buffer = buffer; self.wait_for_readable(); diff --git a/src/sys/lib.rs b/src/sys/lib.rs index 8b704ae5b690..977d52985626 100644 --- a/src/sys/lib.rs +++ b/src/sys/lib.rs @@ -1411,6 +1411,8 @@ impl Tag { #[cfg(not(windows))] pub(crate) const setrlimit: Tag = Tag(106); pub const clone3: Tag = Tag(107); + #[cfg(target_os = "macos")] + pub(crate) const select: Tag = Tag(108); // `inotify_init1`/`inotify_add_watch` fold under the generic `.watch` // tag; `INotifyWatcher.rs` spells it `.inotify`. Alias to `.watch` // so the JS-facing `err.syscall == "watch"` string stays node-compatible. @@ -1418,7 +1420,7 @@ impl Tag { /// The tag name — spelling is frozen (JS-facing /// `err.syscall` string; node-compat code matches on it). pub fn name(self) -> &'static str { - const NAMES: [&str; 108] = [ + const NAMES: [&str; 109] = [ "TODO", "dup", "access", @@ -1528,6 +1530,7 @@ impl Tag { "getrlimit", "setrlimit", "clone3", + "select", ]; NAMES.get(self.0 as usize).copied().unwrap_or("unknown") } @@ -1708,6 +1711,19 @@ mod nocancel { ) -> isize; #[link_name = "poll$NOCANCEL"] pub(crate) fn poll(fds: *mut libc::pollfd, nfds: libc::nfds_t, timeout: c_int) -> c_int; + // The `_DARWIN_UNLIMITED_SELECT` variant of select(2) (same + // `$DARWIN_EXTSN` scheme as `realpath` in `posix_impl`). libsyscall + // maps it straight onto the syscall, so there is no FD_SETSIZE check + // and each set is a bitmap of ceil(nfds / 32) 32-bit words rather + // than a `libc::fd_set`; hence the word pointers. + #[link_name = "select$DARWIN_EXTSN$NOCANCEL"] + pub(crate) fn select( + nfds: c_int, + readfds: *mut u32, + writefds: *mut u32, + errorfds: *mut u32, + timeout: *mut libc::timeval, + ) -> c_int; // Remaining `$NOCANCEL` variants Bun links against. // safe: by-value `c_int` fd; bad fd → -1/EBADF, no UB. #[link_name = "close$NOCANCEL"] @@ -7607,6 +7623,50 @@ pub fn kevent( } } +/// Blocks the calling thread in `select(2)` until `fd` is readable, where +/// readable includes EOF. Retries on EINTR. +/// +/// This is how a named pipe (a FIFO opened by path, or one inherited as +/// stdin) has to be waited on under macOS. XNU attaches kqueue `EVFILT_READ` +/// filters for a FIFO to its vnode, and that filter only fires while bytes are +/// buffered (`vnode_readable_data_count`); the last writer closing posts +/// nothing to the vnode, so a kqueue registration never wakes a reader up for +/// EOF, and neither does `poll(2)`, which XNU implements on top of kqueue. +/// `select(2)` instead goes through `fifo_select` to the FIFO's underlying +/// socket, whose readability includes the `SS_CANTRCVMORE` state that the last +/// writer's close sets (and that a writer which has not connected yet leaves +/// clear, so a FIFO that is still waiting for its first writer blocks here +/// rather than reading as empty). `pipe(2)` pipes are not affected: their own +/// kqueue filter reports `EV_EOF`. +#[cfg(target_os = "macos")] +pub fn block_until_readable(fd: Fd) -> Maybe<()> { + debug_assert!(fd.is_valid()); + let index = fd.native() as usize; + let word = index / 32; + let mut read_set = vec![0u32; word + 1]; + loop { + read_set.fill(0); + read_set[word] = 1 << (index % 32); + // SAFETY: `read_set` holds the ceil(nfds / 32) words the kernel reads + // and writes back for `nfds = fd + 1`; the write and error sets and the + // timeout (wait indefinitely) may be null. + let rc = unsafe { + nocancel::select( + fd.native() + 1, + read_set.as_mut_ptr(), + core::ptr::null_mut(), + core::ptr::null_mut(), + core::ptr::null_mut(), + ) + }; + match get_errno(rc) { + E::SUCCESS => return Ok(()), + E::EINTR => continue, + e => return Err(Error::from_code(e, Tag::select).with_fd(fd)), + } + } +} + /// `clonefile` — macOS-only CoW copy. On non-Darwin returns ENOTSUP so /// callers can fall back to `copy_file`. #[cfg(not(target_os = "macos"))] diff --git a/test/js/bun/util/bun-file-read.test.ts b/test/js/bun/util/bun-file-read.test.ts index 40c6309a0c13..b0127432fbfd 100644 --- a/test/js/bun/util/bun-file-read.test.ts +++ b/test/js/bun/util/bun-file-read.test.ts @@ -1,5 +1,8 @@ import { describe, expect, it } from "bun:test"; -import { tempDir } from "harness"; +import { bunEnv, bunExe, isWindows, tempDir } from "harness"; +import { mkfifo } from "mkfifo"; +import { randomBytes } from "node:crypto"; +import { closeSync, constants, openSync, writeSync } from "node:fs"; import { tmpdir } from "node:os"; import path from "node:path"; @@ -54,3 +57,134 @@ describe("Bun.file read-loop target selection", () => { expect(Bun.hash(buf)).toBe(Bun.hash(bytes.subarray(start, end))); }); }); + +// Whole-file reads of a named pipe. Every one of these ends with the reader +// having drained the pipe and then learning that the last writer closed; on +// macOS that EOF is invisible to kqueue and poll(2) (they only see buffered +// bytes on a FIFO), so the reader has to wait for it differently than it does +// for a pipe(2) pipe, and each of these used to leave the child blocked +// forever there. +describe.skipIf(isWindows)("reading a named pipe to EOF", () => { + function readFifoInChild(script: string, fifo: string, stdin: number | "ignore" = "ignore") { + return Bun.spawn({ + cmd: [bunExe(), "-e", script], + env: { ...bunEnv, FIFO: fifo }, + stdin, + stdout: "pipe", + stderr: "pipe", + }); + } + + // The child must already have the FIFO open for reading before a writer can + // connect to it: a non-blocking open for writing fails with ENXIO until then. + async function openWriterOnceChildIsReading(fifo: string, child: Bun.Subprocess): Promise { + while (true) { + try { + return openSync(fifo, constants.O_WRONLY | constants.O_NONBLOCK); + } catch (err: any) { + if (err.code !== "ENXIO") throw err; + } + if (child.exitCode !== null || child.signalCode !== null) { + throw new Error(`child exited (${child.exitCode ?? child.signalCode}) without opening the FIFO`); + } + await Bun.sleep(5); + } + } + + it.concurrent("bytes() collects a payload that arrives in pieces and ends when the writer closes", async () => { + const payload = randomBytes(256 * 1024); + using dir = tempDir("bun-file-read-fifo", {}); + const fifo = path.join(String(dir), "in.fifo"); + mkfifo(fifo); + // The write end can only be opened, and written to without EPIPE, while + // some reader has the FIFO open; `holder` is that reader until the child + // has opened its own. It never reads, so every byte goes to the child. + let holder = openSync(fifo, constants.O_RDONLY | constants.O_NONBLOCK); + const closeHolder = () => { + if (holder !== -1) closeSync(holder); + holder = -1; + }; + let writer = openSync(fifo, "w"); + try { + await using proc = readFifoInChild( + `const bytes = await Bun.file(process.env.FIFO).bytes(); process.stdout.write(bytes.length + " " + Bun.hash(bytes));`, + fifo, + ); + const stderr = proc.stderr.text(); + // The write end is blocking, so this write only completes as the child + // drains the pipe, and closing it afterwards is what ends the child's + // read. The child cannot exit before that unless it failed; dropping + // `holder` then leaves the pipe without readers, so the blocked write + // fails with EPIPE instead of waiting forever. + const childDied = proc.exited.then(async exitCode => { + closeHolder(); + throw new Error(`child exited with ${exitCode} before the payload was written: ${await stderr}`); + }); + const written = await Promise.race([Bun.write(Bun.file(writer), payload), childDied]); + closeSync(writer); + writer = -1; + const [stdout, stderrText, exitCode] = await Promise.all([proc.stdout.text(), stderr, proc.exited]); + + expect({ written, stdout, stderr: stderrText }).toEqual({ + written: payload.length, + stdout: `${payload.length} ${Bun.hash(payload)}`, + stderr: "", + }); + expect(exitCode).toBe(0); + } finally { + if (writer !== -1) closeSync(writer); + closeHolder(); + } + }); + + it.concurrent("text() waits for a writer that connects after the read started", async () => { + using dir = tempDir("bun-file-read-fifo-late-writer", {}); + const fifo = path.join(String(dir), "late.fifo"); + mkfifo(fifo); + + await using proc = readFifoInChild(`process.stdout.write(await Bun.file(process.env.FIFO).text());`, fifo); + const writer = await openWriterOnceChildIsReading(fifo, proc); + writeSync(writer, "written after the reader opened\n"); + closeSync(writer); + + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout, stderr }).toEqual({ stdout: "written after the reader opened\n", stderr: "" }); + expect(exitCode).toBe(0); + }); + + it.concurrent("text() resolves empty when the writer connects and closes without writing", async () => { + using dir = tempDir("bun-file-read-fifo-empty", {}); + const fifo = path.join(String(dir), "empty.fifo"); + mkfifo(fifo); + + await using proc = readFifoInChild( + `const text = await Bun.file(process.env.FIFO).text(); process.stdout.write(JSON.stringify(text));`, + fifo, + ); + closeSync(await openWriterOnceChildIsReading(fifo, proc)); + + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout, stderr }).toEqual({ stdout: '""', stderr: "" }); + expect(exitCode).toBe(0); + }); + + it.concurrent("Bun.stdin.text() reads a FIFO inherited as stdin to EOF", async () => { + using dir = tempDir("bun-file-read-fifo-stdin", {}); + const fifo = path.join(String(dir), "stdin.fifo"); + mkfifo(fifo); + + // Same dance as above: a reader has to exist before the write end can be + // opened; here that reader becomes the child's stdin. + const readEnd = openSync(fifo, constants.O_RDONLY | constants.O_NONBLOCK); + const writer = openSync(fifo, "w"); + await using proc = readFifoInChild(`process.stdout.write(JSON.stringify(await Bun.stdin.text()));`, fifo, readEnd); + // The child has its own descriptor for the read end now. + closeSync(readEnd); + writeSync(writer, "stdin is a named pipe\n"); + closeSync(writer); + + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout, stderr }).toEqual({ stdout: JSON.stringify("stdin is a named pipe\n"), stderr: "" }); + expect(exitCode).toBe(0); + }); +}); From f4c6d7951c78aecb77526917f544482767269d2a Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 12 Aug 2026 12:41:12 +0000 Subject: [PATCH 2/7] read_file: skip the poll() pre-check for named pipes, select() answers it --- src/runtime/webcore/blob/read_file.rs | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/src/runtime/webcore/blob/read_file.rs b/src/runtime/webcore/blob/read_file.rs index 02d76a17bcd1..51d9719a9e94 100644 --- a/src/runtime/webcore/blob/read_file.rs +++ b/src/runtime/webcore/blob/read_file.rs @@ -897,6 +897,13 @@ impl ReadFile { // call. We already know it's done. && !self.read_eof) { + #[cfg(target_os = "macos")] + if self.is_named_pipe { + if self.block_until_readable() { + continue; + } + break; + } if self.could_block // If we received EOF, we can skip the poll() system // call. We already know it's done. @@ -907,13 +914,6 @@ impl ReadFile { bun_core::Pollable::Ready | bun_core::Pollable::Hup => continue, } } - #[cfg(target_os = "macos")] - if self.is_named_pipe { - if self.block_until_readable() { - continue; - } - break; - } self.read_eof = false; self.buffer = buffer; self.wait_for_readable(); From 64e6b546975fbfe3e52ed8a6d3f01164fa92ac8f Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 12 Aug 2026 12:57:13 +0000 Subject: [PATCH 3/7] Tighten the named-pipe comments; close test fds through using --- src/runtime/webcore/blob/read_file.rs | 16 ++-- src/sys/lib.rs | 29 +++---- test/js/bun/util/bun-file-read.test.ts | 109 ++++++++++++++----------- 3 files changed, 76 insertions(+), 78 deletions(-) diff --git a/src/runtime/webcore/blob/read_file.rs b/src/runtime/webcore/blob/read_file.rs index 51d9719a9e94..31f4a6149296 100644 --- a/src/runtime/webcore/blob/read_file.rs +++ b/src/runtime/webcore/blob/read_file.rs @@ -285,9 +285,7 @@ pub struct ReadFile { pub(crate) io_request: io::Request, #[cfg(not(windows))] pub(crate) could_block: bool, - /// A FIFO vnode as opposed to a `pipe(2)` pipe. kqueue never reports EOF - /// for these, so their waits happen in `block_until_readable` instead of - /// going through the io thread; see `bun_sys::block_until_readable`. + /// FIFO vnode (not a `pipe(2)` pipe); see `bun_sys::block_until_readable`. #[cfg(target_os = "macos")] pub(crate) is_named_pipe: bool, pub(crate) close_after_io: bool, @@ -472,9 +470,8 @@ impl ReadFile { } } - /// The named-pipe counterpart of `wait_for_readable`: waits right here on - /// the pool thread, like `fs.readFile` does for every FIFO. Returns `false` - /// when the wait itself failed, with the error recorded for `then()`. + /// `wait_for_readable` for named pipes: waits on this (pool) thread. + /// Returns `false` if the wait failed, with the error recorded for `then()`. #[cfg(target_os = "macos")] fn block_until_readable(&mut self) -> bool { bloblog!("ReadFile.blockUntilReadable"); @@ -707,9 +704,7 @@ impl ReadFile { self.could_block = !bun_sys::is_regular_file(stat.st_mode as _); #[cfg(target_os = "macos")] { - // `pipe(2)` pipes are S_IFIFO too; XNU's pipe_stat leaves their - // st_dev 0, whereas a FIFO on a filesystem carries its volume's - // device number. + // pipe(2) pipes are S_IFIFO with st_dev == 0 (XNU pipe_stat). self.is_named_pipe = bun_sys::S::ISFIFO(stat.st_mode as _) && stat.st_dev != 0; } self.total_size = @@ -786,8 +781,7 @@ impl ReadFile { if self.could_block { #[cfg(target_os = "macos")] if self.is_named_pipe { - // Waiting before the first read is also what keeps a FIFO - // whose writer has not connected yet from reading as empty. + // Before the first read too: with no writer yet, read() says EOF. if self.block_until_readable() { self.do_read_loop(); } else { diff --git a/src/sys/lib.rs b/src/sys/lib.rs index 977d52985626..17e258ce4e85 100644 --- a/src/sys/lib.rs +++ b/src/sys/lib.rs @@ -1711,11 +1711,8 @@ mod nocancel { ) -> isize; #[link_name = "poll$NOCANCEL"] pub(crate) fn poll(fds: *mut libc::pollfd, nfds: libc::nfds_t, timeout: c_int) -> c_int; - // The `_DARWIN_UNLIMITED_SELECT` variant of select(2) (same - // `$DARWIN_EXTSN` scheme as `realpath` in `posix_impl`). libsyscall - // maps it straight onto the syscall, so there is no FD_SETSIZE check - // and each set is a bitmap of ceil(nfds / 32) 32-bit words rather - // than a `libc::fd_set`; hence the word pointers. + // `_DARWIN_UNLIMITED_SELECT` select(2): the bare syscall, so no + // FD_SETSIZE check, and a set is just ceil(nfds / 32) 32-bit words. #[link_name = "select$DARWIN_EXTSN$NOCANCEL"] pub(crate) fn select( nfds: c_int, @@ -7623,21 +7620,15 @@ pub fn kevent( } } -/// Blocks the calling thread in `select(2)` until `fd` is readable, where -/// readable includes EOF. Retries on EINTR. +/// Blocks in `select(2)` until `fd` is readable or at EOF; retries on EINTR. /// -/// This is how a named pipe (a FIFO opened by path, or one inherited as -/// stdin) has to be waited on under macOS. XNU attaches kqueue `EVFILT_READ` -/// filters for a FIFO to its vnode, and that filter only fires while bytes are -/// buffered (`vnode_readable_data_count`); the last writer closing posts -/// nothing to the vnode, so a kqueue registration never wakes a reader up for -/// EOF, and neither does `poll(2)`, which XNU implements on top of kqueue. -/// `select(2)` instead goes through `fifo_select` to the FIFO's underlying -/// socket, whose readability includes the `SS_CANTRCVMORE` state that the last -/// writer's close sets (and that a writer which has not connected yet leaves -/// clear, so a FIFO that is still waiting for its first writer blocks here -/// rather than reading as empty). `pipe(2)` pipes are not affected: their own -/// kqueue filter reports `EV_EOF`. +/// For named pipes. XNU hooks a FIFO's `EVFILT_READ` (and `poll(2)`, which it +/// implements with kqueue) to the vnode, where it only fires while bytes are +/// buffered, and `fifo_close` posts nothing, so neither ever reports the last +/// writer closing. `select(2)` goes through `fifo_select` to the FIFO's socket, +/// whose readability includes `SS_CANTRCVMORE`: set by that close, clear while +/// no writer has connected yet. `pipe(2)` pipes have their own filter with +/// `EV_EOF` and do not need this. #[cfg(target_os = "macos")] pub fn block_until_readable(fd: Fd) -> Maybe<()> { debug_assert!(fd.is_valid()); diff --git a/test/js/bun/util/bun-file-read.test.ts b/test/js/bun/util/bun-file-read.test.ts index b0127432fbfd..ba4cacc8f111 100644 --- a/test/js/bun/util/bun-file-read.test.ts +++ b/test/js/bun/util/bun-file-read.test.ts @@ -75,12 +75,30 @@ describe.skipIf(isWindows)("reading a named pipe to EOF", () => { }); } + // When a FIFO end gets closed is what these tests are about, so each one is + // closed explicitly at the right moment; `using` only covers the failure paths. + function openFd(file: string, flags: number | string) { + let fd = openSync(file, flags); + return { + get fd() { + return fd; + }, + close() { + if (fd !== -1) closeSync(fd); + fd = -1; + }, + [Symbol.dispose]() { + this.close(); + }, + }; + } + // The child must already have the FIFO open for reading before a writer can // connect to it: a non-blocking open for writing fails with ENXIO until then. - async function openWriterOnceChildIsReading(fifo: string, child: Bun.Subprocess): Promise { + async function openWriterOnceChildIsReading(fifo: string, child: Bun.Subprocess) { while (true) { try { - return openSync(fifo, constants.O_WRONLY | constants.O_NONBLOCK); + return openFd(fifo, constants.O_WRONLY | constants.O_NONBLOCK); } catch (err: any) { if (err.code !== "ENXIO") throw err; } @@ -99,42 +117,32 @@ describe.skipIf(isWindows)("reading a named pipe to EOF", () => { // The write end can only be opened, and written to without EPIPE, while // some reader has the FIFO open; `holder` is that reader until the child // has opened its own. It never reads, so every byte goes to the child. - let holder = openSync(fifo, constants.O_RDONLY | constants.O_NONBLOCK); - const closeHolder = () => { - if (holder !== -1) closeSync(holder); - holder = -1; - }; - let writer = openSync(fifo, "w"); - try { - await using proc = readFifoInChild( - `const bytes = await Bun.file(process.env.FIFO).bytes(); process.stdout.write(bytes.length + " " + Bun.hash(bytes));`, - fifo, - ); - const stderr = proc.stderr.text(); - // The write end is blocking, so this write only completes as the child - // drains the pipe, and closing it afterwards is what ends the child's - // read. The child cannot exit before that unless it failed; dropping - // `holder` then leaves the pipe without readers, so the blocked write - // fails with EPIPE instead of waiting forever. - const childDied = proc.exited.then(async exitCode => { - closeHolder(); - throw new Error(`child exited with ${exitCode} before the payload was written: ${await stderr}`); - }); - const written = await Promise.race([Bun.write(Bun.file(writer), payload), childDied]); - closeSync(writer); - writer = -1; - const [stdout, stderrText, exitCode] = await Promise.all([proc.stdout.text(), stderr, proc.exited]); - - expect({ written, stdout, stderr: stderrText }).toEqual({ - written: payload.length, - stdout: `${payload.length} ${Bun.hash(payload)}`, - stderr: "", - }); - expect(exitCode).toBe(0); - } finally { - if (writer !== -1) closeSync(writer); - closeHolder(); - } + using holder = openFd(fifo, constants.O_RDONLY | constants.O_NONBLOCK); + using writer = openFd(fifo, "w"); + await using proc = readFifoInChild( + `const bytes = await Bun.file(process.env.FIFO).bytes(); process.stdout.write(bytes.length + " " + Bun.hash(bytes));`, + fifo, + ); + const stderr = proc.stderr.text(); + // The write end is blocking, so this write only completes as the child + // drains the pipe, and closing it afterwards is what ends the child's + // read. The child cannot exit before that unless it failed; dropping + // `holder` then leaves the pipe without readers, so the blocked write + // fails with EPIPE instead of waiting forever. + const childDied = proc.exited.then(async exitCode => { + holder.close(); + throw new Error(`child exited with ${exitCode} before the payload was written: ${await stderr}`); + }); + const written = await Promise.race([Bun.write(Bun.file(writer.fd), payload), childDied]); + writer.close(); + const [stdout, stderrText, exitCode] = await Promise.all([proc.stdout.text(), stderr, proc.exited]); + + expect({ written, stdout, stderr: stderrText }).toEqual({ + written: payload.length, + stdout: `${payload.length} ${Bun.hash(payload)}`, + stderr: "", + }); + expect(exitCode).toBe(0); }); it.concurrent("text() waits for a writer that connects after the read started", async () => { @@ -143,9 +151,9 @@ describe.skipIf(isWindows)("reading a named pipe to EOF", () => { mkfifo(fifo); await using proc = readFifoInChild(`process.stdout.write(await Bun.file(process.env.FIFO).text());`, fifo); - const writer = await openWriterOnceChildIsReading(fifo, proc); - writeSync(writer, "written after the reader opened\n"); - closeSync(writer); + using writer = await openWriterOnceChildIsReading(fifo, proc); + writeSync(writer.fd, "written after the reader opened\n"); + writer.close(); const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); expect({ stdout, stderr }).toEqual({ stdout: "written after the reader opened\n", stderr: "" }); @@ -161,7 +169,8 @@ describe.skipIf(isWindows)("reading a named pipe to EOF", () => { `const text = await Bun.file(process.env.FIFO).text(); process.stdout.write(JSON.stringify(text));`, fifo, ); - closeSync(await openWriterOnceChildIsReading(fifo, proc)); + using writer = await openWriterOnceChildIsReading(fifo, proc); + writer.close(); const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); expect({ stdout, stderr }).toEqual({ stdout: '""', stderr: "" }); @@ -175,13 +184,17 @@ describe.skipIf(isWindows)("reading a named pipe to EOF", () => { // Same dance as above: a reader has to exist before the write end can be // opened; here that reader becomes the child's stdin. - const readEnd = openSync(fifo, constants.O_RDONLY | constants.O_NONBLOCK); - const writer = openSync(fifo, "w"); - await using proc = readFifoInChild(`process.stdout.write(JSON.stringify(await Bun.stdin.text()));`, fifo, readEnd); + using readEnd = openFd(fifo, constants.O_RDONLY | constants.O_NONBLOCK); + using writer = openFd(fifo, "w"); + await using proc = readFifoInChild( + `process.stdout.write(JSON.stringify(await Bun.stdin.text()));`, + fifo, + readEnd.fd, + ); // The child has its own descriptor for the read end now. - closeSync(readEnd); - writeSync(writer, "stdin is a named pipe\n"); - closeSync(writer); + readEnd.close(); + writeSync(writer.fd, "stdin is a named pipe\n"); + writer.close(); const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); expect({ stdout, stderr }).toEqual({ stdout: JSON.stringify("stdin is a named pipe\n"), stderr: "" }); From 7afcfadabb0d16f5b88c519b5d7e1ed044049989 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 12 Aug 2026 13:03:12 +0000 Subject: [PATCH 4/7] Shorten the select() comments further --- src/runtime/webcore/blob/read_file.rs | 3 +-- src/sys/lib.rs | 14 +++++--------- 2 files changed, 6 insertions(+), 11 deletions(-) diff --git a/src/runtime/webcore/blob/read_file.rs b/src/runtime/webcore/blob/read_file.rs index 31f4a6149296..58f08a140642 100644 --- a/src/runtime/webcore/blob/read_file.rs +++ b/src/runtime/webcore/blob/read_file.rs @@ -470,8 +470,7 @@ impl ReadFile { } } - /// `wait_for_readable` for named pipes: waits on this (pool) thread. - /// Returns `false` if the wait failed, with the error recorded for `then()`. + /// `wait_for_readable` for named pipes, on this thread; `false` if the wait itself failed. #[cfg(target_os = "macos")] fn block_until_readable(&mut self) -> bool { bloblog!("ReadFile.blockUntilReadable"); diff --git a/src/sys/lib.rs b/src/sys/lib.rs index 17e258ce4e85..2cc3b71bedf0 100644 --- a/src/sys/lib.rs +++ b/src/sys/lib.rs @@ -1711,8 +1711,7 @@ mod nocancel { ) -> isize; #[link_name = "poll$NOCANCEL"] pub(crate) fn poll(fds: *mut libc::pollfd, nfds: libc::nfds_t, timeout: c_int) -> c_int; - // `_DARWIN_UNLIMITED_SELECT` select(2): the bare syscall, so no - // FD_SETSIZE check, and a set is just ceil(nfds / 32) 32-bit words. + // `_DARWIN_UNLIMITED_SELECT` select(2): no FD_SETSIZE limit; a set is ceil(nfds/32) words. #[link_name = "select$DARWIN_EXTSN$NOCANCEL"] pub(crate) fn select( nfds: c_int, @@ -7622,13 +7621,10 @@ pub fn kevent( /// Blocks in `select(2)` until `fd` is readable or at EOF; retries on EINTR. /// -/// For named pipes. XNU hooks a FIFO's `EVFILT_READ` (and `poll(2)`, which it -/// implements with kqueue) to the vnode, where it only fires while bytes are -/// buffered, and `fifo_close` posts nothing, so neither ever reports the last -/// writer closing. `select(2)` goes through `fifo_select` to the FIFO's socket, -/// whose readability includes `SS_CANTRCVMORE`: set by that close, clear while -/// no writer has connected yet. `pipe(2)` pipes have their own filter with -/// `EV_EOF` and do not need this. +/// For named pipes: XNU's kqueue filter for a FIFO (`filt_vnode_common`, which +/// `poll(2)` uses too) fires only while bytes are buffered, never for the last +/// writer closing; `fifo_select` consults the FIFO's socket, which reports that +/// close (and not a writer that has yet to connect). `pipe(2)` pipes report `EV_EOF`. #[cfg(target_os = "macos")] pub fn block_until_readable(fd: Fd) -> Maybe<()> { debug_assert!(fd.is_valid()); From 8bf43ec9de6502eded88d74ae883da068323b536 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 12 Aug 2026 13:15:42 +0000 Subject: [PATCH 5/7] test: bound the wait for the child's FIFO open and name what did not happen --- test/js/bun/util/bun-file-read.test.ts | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/test/js/bun/util/bun-file-read.test.ts b/test/js/bun/util/bun-file-read.test.ts index ba4cacc8f111..8d5f53b86173 100644 --- a/test/js/bun/util/bun-file-read.test.ts +++ b/test/js/bun/util/bun-file-read.test.ts @@ -96,14 +96,18 @@ describe.skipIf(isWindows)("reading a named pipe to EOF", () => { // The child must already have the FIFO open for reading before a writer can // connect to it: a non-blocking open for writing fails with ENXIO until then. async function openWriterOnceChildIsReading(fifo: string, child: Bun.Subprocess) { + const deadline = performance.now() + 10_000; while (true) { try { return openFd(fifo, constants.O_WRONLY | constants.O_NONBLOCK); } catch (err: any) { if (err.code !== "ENXIO") throw err; } - if (child.exitCode !== null || child.signalCode !== null) { - throw new Error(`child exited (${child.exitCode ?? child.signalCode}) without opening the FIFO`); + const exited = child.exitCode ?? child.signalCode; + if (exited !== null || performance.now() > deadline) { + throw new Error( + `nothing opened ${fifo} for reading; child ${exited === null ? "is still running" : `exited (${exited})`}`, + ); } await Bun.sleep(5); } From fc757cdd3833262b6442c3be747eb3ecfbe9eabd Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 12 Aug 2026 17:43:00 +0000 Subject: [PATCH 6/7] read_file: run a named pipe's read as its own pool task, not under the job's VM borrow JobContext::run holds a VM borrow for the duration of the call and VM teardown waits for borrows, so blocking in select() there kept a worker with a read waiting on a FIFO writer from ever being terminated. run_async_with_fd now schedules read_named_pipe_task and returns; the task does the waits and the read loop with no borrow held, like the io-thread hand-off does. Tests: a worker parked on an idle FIFO writer can be terminated; a FIFO read through a descriptor above FD_SETSIZE, which is what the unlimited select variant and the hand-built bitmap are for. --- src/runtime/webcore/blob/read_file.rs | 28 +++++++--- test/js/bun/util/bun-file-read.test.ts | 74 +++++++++++++++++++++++++- 2 files changed, 94 insertions(+), 8 deletions(-) diff --git a/src/runtime/webcore/blob/read_file.rs b/src/runtime/webcore/blob/read_file.rs index 58f08a140642..3876937d7251 100644 --- a/src/runtime/webcore/blob/read_file.rs +++ b/src/runtime/webcore/blob/read_file.rs @@ -484,6 +484,23 @@ impl ReadFile { } } + /// A named pipe's whole read, as its own pool task: `JobContext::run` holds a + /// VM borrow and VM teardown waits for borrows, so waiting for a writer must + /// not happen inside it. Waiting before the first read also keeps a FIFO + /// with no writer yet from reading as empty. + #[cfg(target_os = "macos")] + fn read_named_pipe_task(task: *mut WorkPoolTask) { + // SAFETY: only reached via `WorkPoolTask::callback` with `task` = + // `&mut self.task` (intrusive) scheduled by `run_async_with_fd`; + // recover parent. + let this = unsafe { &mut *ReadFile::from_task_ptr(task) }; + if this.block_until_readable() { + this.do_read_loop(); + } else { + this.on_finish(); + } + } + /// Pick the read target: `buffer`'s spare capacity if it is at least as /// large as `stack_buffer`, otherwise `stack_buffer`; capped by /// `max_length - read_off`. Returns `(use_stack, target)` so the caller @@ -780,12 +797,11 @@ impl ReadFile { if self.could_block { #[cfg(target_os = "macos")] if self.is_named_pipe { - // Before the first read too: with no writer yet, read() says EOF. - if self.block_until_readable() { - self.do_read_loop(); - } else { - self.on_finish(); - } + self.task = WorkPoolTask { + node: Default::default(), + callback: Self::read_named_pipe_task, + }; + WorkPool::schedule(&raw mut self.task); return; } diff --git a/test/js/bun/util/bun-file-read.test.ts b/test/js/bun/util/bun-file-read.test.ts index 8d5f53b86173..1da3edb1d38b 100644 --- a/test/js/bun/util/bun-file-read.test.ts +++ b/test/js/bun/util/bun-file-read.test.ts @@ -58,13 +58,13 @@ describe("Bun.file read-loop target selection", () => { }); }); -// Whole-file reads of a named pipe. Every one of these ends with the reader +// Whole-file reads of a named pipe. All but the last end with the reader // having drained the pipe and then learning that the last writer closed; on // macOS that EOF is invisible to kqueue and poll(2) (they only see buffered // bytes on a FIFO), so the reader has to wait for it differently than it does // for a pipe(2) pipe, and each of these used to leave the child blocked // forever there. -describe.skipIf(isWindows)("reading a named pipe to EOF", () => { +describe.skipIf(isWindows)("reading a named pipe", () => { function readFifoInChild(script: string, fifo: string, stdin: number | "ignore" = "ignore") { return Bun.spawn({ cmd: [bunExe(), "-e", script], @@ -204,4 +204,74 @@ describe.skipIf(isWindows)("reading a named pipe to EOF", () => { expect({ stdout, stderr }).toEqual({ stdout: JSON.stringify("stdin is a named pipe\n"), stderr: "" }); expect(exitCode).toBe(0); }); + + // The macOS wait is a select(2); this is the case where a plain fd_set + // would not hold the descriptor (FD_SETSIZE is 1024 there). + it.concurrent("Bun.file(fd) reads a FIFO whose descriptor number is above FD_SETSIZE", async () => { + using dir = tempDir("bun-file-read-fifo-high-fd", {}); + const fifo = path.join(String(dir), "high.fifo"); + mkfifo(fifo); + + await using proc = readFifoInChild( + `import { constants, openSync } from "node:fs"; + while (openSync("/dev/null", "r") < 1024) {} + const fd = openSync(process.env.FIFO, constants.O_RDONLY | constants.O_NONBLOCK); + const text = await Bun.file(fd).text(); + process.stdout.write(JSON.stringify({ fd, text }));`, + fifo, + ); + using writer = await openWriterOnceChildIsReading(fifo, proc); + writeSync(writer.fd, "read through a high fd\n"); + writer.close(); + + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect(stderr).toBe(""); + const { fd, text } = JSON.parse(stdout); + expect(fd).toBeGreaterThanOrEqual(1024); + expect(text).toBe("read through a high fd\n"); + expect(exitCode).toBe(0); + }); + + // A read that is waiting for a writer which never shows up must not keep + // its VM from shutting down: the wait has to happen outside the job's VM + // borrow, which terminate() waits for. + it.concurrent("terminate() completes while a worker's read is waiting on an idle writer", async () => { + using dir = tempDir("bun-file-read-fifo-worker", { + "main.ts": ` + const worker = new Worker(new URL("./worker.ts", import.meta.url).href); + const closed = new Promise(resolve => worker.addEventListener("close", resolve)); + worker.addEventListener("message", ({ data }) => console.log(data)); + // Our stdin is closed once the test has connected a writer to the FIFO. + await Bun.stdin.text(); + worker.terminate(); + await closed; + console.log("terminated"); + `, + "worker.ts": ` + const read = Bun.file(process.env.FIFO!).text(); + postMessage("reading"); + await read; + postMessage("the read finished, which it should not have"); + `, + }); + const fifo = path.join(String(dir), "idle.fifo"); + mkfifo(fifo); + + await using proc = Bun.spawn({ + cmd: [bunExe(), "main.ts"], + cwd: String(dir), + env: { ...bunEnv, FIFO: fifo }, + stdin: "pipe", + stdout: "pipe", + stderr: "pipe", + }); + // Connecting a writer proves the worker's read has the FIFO open; never + // writing to it keeps that read waiting for the rest of the test. + using _idleWriter = await openWriterOnceChildIsReading(fifo, proc); + proc.stdin.end(); + + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout, stderr }).toEqual({ stdout: "reading\nterminated\n", stderr: "" }); + expect(exitCode).toBe(0); + }); }); From dc62568ae1642dbe36ece8aa152b777499fc0af5 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 12 Aug 2026 22:16:42 +0000 Subject: [PATCH 7/7] test: name the worker test's spawned files *.fixture.ts --- test/js/bun/util/bun-file-read.test.ts | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/test/js/bun/util/bun-file-read.test.ts b/test/js/bun/util/bun-file-read.test.ts index 1da3edb1d38b..c7785cf8738a 100644 --- a/test/js/bun/util/bun-file-read.test.ts +++ b/test/js/bun/util/bun-file-read.test.ts @@ -237,8 +237,8 @@ describe.skipIf(isWindows)("reading a named pipe", () => { // borrow, which terminate() waits for. it.concurrent("terminate() completes while a worker's read is waiting on an idle writer", async () => { using dir = tempDir("bun-file-read-fifo-worker", { - "main.ts": ` - const worker = new Worker(new URL("./worker.ts", import.meta.url).href); + "main.fixture.ts": ` + const worker = new Worker(new URL("./worker.fixture.ts", import.meta.url).href); const closed = new Promise(resolve => worker.addEventListener("close", resolve)); worker.addEventListener("message", ({ data }) => console.log(data)); // Our stdin is closed once the test has connected a writer to the FIFO. @@ -247,7 +247,7 @@ describe.skipIf(isWindows)("reading a named pipe", () => { await closed; console.log("terminated"); `, - "worker.ts": ` + "worker.fixture.ts": ` const read = Bun.file(process.env.FIFO!).text(); postMessage("reading"); await read; @@ -258,7 +258,7 @@ describe.skipIf(isWindows)("reading a named pipe", () => { mkfifo(fifo); await using proc = Bun.spawn({ - cmd: [bunExe(), "main.ts"], + cmd: [bunExe(), "main.fixture.ts"], cwd: String(dir), env: { ...bunEnv, FIFO: fifo }, stdin: "pipe",