diff --git a/CHANGELOG.md b/CHANGELOG.md index 1e60aea4ff71..7405e6d15942 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,8 @@ ## Unreleased +- Publish ready POSIX subprocess exits and drain their owned pipes before stalled-host retirement observers, without changing the global event-loop phase order; wait for extra stdio pipes before emitting child `close`. + - Release retained heap-snapshot metadata when a local inspector Session disables `HeapProfiler` or disconnects, including sessions that never enabled the domain. - Follow directory symlinks in Windows asynchronous recursive `fs.readdir` callback results and promise string results while preserving promise Dirent traversal boundaries. diff --git a/Cargo.lock b/Cargo.lock index cd424128a7f7..7888e3ac3b6c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1451,6 +1451,7 @@ version = "0.0.0" dependencies = [ "bun_alloc", "bun_analytics", + "bun_collections", "bun_core", "bun_crash_handler", "bun_dispatch", diff --git a/docs/runtime/nodejs-compat.mdx b/docs/runtime/nodejs-compat.mdx index f6c15bd56ec5..7e776e331f9c 100644 --- a/docs/runtime/nodejs-compat.mdx +++ b/docs/runtime/nodejs-compat.mdx @@ -143,6 +143,8 @@ Child IPC dispatch calls the JavaScript `process.emit` property for `message`, ` Child stdout/stderr `ref()` and `unref()` set the pipe reference state and return the stream. Repeated calls are idempotent, so one `unref()` releases the reference after any number of `ref()` calls. +On POSIX, ready child exits and owned pipe EOF remain observable after a blocked I/O callback, before a deadline timer's deferred observer. Extra stdio pipes participate in child `close` accounting. A descendant that retains a pipe keeps it open until real EOF; child exit does not synthesize EOF. + Child construction publishes the `child_process` diagnostics channel with the exact child object. Spawning publishes `child_process.spawn` tracing start, end, and error events before the corresponding child events. Unspawned stdio getters are safe; failed spawns drain readable pipes to EOF and retain fork IPC methods. Exit listeners run before stdio resumes, and POSIX PATH permission failures report `EACCES`. 🟡 IPC can send `net.Socket`, `net.Server` and `dgram.Socket` handles (including to and from Node.js processes), but not `http` server sockets. `serialization: "advanced"` only works between Bun processes, so use JSON serialization for Node.js ↔ Bun IPC. Missing `subprocess.channel.ref()`/`unref()`. You cannot pass a child's `stdout`/`stderr` as another child's `stdio`, and `spawnSync` does not return extra `stdio` pipes in `output`. diff --git a/packages/bun-usockets/src/internal/internal.h b/packages/bun-usockets/src/internal/internal.h index af0274a77513..b55889f1e342 100644 --- a/packages/bun-usockets/src/internal/internal.h +++ b/packages/bun-usockets/src/internal/internal.h @@ -173,6 +173,7 @@ extern struct addrinfo_result *Bun__addrinfo_getRequestResult(struct addrinfo_re * LIBUS_POLL_HANGUP = epoll EPOLLHUP, both directions down, re-reported until the fd is closed. */ #define LIBUS_POLL_EOF 1 #define LIBUS_POLL_HANGUP 2 +#define LIBUS_SOCKET_OWNER_READ 0x10000 void us_internal_dispatch_ready_poll(struct us_poll_t *p, int error, int eof, int events); void us_internal_timer_sweep(us_loop_r loop); void us_internal_enable_sweep_timer(struct us_loop_t *loop); diff --git a/packages/bun-usockets/src/libusockets.h b/packages/bun-usockets/src/libusockets.h index 430f8f1a808c..72c941b8a489 100644 --- a/packages/bun-usockets/src/libusockets.h +++ b/packages/bun-usockets/src/libusockets.h @@ -724,6 +724,7 @@ void us_socket_remote_address(us_socket_r s, char *nonnull_arg buf, int *nonnull void us_socket_local_address(us_socket_r s, char *nonnull_arg buf, int *nonnull_arg length) nonnull_fn_decl; int us_socket_ipc_write_fd(us_socket_r s, const char *data, int length, int fd) nonnull_fn_decl; +void us_socket_drain_readable(us_socket_r s) nonnull_fn_decl; void us_socket_sendfile_needs_more(us_socket_r s) nonnull_fn_decl; void *us_listen_socket_ext(struct us_listen_socket_t *ls) nonnull_fn_decl; LIBUS_SOCKET_DESCRIPTOR us_listen_socket_get_fd(struct us_listen_socket_t *ls) nonnull_fn_decl; diff --git a/packages/bun-usockets/src/loop.c b/packages/bun-usockets/src/loop.c index 7b856f404f2c..0470a526c033 100644 --- a/packages/bun-usockets/src/loop.c +++ b/packages/bun-usockets/src/loop.c @@ -802,6 +802,14 @@ void us_internal_dispatch_ready_poll(struct us_poll_t *p, int error, int eof, in : us_dispatch_data(s, loop->data.recv_buf + LIBUS_RECV_BUFFER_PADDING, length); /* After socket adoption, track the new socket; the old one becomes invalid */ s = us_internal_socket_follow_adopted(s); + /* A descendant can keep the pipe writable after its + * direct parent exits; bound this owner-only drain. */ + if (events & LIBUS_SOCKET_OWNER_READ) { + if (s && !us_socket_is_closed(s) && !s->flags.is_paused && ++repeat_recv_count < 64) { + continue; + } + break; + } // loop->num_ready_polls isn't accessible on Windows. #ifndef WIN32 // rare case: we're reading a lot of data, there's more to be read, and either: diff --git a/packages/bun-usockets/src/socket.c b/packages/bun-usockets/src/socket.c index 284dab2bf9c3..932cc95d389d 100644 --- a/packages/bun-usockets/src/socket.c +++ b/packages/bun-usockets/src/socket.c @@ -44,6 +44,18 @@ int us_socket_remote_port(struct us_socket_t *s) { } } +/* An exited subprocess may still have unread pipe bytes. A descendant can + * retain the peer: only recv()==0 is EOF, never would-block. */ +void us_socket_drain_readable(struct us_socket_t *s) { + if (us_socket_is_closed(s) || s->ssl || s->flags.is_paused || s->read_eof) { + return; + } + struct us_loop_t *loop = s->group->loop; + loop->data.tick_depth++; + us_internal_dispatch_ready_poll(&s->p, 0, 0, LIBUS_SOCKET_READABLE | LIBUS_SOCKET_OWNER_READ); + loop->data.tick_depth--; +} + void us_socket_shutdown_read(struct us_socket_t *s) { /* This syscall is idempotent so no extra check is needed */ bsd_shutdown_socket_read(us_poll_fd((struct us_poll_t *) s)); diff --git a/src/js/node/child_process.ts b/src/js/node/child_process.ts index b1d7bc96db24..2be300b62047 100644 --- a/src/js/node/child_process.ts +++ b/src/js/node/child_process.ts @@ -1238,8 +1238,17 @@ class ChildProcess extends EventEmitter { #flushStdio() { const stdio = this.stdio; if (stdio === undefined) return; - for (const stream of stdio) { - if (stream?.readable) stream.resume(); + for (let i = 0; i < stdio.length; i++) { + const stream = stdio[i]; + if (stream?.readable) { + stream.resume(); + if (i > 2 && process.platform !== "win32") { + const handle = this.#extraStdioHandles[i]; + if (handle && stream._handle === handle) { + $rust("node_net_binding.rs", "drainSubprocessSocket")(handle); + } + } + } } } @@ -1364,14 +1373,23 @@ class ChildProcess extends EventEmitter { default: switch (io) { case "pipe": - case "socket-fd": + case "socket-fd": { if (!NetModule) NetModule = require("node:net"); // #spawn mapped "pipe" at i>=3 to "socket-fd", so the parent-end // fd in handle.stdio[i] is UnownedFd: we own it and // net.connect({fd}) -> usockets will close it on socket close. const fd = handle && handle.stdio[i]; if (fd == null) return null; - return NetModule.connect({ fd }); + const socket = NetModule.connect({ fd }); + const nativeHandle = socket._handle; + this.#extraStdioHandles[i] = nativeHandle; + this.#closesNeeded++; + socket.once("close", () => { + this.#extraStdioHandles[i] = undefined; + this.#maybeClose(); + }); + return socket; + } } return null; } @@ -1382,6 +1400,7 @@ class ChildProcess extends EventEmitter { #stderr; #stdioObject; #stdioOptions; + #extraStdioHandles: unknown[] = []; #createStdioObject() { const opts = this.#stdioOptions; diff --git a/src/jsc/event_loop.rs b/src/jsc/event_loop.rs index 6d22a86c1bb8..912e6aca436a 100644 --- a/src/jsc/event_loop.rs +++ b/src/jsc/event_loop.rs @@ -47,6 +47,8 @@ pub type Queue = pub struct EventLoop { pub tasks: Queue, + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + process_completions: Option>, /// Set when teardown releases the queue: from then on `enqueue_task` /// releases instead of parking (nothing will tick this loop again). closed_for_tasks: bool, @@ -116,6 +118,8 @@ impl Default for EventLoop { fn default() -> Self { Self { tasks: Queue::init(), + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + process_completions: None, closed_for_tasks: false, immediate_tasks: Vec::new(), next_immediate_tasks: Vec::new(), @@ -772,10 +776,28 @@ impl EventLoop { pub fn tick(&mut self) { jsc::mark_binding(); crate::top_scope!(scope, self.global_ref()); + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + let ready = if self.entered_event_loop_count == 0 + && self.uws_loop.is_none_or(|loop_| { + // SAFETY: this event loop owns the native loop on this thread. + unsafe { (*loop_.as_ptr()).internal_loop_data.tick_depth == 0 } + }) { + self.process_completions + .as_ref() + .map(|queue| queue.capture()) + .unwrap_or_default() + } else { + Vec::new() + }; self.entered_event_loop_count += 1; // `Err(Stopped)`: a fold or checkpoint met the VM's termination; the turn is over. let _ = self.tick_turn(&mut scope); self.entered_event_loop_count -= 1; + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + for process in ready { + // SAFETY: capture retained this owner until the snapshot ref drops. + unsafe { bun_spawn::Process::publish_captured_completion(process.as_ptr()) }; + } } fn tick_turn(&mut self, scope: &mut crate::TopExceptionScope) -> Result<(), Stopped> { @@ -947,6 +969,10 @@ impl EventLoop { } pub fn deinit(&mut self) { + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + { + self.process_completions = None; + } // Everything queued was released by `release_queued_tasks` (which // also made later enqueues release on arrival) and refused posts never // reach `concurrent_tasks`; nothing can be left to leak with the VM box. @@ -1226,6 +1252,31 @@ impl EventLoop { } } +#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] +#[unsafe(no_mangle)] +unsafe fn __bun_watch_process_completion( + process: *mut bun_spawn::Process, + event_loop: EventLoopHandle, +) -> bun_sys::Result<()> { + let (tag, pointer) = event_loop.into_tag_ptr(); + if tag != 1 { + return Ok(()); + } + let event_loop = pointer.cast::(); + // SAFETY: tag 1 denotes the live JS loop passed by Process::watch. + // spawnSync's private loop cannot run an asynchronous retirement observer. + if unsafe { (*event_loop).isolated_poster.is_some() } { + return Ok(()); + } + // SAFETY: registration runs on the owning thread, without JS callbacks. + let slot = unsafe { &mut (*event_loop).process_completions }; + if slot.is_none() { + *slot = Some(bun_spawn::completion::CompletionQueue::new()?); + } + // SAFETY: the primary watch retains process; the event loop retains this queue. + unsafe { slot.as_ref().unwrap().add(process) } +} + impl EventLoop { /// # Safety /// `done` must point to a live `bool`; C++ writes `true` through it from a diff --git a/src/runtime/api/bun/subprocess.rs b/src/runtime/api/bun/subprocess.rs index 56d018ae01a5..c0023edde5d2 100644 --- a/src/runtime/api/bun/subprocess.rs +++ b/src/runtime/api/bun/subprocess.rs @@ -1087,6 +1087,46 @@ impl Subprocess<'_> { } } + // Accessing stdout/stderr transfers the reader to the cached stream. + // The original Readable is then Closed, but this is still our pipe. + if !this_jsvalue.is_empty() { + for value in [ + js::stdout_get_cached(this_jsvalue), + js::stderr_get_cached(this_jsvalue), + ] + .into_iter() + .flatten() + { + if let Some(stream) = crate::webcore::ReadableStream::from_js_direct(value) { + if let crate::webcore::readable_stream::Source::File(file) = stream.ptr { + // SAFETY: the cached stream roots its source. BufferedReader's + // read pins that source while callbacks can re-enter it. + let reader = unsafe { (*file).reader.get() }; + // SAFETY: the rooted source owns this reader; read pins it across callbacks. + unsafe { + if !(*reader).is_done() { + (*reader).unpause(); + bun_io::BufferedReader::read(reader); + } + } + } + } + value.ensure_still_alive(); + } + } + + #[cfg(unix)] + if let Some(ipc) = self.ipc_data.get().clone() { + let socket = match *ipc.socket.get() { + IPC::SocketUnion::Open(socket) => socket.socket.get(), + _ => None, + }; + if let Some(socket) = socket { + // SAFETY: the cloned IPC owner retains the socket across its final callbacks. + unsafe { bun_uws::us_socket_t::drain_readable(socket) }; + } + } + // When Bun itself killed the child (timeout/maxBuffer/AbortSignal) stop // waiting on pipe EOF after the drain above: a grandchild may still // hold the write end and the caller already opted into a bounded wait. diff --git a/src/runtime/node/node_net_binding.rs b/src/runtime/node/node_net_binding.rs index 34efbf801d68..d13951573719 100644 --- a/src/runtime/node/node_net_binding.rs +++ b/src/runtime/node/node_net_binding.rs @@ -13,6 +13,39 @@ use crate::socket::{Listener, NativeCallbacks, NewSocket, SocketFlags, TCPSocket static AUTO_SELECT_FAMILY_DEFAULT: AtomicBool = AtomicBool::new(true); +pub(crate) fn drain_subprocess_socket(global: &JSGlobalObject) -> JSValue { + #[bun_jsc::host_fn] + fn drain(_global: &JSGlobalObject, frame: &CallFrame) -> JsResult { + #[cfg(unix)] + { + let [value] = frame.arguments_as_array::<1>(); + if let Some(socket) = value.as_::() { + // SAFETY: the argument is a live wrapper; copy its handle before + // entering callbacks, which may detach or replace it. + let (flags, handle) = unsafe { ((*socket).flags.get(), (*socket).socket.get()) }; + if !flags.contains(SocketFlags::BYPASS_TLS) { + if let Some(raw) = handle.socket.get() { + // SAFETY: the JS argument keeps the wrapper live; uSockets + // retains a closed native socket across nested callbacks. + unsafe { uws::us_socket_t::drain_readable(raw) }; + value.ensure_still_alive(); + } + } + } + } + #[cfg(windows)] + let _ = frame; + Ok(JSValue::UNDEFINED) + } + JSFunction::create( + global, + "drainSubprocessSocket", + __jsc_host_drain, + 1, + Default::default(), + ) +} + // This is only used to provide the getDefaultAutoSelectFamilyAttemptTimeout and // setDefaultAutoSelectFamilyAttemptTimeout functions, not currently read by any other code. It's // `threadlocal` because Node.js expects each Worker to have its own copy of this, and currently diff --git a/src/spawn/Cargo.toml b/src/spawn/Cargo.toml index 588f1e18d15f..2018e60d60c6 100644 --- a/src/spawn/Cargo.toml +++ b/src/spawn/Cargo.toml @@ -27,3 +27,6 @@ bun_ptr.workspace = true bun_spawn_sys.workspace = true bun_sys.workspace = true bun_threading.workspace = true + +[target.'cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))'.dependencies] +bun_collections.workspace = true diff --git a/src/spawn/completion.rs b/src/spawn/completion.rs new file mode 100644 index 000000000000..4227920da619 --- /dev/null +++ b/src/spawn/completion.rs @@ -0,0 +1,310 @@ +//! A process-only readiness snapshot. No socket or file callbacks run here. + +use std::cell::{Cell, RefCell}; +use std::rc::Rc; + +use bun_collections::HashMap; +use bun_ptr::RefPtr; +use bun_sys::{Fd, FdExt as _}; + +use crate::Process; + +pub struct CompletionQueue { + fd: Cell>, + watched: RefCell>, + retry: RefCell>, +} + +impl CompletionQueue { + pub fn new() -> bun_sys::Result> { + Ok(Rc::new(Self { + fd: Cell::new(Some(open_descriptor()?)), + watched: RefCell::new(HashMap::new()), + retry: RefCell::new(HashMap::new()), + })) + } + + /// # Safety + /// The primary watch retains `process` until it removes this registration. + /// The queue is held by its event loop, on this thread, until teardown. + pub unsafe fn add(&self, process: *mut Process) -> bun_sys::Result<()> { + if self.watched.borrow().contains(&process) { + return Ok(()); + } + if self.fd.get().is_none() { + self.fd.set(Some(open_descriptor()?)); + } + #[cfg(target_os = "macos")] + let result = { + let change = libc::kevent64_s { + // SAFETY: the caller holds the primary watch ref. + ident: unsafe { (*process).pid } as u64, + filter: libc::EVFILT_PROC, + flags: libc::EV_ADD | libc::EV_ONESHOT, + fflags: libc::NOTE_EXIT, + data: 0, + udata: process as u64, + ext: [0; 2], + }; + // SAFETY: live queue, stack-local changelist, no output events. + retry_interrupted(|| unsafe { + libc::kevent64( + self.fd.get().unwrap().native(), + &raw const change, + 1, + std::ptr::null_mut(), + 0, + 0, + std::ptr::null(), + ) + }) + }; + #[cfg(any(target_os = "linux", target_os = "android"))] + let result = { + // Only this private queue is one-shot: capture retains every + // returned owner before any callback can re-enter the loop. + let mut event = libc::epoll_event { + events: (libc::EPOLLIN | libc::EPOLLONESHOT) as u32, + u64: process as u64, + }; + // SAFETY: primary watch holds process/pidfd; event is stack-local. + retry_interrupted(|| unsafe { + libc::epoll_ctl( + self.fd.get().unwrap().native(), + libc::EPOLL_CTL_ADD, + (*process).pidfd, + &raw mut event, + ) + }) + }; + if result < 0 { + let error = last_error(); + #[cfg(target_os = "macos")] + if error.get_errno() == bun_sys::E::ESRCH { + self.retry.borrow_mut().insert(process, ()); + self.watched.borrow_mut().insert(process, ()); + // SAFETY: caller contract; Rc keeps this shared backref stable. + unsafe { (*process).completion_queue = self }; + return Ok(()); + } + if self.watched.borrow().is_empty() { + if let Some(fd) = self.fd.take() { + close_owned_descriptor(fd); + } + } + return Err(error); + } + self.watched.borrow_mut().insert(process, ()); + // SAFETY: caller contract; removed before process or queue destruction. + unsafe { (*process).completion_queue = self }; + Ok(()) + } + + /// # Safety + /// `process` is live and its registration, if any, belongs to this queue. + pub unsafe fn remove(&self, process: *mut Process) { + if self.watched.borrow_mut().remove(&process).is_none() { + return; + } + // SAFETY: caller keeps process live during deregistration. + unsafe { (*process).completion_queue = std::ptr::null() }; + self.retry.borrow_mut().remove(&process); + #[cfg(target_os = "macos")] + { + let change = libc::kevent64_s { + // SAFETY: caller contract. + ident: unsafe { (*process).pid } as u64, + filter: libc::EVFILT_PROC, + flags: libc::EV_DELETE, + fflags: 0, + data: 0, + udata: 0, + ext: [0; 2], + }; + // SAFETY: live queue and stack changelist. EV_ONESHOT may already + // have deleted the registration; no live pointer remains afterward. + retry_interrupted(|| unsafe { + libc::kevent64( + self.fd.get().unwrap().native(), + &raw const change, + 1, + std::ptr::null_mut(), + 0, + 0, + std::ptr::null(), + ) + }); + } + #[cfg(any(target_os = "linux", target_os = "android"))] + // SAFETY: process is live and still owns its pidfd at this point. + retry_interrupted(|| unsafe { + libc::epoll_ctl( + self.fd.get().unwrap().native(), + libc::EPOLL_CTL_DEL, + (*process).pidfd, + std::ptr::null_mut(), + ) + }); + if self.watched.borrow().is_empty() { + if let Some(fd) = self.fd.take() { + close_owned_descriptor(fd); + } + } + } + + pub(crate) fn retry_reap(&self, process: *mut Process) { + debug_assert!(self.watched.borrow().contains(&process)); + self.retry.borrow_mut().insert(process, ()); + } + + pub fn capture(&self) -> Vec> { + let registered = self.watched.borrow().len(); + if registered == 0 { + return Vec::new(); + } + let mut ready = Vec::new(); + { + let mut retry = self.retry.borrow_mut(); + for &process in retry.keys() { + // SAFETY: retry contains only registered owners; no callbacks run + // during capture, and each owner still holds its primary watch ref. + unsafe { + (*process).ref_(); + ready.push(RefPtr::from_raw(process)); + } + } + retry.clear(); + } + #[cfg(target_os = "macos")] + let mut events = [libc::kevent64_s { + ident: 0, + filter: 0, + flags: 0, + fflags: 0, + data: 0, + udata: 0, + ext: [0; 2], + }; 64]; + #[cfg(any(target_os = "linux", target_os = "android"))] + let mut events = [libc::epoll_event { events: 0, u64: 0 }; 64]; + // One-shot notifications cannot repeat within this capture. Collect the + // whole ready set, bounded by registrations, without scanning live PIDs. + while ready.len() < registered { + let capacity = events.len().min(registered - ready.len()) as i32; + #[cfg(target_os = "macos")] + // SAFETY: live queue; output capacity fits events; zero timeout. + let count = unsafe { + let timeout = libc::timespec { + tv_sec: 0, + tv_nsec: 0, + }; + libc::kevent64( + self.fd.get().unwrap().native(), + std::ptr::null(), + 0, + events.as_mut_ptr(), + capacity, + 0, + &raw const timeout, + ) + }; + #[cfg(any(target_os = "linux", target_os = "android"))] + // SAFETY: live queue and output buffer; zero timeout. + let count = unsafe { + libc::epoll_wait( + self.fd.get().unwrap().native(), + events.as_mut_ptr(), + capacity, + 0, + ) + }; + if count < 0 && bun_sys::get_errno(count) == bun_sys::E::EINTR { + continue; + } + if count <= 0 { + break; + } + for event in &events[..count as usize] { + #[cfg(target_os = "macos")] + let process = event.udata as *mut Process; + #[cfg(any(target_os = "linux", target_os = "android"))] + let process = event.u64 as *mut Process; + // SAFETY: registration retains the primary watch ref. Retain + // this snapshot before permitting tasks to detach an owner. + unsafe { + (*process).ref_(); + ready.push(RefPtr::from_raw(process)); + } + } + if count < capacity { + break; + } + } + ready + } +} + +impl Drop for CompletionQueue { + fn drop(&mut self) { + for &process in self.watched.get_mut().keys() { + // SAFETY: every registered owner holds its primary watch ref; + // destruction runs no callbacks and invalidates every backref. + unsafe { (*process).completion_queue = std::ptr::null() }; + } + if let Some(fd) = self.fd.take() { + close_owned_descriptor(fd); + } + } +} + +fn open_descriptor() -> bun_sys::Result { + #[cfg(target_os = "macos")] + // SAFETY: kqueue has no pointer arguments. + let fd = retry_interrupted(|| unsafe { libc::kqueue() }); + #[cfg(any(target_os = "linux", target_os = "android"))] + // SAFETY: epoll_create1 has no pointer arguments. + let fd = retry_interrupted(|| unsafe { libc::epoll_create1(libc::EPOLL_CLOEXEC) }); + if fd < 0 { + return Err(last_error()); + } + let mut fd = Fd::from_native(fd); + if fd.stdio_tag().is_some() { + let moved = loop { + match bun_sys::dup_at_least(fd, 3) { + Err(error) if error.get_errno() == bun_sys::E::EINTR => continue, + result => break result, + } + }; + close_owned_descriptor(fd); + fd = moved?; + } + #[cfg(target_os = "macos")] + if let Err(error) = bun_sys::set_close_on_exec(fd) { + close_owned_descriptor(fd); + return Err(error); + } + Ok(fd) +} + +fn close_owned_descriptor(fd: Fd) { + // A private queue may occupy 0–2 after the application closes stdio. + let error = fd.close_allowing_standard_io(None); + debug_assert!(error.is_none()); +} + +fn last_error() -> bun_sys::Error { + #[cfg(target_os = "macos")] + let tag = bun_sys::Tag::kqueue; + #[cfg(any(target_os = "linux", target_os = "android"))] + let tag = bun_sys::Tag::epoll_ctl; + bun_sys::Error::from_code(bun_sys::get_errno(-1i32), tag) +} + +fn retry_interrupted(mut operation: impl FnMut() -> i32) -> i32 { + loop { + let result = operation(); + if result >= 0 || bun_sys::get_errno(result) != bun_sys::E::EINTR { + return result; + } + } +} diff --git a/src/spawn/lib.rs b/src/spawn/lib.rs index 70eda17189d9..f548a0c2d580 100644 --- a/src/spawn/lib.rs +++ b/src/spawn/lib.rs @@ -24,6 +24,9 @@ pub mod ctrl_c; #[path = "process.rs"] pub mod process; +#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] +pub mod completion; + /// Generic `StaticPipeWriter

