diff --git a/src/runtime/webcore/Blob.rs b/src/runtime/webcore/Blob.rs index ec3d76a55648..29744cdc05a0 100644 --- a/src/runtime/webcore/Blob.rs +++ b/src/runtime/webcore/Blob.rs @@ -5445,8 +5445,19 @@ fn write_string_to_file_fast( bun_sys::Result::Err(err) => { truncate.set(false); if err.get_errno() == bun_sys::E::EAGAIN { - *needs_async = true; - return JSValue::ZERO; + if written.get() == 0 { + *needs_async = true; + return JSValue::ZERO; + } + // The async path re-sends the whole input, so after a partial + // write finish synchronously via poll(POLLOUT). + let mut pfd = [bun_sys::posix::PollFd { + fd: fd.native(), + events: bun_sys::posix::POLL_OUT, + revents: 0, + }]; + let _ = bun_sys::posix::poll(&mut pfd, -1); + continue; } let err_js = if !NEEDS_OPEN { err.to_js(global_this) @@ -5520,8 +5531,17 @@ fn write_bytes_to_file_fast( bun_sys::Result::Err(err) => { #[cfg(not(windows))] if err.get_errno() == bun_sys::E::EAGAIN { - *_needs_async = true; - return JSValue::ZERO; + if written == 0 { + *_needs_async = true; + return JSValue::ZERO; + } + let mut pfd = [bun_sys::posix::PollFd { + fd: fd.native(), + events: bun_sys::posix::POLL_OUT, + revents: 0, + }]; + let _ = bun_sys::posix::poll(&mut pfd, -1); + continue; } let err_js = if !NEEDS_OPEN { err.to_js(global_this) diff --git a/src/runtime/webcore/blob/read_file.rs b/src/runtime/webcore/blob/read_file.rs index c7469512233d..5b663901938e 100644 --- a/src/runtime/webcore/blob/read_file.rs +++ b/src/runtime/webcore/blob/read_file.rs @@ -508,54 +508,41 @@ impl ReadFile { /// Never touches `self.buffer`; the caller moves it out for the duration. pub fn do_read(&mut self, buf: &mut [u8], read_len: &mut usize, retry: &mut bool) -> bool { - let result: bun_sys::Result = 'brk: { - if bun_sys::S::ISSOCK(self.file_store.mode) { - break 'brk bun_sys::recv_non_block(self.opened_fd, buf); - } - break 'brk bun_sys::read(self.opened_fd, buf); + let result: bun_sys::Result = if bun_sys::S::ISSOCK(self.file_store.mode) { + bun_sys::recv_non_block(self.opened_fd, buf) + } else { + bun_sys::read(self.opened_fd, buf) }; - loop { - match &result { - Ok(res) => { - *read_len = *res as usize; // @truncate — usize→usize is identity here - self.read_eof = *res == 0; + match result { + Ok(res) => { + *read_len = res; + self.read_eof = res == 0; + true + } + Err(err) => match err.get_errno() { + e if e == io::RETRY => { + self.could_block = true; + *retry = true; + self.read_eof = false; + true } - Err(err) => { - match err.get_errno() { - e if e == io::RETRY => { - if !self.could_block { - // regular files cannot use epoll. - // this is fine on kqueue, but not on epoll. - continue; - } - *retry = true; - self.read_eof = false; - return true; - } - _ => { - self.errno = Some(bun_errno::from_errno(err.errno as i32).into()); - self.system_error = Some(err.to_system_error().into()); - if self.system_error.as_ref().unwrap().path.is_empty() { - self.system_error.as_mut().unwrap().path = - if self.file_store.pathlike.is_path() { - BunString::clone_utf8( - self.file_store.pathlike.path().slice(), - ) - .into() - } else { - BunString::EMPTY.into() - }; - } - return false; - } + _ => { + self.errno = Some(bun_errno::from_errno(err.errno as i32).into()); + self.system_error = Some(err.to_system_error().into()); + if self.system_error.as_ref().unwrap().path.is_empty() { + self.system_error.as_mut().unwrap().path = + if self.file_store.pathlike.is_path() { + BunString::clone_utf8(self.file_store.pathlike.path().slice()) + .into() + } else { + BunString::EMPTY.into() + }; } + false } - } - break; + }, } - - true } pub fn then(this: Box, _: &JSGlobalObject) -> jsc::JsTerminatedResult<()> { diff --git a/src/runtime/webcore/blob/write_file.rs b/src/runtime/webcore/blob/write_file.rs index 1f6795af3729..f16c72b45fb1 100644 --- a/src/runtime/webcore/blob/write_file.rs +++ b/src/runtime/webcore/blob/write_file.rs @@ -338,35 +338,24 @@ impl WriteFile { // // On macOS, it is an error to use pwrite() on a // non-seekable file. - let result: bun_sys::Result = - sys::write(fd, &self.bytes_blob.shared_view()[off..off + len]); - - loop { - match &result { - bun_sys::Result::Ok(res) => { - *wrote = *res; - self.total_written += *res; - } - bun_sys::Result::Err(err) => { - if err.get_errno() == io::RETRY { - if !self.could_block { - // regular files cannot use epoll. - // this is fine on kqueue, but not on epoll. - continue; - } - self.wait_for_writable(); - return false; - } else { - self.errno = Some(bun_errno::from_errno(err.errno as i32).into()); - self.system_error = Some(err.to_system_error().into()); - return false; - } + match sys::write(fd, &self.bytes_blob.shared_view()[off..off + len]) { + bun_sys::Result::Ok(res) => { + *wrote = res; + self.total_written += res; + true + } + bun_sys::Result::Err(err) => { + if err.get_errno() == io::RETRY { + // EAGAIN is impossible on a regular file, so the fd is pollable. + self.could_block = true; + self.wait_for_writable(); + } else { + self.errno = Some(bun_errno::from_errno(err.errno as i32).into()); + self.system_error = Some(err.to_system_error().into()); } + false } - break; } - - true } pub fn then(mut this: Box, _global: &JSGlobalObject) -> Result<(), JsTerminated> { @@ -414,7 +403,8 @@ impl WriteFile { } pub fn is_allowed_to_close(&self) -> bool { - self.file_blob + match self + .file_blob .store .get() .as_ref() @@ -422,7 +412,11 @@ impl WriteFile { .data .as_file() .pathlike - .is_path() + { + PathOrFileDescriptor::Path(_) => true, + // opened_fd differs from the caller's fd only when run_with_fd duped it. + PathOrFileDescriptor::Fd(fd) => self.opened_fd != Fd::INVALID && self.opened_fd != fd, + } } #[cfg(not(windows))] @@ -450,26 +444,37 @@ impl WriteFile { let fd = self.opened_fd; - self.could_block = 'brk: { - if let Some(store) = self.file_blob.store.get().as_ref() { - if let blob::store::Data::File(file) = &store.data { - if file.pathlike.is_fd() { - // If seekable was set, then so was mode - if file.seekable.is_some() { - // This is mostly to handle pipes which were passsed to the process somehow - // such as stderr, stdout. Bun.stdin and Bun.stderr will automatically set `mode` for us. - break 'brk !bun_sys::is_regular_file(file.mode); - } + let caller_supplied_fd = match self.file_blob.store.get().as_ref() { + Some(store) => match &store.data { + blob::store::Data::File(file) => match file.pathlike { + PathOrFileDescriptor::Fd(_) => { + self.could_block = file.mode != 0 && !bun_sys::is_regular_file(file.mode); + true } + PathOrFileDescriptor::Path(_) => { + // Opened with O_NONBLOCK; don't fstat. + self.could_block = false; + false + } + }, + _ => false, + }, + None => false, + }; + + // The IO thread's epoll/kqueue keys interest by fd number, so concurrent WriteFiles on + // one caller-supplied fd would collide; dup so this instance polls a private fd number. + if caller_supplied_fd { + match bun_sys::dup(fd) { + bun_sys::Result::Ok(duped) => self.opened_fd = duped, + bun_sys::Result::Err(err) => { + self.errno = Some(bun_errno::from_errno(err.errno as i32).into()); + self.system_error = Some(err.to_system_error().into()); + self.on_finish(); + return; } } - - // We opened the file descriptor with O_NONBLOCK, so we - // shouldn't have to worry about blocking reads/writes - // - // We do not call fstat() because that is very expensive. - false - }; + } // We have never supported offset in Bun.write(). // and properly adding support means we need to also support it diff --git a/test/js/bun/io/bun-write.test.js b/test/js/bun/io/bun-write.test.js index a4c3fb2f551b..3dd4d8a6932d 100644 --- a/test/js/bun/io/bun-write.test.js +++ b/test/js/bun/io/bun-write.test.js @@ -728,3 +728,51 @@ int posix_fadvise(int fd, off_t offset, off_t len, int advice) { expect(f.name).toBe(filePath); }); }); + +// Once fd 1 is O_NONBLOCK (process.stdout.write's FileSink sets it on the shared open file +// description via a dup), a Bun.write that overflowed the pipe buffer used to wedge a +// thread-pool worker at 100% CPU: the async WriteFile path's could_block stayed false for +// stdio and its EAGAIN handler re-matched a cached write() result without re-issuing the +// syscall. Runs outside describe.concurrent so an unfixed build's spin doesn't starve +// neighbours into their default timeout. +it.skipIf(isWindows)( + "Bun.write(Bun.stdout, ...) to a full nonblocking pipe completes", + async () => { + // Below 256 KiB exercises the sync fast path's EAGAIN -> needs_async fallback; at/above + // 256 KiB goes straight to the thread-pool WriteFile path. + const script = ` + process.stdout.write("x"); // constructs the fd 1 FileSink, which flips the pipe O_NONBLOCK + const fdDir = process.platform === "darwin" ? "/dev/fd" : "/proc/self/fd"; + const fdCount = () => require("fs").readdirSync(fdDir).length; + const small = Buffer.alloc(64 * 1024, 65).toString(); + const large = Buffer.alloc(256 * 1024, 66).toString(); + let wrote = 0, fdsBefore, fdsAfter; + for (let round = 0; round < 2; round++) { + fdsBefore = fdCount(); + const ps = []; + for (let i = 0; i < 16; i++) ps.push(Bun.write(Bun.stdout, small)); + for (let i = 0; i < 4; i++) ps.push(Bun.write(Bun.stdout, large)); + for (const n of await Promise.all(ps)) wrote += n; + fdsAfter = fdCount(); + } + process.stderr.write("wrote=" + wrote + " fdDelta=" + (fdsAfter - fdsBefore)); + `; + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", script], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + timeout: 20_000, + killSignal: "SIGKILL", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.bytes(), proc.stderr.text(), proc.exited]); + const expected = 1 + 2 * (16 * 64 * 1024 + 4 * 256 * 1024); + expect({ length: stdout.length, stderr, exitCode, signalCode: proc.signalCode }).toEqual({ + length: expected, + stderr: "wrote=" + (expected - 1) + " fdDelta=0", + exitCode: 0, + signalCode: null, + }); + }, + 30_000, +);