Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions docs/runtime/nodejs-compat.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Expand Down
1 change: 1 addition & 0 deletions packages/bun-usockets/src/internal/internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
1 change: 1 addition & 0 deletions packages/bun-usockets/src/libusockets.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
8 changes: 8 additions & 0 deletions packages/bun-usockets/src/loop.c
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
12 changes: 12 additions & 0 deletions packages/bun-usockets/src/socket.c
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down
27 changes: 23 additions & 4 deletions src/js/node/child_process.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
}
}
}

Expand Down Expand Up @@ -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;
}
Expand All @@ -1382,6 +1400,7 @@ class ChildProcess extends EventEmitter {
#stderr;
#stdioObject;
#stdioOptions;
#extraStdioHandles: unknown[] = [];

#createStdioObject() {
const opts = this.#stdioOptions;
Expand Down
51 changes: 51 additions & 0 deletions src/jsc/event_loop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::rc::Rc<bun_spawn::completion::CompletionQueue>>,
/// 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,
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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> {
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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::<EventLoop>();
// 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
Expand Down
40 changes: 40 additions & 0 deletions src/runtime/api/bun/subprocess.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
33 changes: 33 additions & 0 deletions src/runtime/node/node_net_binding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<JSValue> {
#[cfg(unix)]
{
let [value] = frame.arguments_as_array::<1>();
if let Some(socket) = value.as_::<TCPSocket>() {
// 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
Expand Down
3 changes: 3 additions & 0 deletions src/spawn/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading
Loading