`. #[path = "static_pipe_writer.rs"] pub mod static_pipe_writer; diff --git a/src/spawn/process.rs b/src/spawn/process.rs index b9cf0a9b1af4..abff0e467580 100644 --- a/src/spawn/process.rs +++ b/src/spawn/process.rs @@ -118,6 +118,8 @@ fn call_exit_handler( #[derive(bun_ptr::ThreadSafeRefCounted)] pub struct Process { pub pid: PidT, + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + pub(crate) completion_queue: *const crate::completion::CompletionQueue, #[cfg(any(target_os = "linux", target_os = "android"))] pub(crate) pidfd: PidFdType, pub status: Status, @@ -135,6 +137,11 @@ impl Drop for Process { /// The allocation itself is freed by the `heap::take` in `destructor` /// above; this `Drop` body covers the `poller.deinit()` call. fn drop(&mut self) { + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + if !self.completion_queue.is_null() { + // SAFETY: queue teardown clears the backref; self is live until Drop returns. + unsafe { (*self.completion_queue).remove(self) }; + } self.poller.deinit(); } } @@ -283,6 +290,8 @@ impl Process { bun_core::heap::into_raw(Box::new(Process { ref_count: bun_ptr::ThreadSafeRefCount::init(), pid: posix.pid, + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + completion_queue: core::ptr::null(), #[cfg(any(target_os = "linux", target_os = "android"))] pidfd: posix.pidfd.unwrap_or(0), js_poster: event_loop.js_poster(), @@ -323,6 +332,43 @@ impl Process { let _ = sync_; } + /// A retained process from a readiness snapshot, independent of the + /// primary poll's dispatch ref. A nested task may have reaped it already. + /// + /// # Safety + /// The caller retains the snapshot ref on the process's owning thread. + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + pub unsafe fn publish_captured_completion(this: *mut Self) { + // SAFETY: the caller's snapshot keeps the owner live on this thread. + if unsafe { (*this).has_exited() || !matches!(&(*this).poller, Poller::Fd(_)) } { + return; + } + let mut rusage = rusage_zeroed(); + let result = posix_spawn::wait4( + // SAFETY: the retained owner still identifies this child. + unsafe { (*this).pid }, + libc::WNOHANG as u32, + Some(&mut rusage), + ); + // SAFETY: wait4 runs no JS callbacks; the snapshot still retains this owner. + let Some(status) = Status::from(unsafe { (*this).pid }, &result) else { + // A Darwin exit notification may precede a reapable status. + // Keep this observed owner dirty after consuming the one-shot. + // SAFETY: the snapshot retains the process; teardown clears this backref. + let queue = unsafe { (*this).completion_queue }; + if !queue.is_null() { + // SAFETY: a non-null backref denotes this owner's live registration. + unsafe { (*queue).retry_reap(this) }; + } + return; + }; + // SAFETY: the snapshot retains the owner across reentrant exit callbacks. + unsafe { (*this).on_exit(status, &rusage) }; + // SAFETY: on_exit detached the undispatched primary watch. Its ref is ours + // to release; the caller still retains the snapshot ref. + unsafe { Self::deref(this) }; + } + /// # Safety /// `this` carries the +1 ref taken when the waiter-thread task was queued. /// `RefPtr::from_raw` releases it on return — which may free `this` — so @@ -450,6 +496,40 @@ impl Process { } { Ok(()) => { self.ref_(); + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + { + unsafe extern "Rust" { + fn __bun_watch_process_completion( + process: *mut Process, + event_loop: EventLoopHandle, + ) -> bun_sys::Result<()>; + } + // Only Subprocess owns the pipe drains paired with this + // publication path; shell/install owners keep their joins. + if self + .exit_handler + .is_some_and(|handler| handler.kind == ProcessExitKind::Subprocess) + { + // SAFETY: the primary watch just retained self and its owning loop. + let completion = + unsafe { __bun_watch_process_completion(self, self.event_loop) }; + if let Err(error) = completion { + // The primary watch is already registered. A + // secondary snapshot must not make a valid spawn + // fail when descriptor/watch resources are full. + if !matches!( + error.get_errno(), + bun_sys::E::EMFILE + | bun_sys::E::ENFILE + | bun_sys::E::ENOMEM + | bun_sys::E::ENOSPC + ) { + self.close(); + return Err(error); + } + } + } + } Ok(()) } Err(err) => { @@ -578,6 +658,11 @@ impl Process { } pub fn close(&mut self) { + #[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))] + if !self.completion_queue.is_null() { + // SAFETY: self is live and queue teardown clears the shared backref. + unsafe { (*self.completion_queue).remove(self) }; + } #[cfg(unix)] { let mut stranded_watch_ref = false; diff --git a/src/uws_sys/us_socket_t.rs b/src/uws_sys/us_socket_t.rs index 09f8e8d86611..e1b44d560a83 100644 --- a/src/uws_sys/us_socket_t.rs +++ b/src/uws_sys/us_socket_t.rs @@ -74,6 +74,18 @@ pub struct UsIoVec { } impl us_socket_t { + /// Read an owned pipe without dispatching the loop's other ready handles. + /// + /// # Safety + /// `socket` is a live socket on this thread; its owner remains live across callbacks. + pub unsafe fn drain_readable(socket: *mut Self) { + unsafe extern "C" { + fn us_socket_drain_readable(socket: *mut us_socket_t); + } + // SAFETY: caller keeps the socket and its owner live across callbacks. + unsafe { us_socket_drain_readable(socket) }; + } + pub(crate) fn pause(&mut self) { bun_core::scoped_log!(uws, "us_socket_pause({:p})", self); c::us_socket_pause(self); diff --git a/test/js/bun/spawn/spawn-pipe-stale-fd-unregister.test.ts b/test/js/bun/spawn/spawn-pipe-stale-fd-unregister.test.ts index f2c1c3e72b86..42baf2a913c6 100644 --- a/test/js/bun/spawn/spawn-pipe-stale-fd-unregister.test.ts +++ b/test/js/bun/spawn/spawn-pipe-stale-fd-unregister.test.ts @@ -35,7 +35,8 @@ function pipeFds() { if (!Number.isInteger(fd)) continue; try { const st = fstatSync(fd); - if (st.isFIFO() || st.isSocket()) out.set(fd, st.ino); + // stdout is a socketpair; Darwin kqueue descriptors also report FIFO. + if (st.isSocket()) out.set(fd, st.ino); } catch { // closed between readdir and fstat (e.g. readdir's own dir fd) } diff --git a/test/js/node/child_process/child_process.test.ts b/test/js/node/child_process/child_process.test.ts index 186538d42892..c427fa37d9be 100644 --- a/test/js/node/child_process/child_process.test.ts +++ b/test/js/node/child_process/child_process.test.ts @@ -34,6 +34,27 @@ import path from "path"; const debug = process.env.DEBUG ? console.log : () => {}; const originalProcessEnv = process.env; + +it.skipIf(!isPosix)("publishes exited children and ready pipe EOF after a blocked I/O callback", async () => { + const nativeNode = nodeExe(); + expect(nativeNode).not.toBeNull(); + using dir = tempDir("subprocess-retirement", {}); + await using child = Bun.spawn({ + cmd: [bunExe(), path.join(import.meta.dir, "fixtures", "retirement-io.cjs")], + env: { ...bunEnv, PROBE_NODE: nativeNode!, TMPDIR: String(dir) }, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, code] = await Promise.all([child.stdout.text(), child.stderr.text(), child.exited]); + expect(stderr).toBe(""); + expect(JSON.parse(stdout)).toMatchObject({ + passed: true, + errors: [], + observerSnapshot: { exit: true, fd3End: true, stdoutEnd: true, stderrEnd: true }, + }); + expect(code).toBe(0); +}); + beforeEach(() => { process.env = { ...bunEnv }; // Github actions might filter these out diff --git a/test/js/node/child_process/fixtures/retirement-io.cjs b/test/js/node/child_process/fixtures/retirement-io.cjs new file mode 100644 index 000000000000..f8771df98cbc --- /dev/null +++ b/test/js/node/child_process/fixtures/retirement-io.cjs @@ -0,0 +1,244 @@ +// A blocked control-I/O callback must not hide an exited child or its ready pipe EOF. +const fs = require("node:fs"); +const path = require("node:path"); +const os = require("node:os"); +const { spawn } = require("node:child_process"); +const { Socket } = require("node:net"); +const [role, dir, node] = process.argv.slice(2); +const atomic = (name, data) => { + const file = path.join(dir, name); + fs.writeFileSync(file + ".tmp", JSON.stringify(data)); + fs.renameSync(file + ".tmp", file); +}; +const lines = (stream, onLine) => { + let pending = ""; + stream.on("data", data => { + pending += data.toString(); + let end; + while ((end = pending.indexOf("\n")) >= 0) { + const line = pending.slice(0, end); + pending = pending.slice(end + 1); + onLine(JSON.parse(line)); + } + }); +}; +if (role === "--anchor") { + setInterval(() => {}, 1000); +} else if (role === "--observer") { + const timer = setInterval(() => { + if (!fs.existsSync(path.join(dir, "observe.json"))) return; + const input = JSON.parse(fs.readFileSync(path.join(dir, "observe.json"), "utf8")); + try { + process.kill(-input.anchorPid, 0); + } catch (error) { + if (process.platform === "darwin" && error.code === "EPERM") return; + if (error.code !== "ESRCH") throw error; + atomic("observed.json", { ...input, retiredAt: Date.now() }); + clearInterval(timer); + process.exit(0); + } + }, 2); + process.stdout.write("ready\n"); +} else if (role === "--relay") { + const anchor = spawn(node, [__filename, "--anchor"], { detached: true, stdio: "ignore" }); + const control = new Socket({ fd: 3, readable: true, writable: true }); + let acknowledged = false; + let cancelling = false; + const cancel = () => { + if (cancelling) return; + cancelling = true; + if (!anchor.kill("SIGKILL")) process.exit(0); + }; + process.on("SIGTERM", cancel); + control.once("error", cancel); + control.once("end", cancel); + lines(control, message => { + if (message.type === "close") { + fs.writeSync(1, "stdout-final\n"); + fs.writeSync(2, "stderr-final\n"); + control.write(JSON.stringify({ type: "closing", result: 0 }) + "\n"); + } else if (message.type === "ack") { + acknowledged = true; + atomic("ack-seen.json", { at: Date.now() }); + anchor.kill("SIGTERM"); + } + }); + anchor.once("spawn", () => control.write(JSON.stringify({ type: "ready", anchorPid: anchor.pid }) + "\n")); + anchor.once("error", error => { + throw error; + }); + anchor.once("exit", () => { + if (cancelling) process.exit(0); + if (!acknowledged) throw new Error("anchor exited before acknowledgement"); + atomic("anchor-reaped.json", { at: Date.now() }); + // There are deliberately no fd3 writes after the acknowledgement. + process.exit(0); + }); +} else { + if (process.platform === "win32") throw new Error("POSIX process-group oracle"); + const nativeNode = process.env.PROBE_NODE; + if (!nativeNode) throw new Error("PROBE_NODE must name the native fixture runtime"); + const root = fs.mkdtempSync(path.join(os.tmpdir(), "retirement-io-")); + const started = Date.now(); + const events = []; + const errors = []; + const counts = {}; + const mark = (name, extra = {}) => { + counts[name] = (counts[name] || 0) + 1; + events.push({ name, ms: Date.now() - started, ...extra }); + }; + const record = (name, data) => { + const file = path.join(root, name); + fs.writeFileSync(file + ".tmp", JSON.stringify(data)); + fs.renameSync(file + ".tmp", file); + }; + const read = name => + fs.existsSync(path.join(root, name)) ? JSON.parse(fs.readFileSync(path.join(root, name), "utf8")) : null; + let relay; + let anchorPid; + let deadline; + let ackAt; + let resumedAt; + let observerSnapshot; + let finished = false; + let stdout = ""; + let stderr = ""; + const snapshot = () => ({ + exit: !!counts["relay-exit"], + close: !!counts["relay-close"], + fd3End: !!counts["fd3-end"], + fd3Close: !!counts["fd3-close"], + stdoutEnd: !!counts["stdout-end"], + stderrEnd: !!counts["stderr-end"], + }); + const observer = spawn(nativeNode, [__filename, "--observer", root], { stdio: ["ignore", "pipe", "pipe"] }); + observer.stderr.on("data", data => errors.push("observer: " + data)); + observer.on("error", error => errors.push(String(error))); + observer.on("exit", code => { + mark("observer-exit", { code }); + maybeFinish(); + }); + const watchdog = setTimeout(() => { + errors.push("probe watchdog"); + finish(); + }, 5000); + observer.stdout.once("data", () => { + mark("observer-ready"); + relay = spawn(nativeNode, [__filename, "--relay", root, nativeNode], { stdio: ["ignore", "pipe", "pipe", "pipe"] }); + relay.on("error", error => errors.push(String(error))); + relay.on("exit", (code, signal) => { + mark("relay-exit", { code, signal }); + maybeFinish(); + }); + relay.on("close", (code, signal) => { + mark("relay-close", { code, signal }); + maybeFinish(); + }); + for (const [name, stream] of [ + ["stdout", relay.stdout], + ["stderr", relay.stderr], + ["fd3", relay.stdio[3]], + ]) { + stream.on("end", () => mark(name + "-end")); + stream.on("close", () => { + mark(name + "-close"); + maybeFinish(); + }); + stream.on("error", error => errors.push(name + ": " + error)); + } + relay.stdout.on("data", data => { + stdout += data; + mark("stdout-data"); + }); + relay.stderr.on("data", data => { + stderr += data; + mark("stderr-data"); + }); + lines(relay.stdio[3], message => { + mark("fd3-" + message.type); + if (message.type === "ready") { + anchorPid = message.anchorPid; + deadline = Date.now() + 100; + setTimeout(() => { + mark("deadline-timer", snapshot()); + setImmediate(() => { + observerSnapshot = snapshot(); + mark("deadline-immediate", observerSnapshot); + maybeFinish(); + }); + }, 100); + relay.stdio[3].write(JSON.stringify({ type: "close" }) + "\n"); + } else if (message.type === "closing") { + relay.stdio[3].write(JSON.stringify({ type: "ack" }) + "\n", error => { + if (error) errors.push(String(error)); + ackAt = Date.now(); + record("observe.json", { anchorPid, ackAt }); + mark("ack-flushed"); + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 300); + resumedAt = Date.now(); + mark("host-resumed", { observation: read("observed.json"), anchorReaped: read("anchor-reaped.json") }); + }); + } + }); + }); + function maybeFinish() { + if (observerSnapshot && counts["relay-close"] && counts["fd3-close"] && counts["observer-exit"]) finish(); + } + function finish() { + if (finished) return; + finished = true; + clearTimeout(watchdog); + const observed = read("observed.json"); + const ackSeen = read("ack-seen.json"); + const reaped = read("anchor-reaped.json"); + if (!(observed && ackAt <= observed.retiredAt && observed.retiredAt < deadline && deadline < resumedAt)) + errors.push("independent extinction deadline precondition failed"); + if (!(reaped && reaped.at < deadline)) errors.push("anchor not reaped before deadline"); + if (stdout !== "stdout-final\n" || stderr !== "stderr-final\n") errors.push("final output differs"); + for (const name of [ + "ack-flushed", + "host-resumed", + "relay-exit", + "relay-close", + "fd3-end", + "fd3-close", + "deadline-immediate", + ]) { + if (counts[name] !== 1) errors.push(name + " count " + (counts[name] || 0)); + } + // This exact I/O-entry fixture's Node oracle exposes native exit and EOF + // before the observer. Socket 'close' belongs to the later closing phase. + if ( + !observerSnapshot?.exit || + !observerSnapshot?.fd3End || + !observerSnapshot?.stdoutEnd || + !observerSnapshot?.stderrEnd + ) + errors.push("exit or readable EOF missing at deadline observer"); + const exit = events.find(event => event.name === "relay-exit"); + if (exit?.code !== 0 || exit?.signal !== null) errors.push("relay exit differs"); + console.log( + JSON.stringify({ + runtime: process.version, + bun: process.versions.bun || null, + revision: typeof Bun === "undefined" ? null : Bun.revision, + platform: process.platform, + started, + ackAt, + deadline, + resumedAt, + observed, + ackSeen, + reaped, + observerSnapshot, + events, + errors, + passed: errors.length === 0, + }), + ); + observer.kill("SIGKILL"); + if (relay) relay.kill("SIGTERM"); + fs.rmSync(root, { recursive: true, force: true }); + process.exit(errors.length ? 1 : 0); + } +} diff --git a/test/js/node/http2/h2-conformance.test.ts b/test/js/node/http2/h2-conformance.test.ts index b121dbb06374..5820d868c317 100644 --- a/test/js/node/http2/h2-conformance.test.ts +++ b/test/js/node/http2/h2-conformance.test.ts @@ -164,7 +164,7 @@ beforeAll(async () => { stream.respond({ ":status": 200 }); stream.end("ok"); }); - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); port = (server.address() as net.AddressInfo).port; }); @@ -376,7 +376,7 @@ describe("CONTINUATION (checklist §3,§7)", () => { stream.respond({ ":status": 200 }); stream.end("ok"); }); - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); const c = await RawH2.connect((server.address() as net.AddressInfo).port); try { @@ -953,7 +953,7 @@ describe("request header and body framing (RFC 9113 §8.1)", () => { }); stream.resume(); }); - deferredServer.listen(0); + deferredServer.listen(0, "127.0.0.1"); await once(deferredServer, "listening"); deferredPort = (deferredServer.address() as net.AddressInfo).port; }); @@ -1333,7 +1333,7 @@ describe("inbound stream lifecycle", () => { stream.resume(); if (refs.length === total) allOpen.resolve(); }); - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); const c = await RawH2.connect((server.address() as net.AddressInfo).port); try { @@ -1516,9 +1516,9 @@ describe("inbound stream lifecycle", () => { }); stream.end("body"); }); - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); - const client = http2.connect(`http://localhost:${(server.address() as net.AddressInfo).port}`); + const client = http2.connect(`http://127.0.0.1:${(server.address() as net.AddressInfo).port}`); client.on("error", e => trailers.reject(e)); try { const req = client.request({ ":path": "/" }); @@ -1541,7 +1541,7 @@ describe("inbound stream lifecycle", () => { stream.respond({ ":status": 200 }); stream.write(Buffer.alloc(1 << 22, "a")); }); - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); const c = await RawH2.connect((server.address() as net.AddressInfo).port); try { @@ -1577,7 +1577,7 @@ describe("inbound stream lifecycle", () => { stream.respond({ ":status": 200 }); stream.write(Buffer.alloc(1 << 22, "a")); }); - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); const c = await RawH2.connect((server.address() as net.AddressInfo).port); try { @@ -1614,7 +1614,7 @@ describe("inbound stream lifecycle", () => { stream.end("ok"); } }); - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); const c = await RawH2.connect((server.address() as net.AddressInfo).port); c.sendPreface(); @@ -1740,7 +1740,7 @@ describe("stream release after a queued END_STREAM", () => { const STALLED_BODY = Buffer.alloc(256 * 1024, "s"); async function listen(server: http2.Http2Server): Promise { - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); return `http://127.0.0.1:${(server.address() as net.AddressInfo).port}`; } @@ -1958,7 +1958,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou } async function withClient(server: http2.Http2Server, body: (c: RawH2) => Promise): Promise { - server.listen(0); + server.listen(0, "127.0.0.1"); await once(server, "listening"); const c = await RawH2.connect((server.address() as net.AddressInfo).port); try {