Skip to content
Open
6 changes: 6 additions & 0 deletions src/js/internal/cluster/RoundRobinHandle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,12 @@ export default class RoundRobinHandle {
return this.all.has(worker.id);
}

// The descriptor that the listener of this handle holds, -1 before it listens and after remove() closed it.
get fd() {
const fd = this.server?._handle?.fd;
return typeof fd === "number" ? fd : -1;
}

// With the channel still up the unacked newconn is settled by its ack; once it is gone, a crashed worker's goes to another worker and a disconnected worker's (already settled by it) is dropped.
remove(worker, channelGone = false) {
if (channelGone) {
Expand Down
10 changes: 8 additions & 2 deletions src/js/internal/cluster/SharedHandle.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
const clusterRawBind = $newRustFunction("node_cluster_binding.rs", "clusterRawBind", 4);
const closeRawHandle = $newRustFunction("node_cluster_binding.rs", "clusterCloseHandle", 1);
const validateFd = $newRustFunction("node_cluster_binding.rs", "clusterValidateFd", 1);
const validateFd = $newRustFunction("node_cluster_binding.rs", "clusterValidateFd", 2);

export default class SharedHandle {
key;
Expand All @@ -19,7 +19,7 @@ export default class SharedHandle {
this.sharedOnly = sharedOnly === true;

if (typeof fd === "number" && fd >= 0) {
const err = validateFd(fd);
const err = validateFd(fd, addressType === "udp4" || addressType === "udp6");
if (err !== 0) {
this.errno = err;
} else {
Expand Down Expand Up @@ -47,6 +47,12 @@ export default class SharedHandle {
return this.workers.has(worker.id);
}

// The descriptor that remove() closes: the number a worker named, or the socket that clusterRawBind made.
get fd() {
const handle = this.handle;
return handle ? handle.fd : -1;
}

remove(worker) {
const workers = this.workers;
if (!workers.has(worker.id)) return false;
Expand Down
28 changes: 26 additions & 2 deletions src/js/internal/cluster/primary.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,13 @@ const { kHandle } = require("internal/shared");

const sendHelper = $newRustFunction("node_cluster_binding.rs", "sendHelperPrimary", 4);
const onInternalMessage = $newRustFunction("node_cluster_binding.rs", "onInternalMessagePrimary", 3);
const { UV_EINVAL, UV_ENOBUFS } = process.binding("uv");
const { UV_EEXIST, UV_EINVAL, UV_ENOBUFS } = process.binding("uv");

let child_process;
let RoundRobinHandle;
let SharedHandle;
let validateFd;
let isFdOfDgramSocket;

const ArrayPrototypeSlice = Array.prototype.slice;
const ObjectValues = Object.values;
Expand Down Expand Up @@ -45,7 +47,7 @@ const SCHED_RR = 2;

export default cluster;

const handles = new Map();
const handles = new Map<string, any>();
cluster.isWorker = false;
cluster.isMaster = true; // Deprecated alias. Must be same as isPrimary.
cluster.isPrimary = true;
Expand Down Expand Up @@ -240,6 +242,21 @@ function exitedAfterDisconnect(worker, message) {
send(worker, { ack: message.seq });
}

// A handle closes its descriptor when its last worker leaves, so a second handle for a held descriptor closes it under the holder. Looked up by number: `udp4` and `udp6` are two keys for one descriptor.
function errnoOfFdQuery(message) {
const fd = message.fd;
const isUdp = message.addressType === "udp4" || message.addressType === "udp6";
let held = isUdp && (isFdOfDgramSocket ??= $newRustFunction("udp_socket.rs", "jsDgramIsFdAdopted", 1))(fd);
if (!held) {
handles.$forEach(handle => {
if (handle.fd === fd) held = true;
});
Comment on lines +246 to +253

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟣 pre-existing, not blocking: pre-existing: a primary that itself listens on an inherited descriptor and whose worker names the same number still loses that listener and closes the number twice, exactly as on the base branch. The holder lookup in errnoOfFdQuery at primary.ts:246-253 only sees the dgram registry and the entries of handles; a net.Server or Bun.listen listener of the primary is not a holder it can see, so a SharedHandle adopts the descriptor and closes it when the worker leaves. Fix: treat a descriptor that a live listener of the primary holds as held too, returning EEXIST like the other holders; the PR description lists this under "Not covered".
A small fix can ride a push you are already making; otherwise a short reply is enough.

Why this was flagged

Trigger: in the primary, net.createServer().listen({ fd: 3 }), then a worker calls server.listen({ fd: 3 }) under SCHED_NONE (or tls under SCHED_RR). queryServer reaches errnoOfFdQuery at src/js/internal/cluster/primary.ts:246; handles.$forEach at primary.ts:251-253 finds no handle because the primary's own net.Server is not in handles. errnoOfFdQuery returns 0, and primary.ts:325 constructs a SharedHandle; SharedHandle.ts:22 validateFd(3, false) passes and SharedHandle.ts:26 records { fd: 3 }. When the worker closes its server, close() at primary.ts:376-383 calls remove(), and SharedHandle.ts:67 closeRawHandle(3) closes the descriptor that the primary's own listener still polls; the primary's server silently stops accepting, and its later server.close() closes number 3 a second time. The base branch behaves the same; the diff adds the holder lookup but it cannot see this holder. No safeguard covers it: validateFd checks kind and connectedness only.

Verification: Trigger: the primary does net.createServer().listen({ fd: N }) and a worker calls listen({ fd: N }) under SCHED_NONE. That listener is never inserted into handles, so the holder lookup in errnoOfFdQuery (src/js/internal/cluster/primary.ts:249-253) finds nothing and returns 0. SharedHandle.remove (SharedHandle.ts:67) calls closeRawHandle(fd), closing the descriptor under the primary's live listener.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This case is real. It is the first entry under "Not covered" in the PR body, and I left it out of this PR on purpose.

I measured it with release builds. The primary listens on descriptor 3 with its own net.Server. One worker under SCHED_NONE calls listen({ fd: 3 }) and closes its server when it listens. Then a client connects.

runtime answer to the worker after that
node v26.3.0 bind EEXIST descriptor 3 is open, the primary serves the client
bun, main listening descriptor 3 of the primary is EBADF, the client gets no answer
bun, this PR listening the same as main

Node has this answer from libuv: uv_tcp_open returns UV_EEXIST for a descriptor that its loop watches already (tcp.c). The cluster code of Bun cannot see a listener that it did not make. A walk over the listen sockets of the loop can see it, but that is native code in packages/bun-usockets. This PR changes three builtin JS files and no native code, so that lookup is a change of its own.

}
if (!held) return 0;
// The kind comes before the holder, as in node. Not as in node v26.3.0: EEXIST also under SCHED_NONE and for udp4 then udp6, where its worker dies on an assertion (nodejs/node#64869) or node serves and closes the number two times.
return (validateFd ??= $newRustFunction("node_cluster_binding.rs", "clusterValidateFd", 2))(fd, isUdp) || UV_EEXIST;
}

function queryServer(worker, message) {
// Stop processing if worker already disconnecting
if (worker.exitedAfterDisconnect) return;
Expand Down Expand Up @@ -294,6 +311,13 @@ function queryServer(worker, message) {
worker.emit("error", error);
return;
}
if (typeof message.fd === "number" && message.fd >= 0) {
const errno = errnoOfFdQuery(message);
if (errno !== 0) {
send(worker, { errno, key, ack: message.seq, data: cachedHandle ? cachedHandle.data : message.data }, null);
return;
}
}
if (
schedulingPolicy !== SCHED_RR ||
message.sharedOnly === true ||
Expand Down
31 changes: 28 additions & 3 deletions src/runtime/node/node_cluster_binding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -731,6 +731,7 @@ pub(crate) fn cluster_raw_bind(global: &JSGlobalObject, frame: &CallFrame) -> Js
}
}

/// `(fd, wantDgram)` → 0 when `fd` can serve a query of that kind (datagram, else stream), or a negative errno.
#[bun_jsc::host_fn]
pub(crate) fn cluster_validate_fd(global: &JSGlobalObject, frame: &CallFrame) -> JsResult<JSValue> {
let _ = global;
Expand All @@ -741,12 +742,21 @@ pub(crate) fn cluster_validate_fd(global: &JSGlobalObject, frame: &CallFrame) ->
#[cfg(not(windows))]
{
let fd = value.to_int32();
// `fd: 3.5` is not descriptor 3: https://github.com/nodejs/node/blob/v26.3.0/lib/net.js#L1904-L1911
if f64::from(fd) != value.as_number() {
return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL));
}
if fd < 0 {
return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::BADF));
}
let wanted = if frame.argument(1).to_boolean() {
libc::SOCK_DGRAM
} else {
libc::SOCK_STREAM
};
let mut ty: libc::c_int = 0;
let mut len = core::mem::size_of::<libc::c_int>() as libc::socklen_t;
// SAFETY: plain getsockopt on a caller-supplied fd; out-params are
// SAFETY: plain getsockopt on a caller-supplied fd; out-params are live locals, and `len` is the size of `ty`.
let rc = unsafe {
libc::getsockopt(
fd,
Expand All @@ -756,8 +766,8 @@ pub(crate) fn cluster_validate_fd(global: &JSGlobalObject, frame: &CallFrame) ->
&raw mut len,
)
};
// node's createServerHandle: EINVAL for anything that cannot listen (e.g. a connected stdio socketpair), fd left untouched.
if rc != 0 || (ty != libc::SOCK_STREAM && ty != libc::SOCK_DGRAM) {
// A descriptor of another kind is EINVAL before node opens it: https://github.com/nodejs/node/blob/v26.3.0/lib/net.js#L1904-L1911 and https://github.com/nodejs/node/blob/v26.3.0/lib/internal/dgram.js#L67-L70
if rc != 0 || ty != wanted {
return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL));
}
if ty == libc::SOCK_STREAM {
Expand All @@ -771,6 +781,21 @@ pub(crate) fn cluster_validate_fd(global: &JSGlobalObject, frame: &CallFrame) ->
if connected {
return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL));
}
} else {
// SAFETY: sockaddr_storage is plain data; getsockname only writes within `name_len`.
let family = unsafe {
let mut name: libc::sockaddr_storage = bun_core::ffi::zeroed_unchecked();
let mut name_len =
core::mem::size_of::<libc::sockaddr_storage>() as libc::socklen_t;
if libc::getsockname(fd, (&raw mut name).cast(), &raw mut name_len) != 0 {
return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL));
}
libc::c_int::from(name.ss_family)
};
// https://github.com/nodejs/node/blob/v26.3.0/deps/uv/src/unix/tty.c#L458-L460
if family != libc::AF_INET && family != libc::AF_INET6 {
return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL));
}
}
Ok(JSValue::js_number_from_int32(0))
}
Expand Down
Loading
Loading