From d4b09f3de3864fa89637456adda7ef6c6bdf896b Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 13 Aug 2026 02:58:20 +0000 Subject: [PATCH 1/5] node:http: listen through the cluster primary's shared handle in workers In a cluster worker, http.Server#listen() bound its own SO_REUSEPORT socket with Bun.serve(), so listen(0) gave every worker a different port, a unix path could only be served by one worker, bind errors had Bun.serve's shape and worker.disconnect() did not know about the server. The worker now asks the primary through cluster._getServer() with sharedOnly, like a TLS net.Server does: the primary binds once per key and ships the descriptor, and the worker's Bun.serve() accepts on it through a new listen_fd path (uws HttpContext::listen_fd -> us_socket_group_listen_fd, ServerConfig.listen_fd, only read for node:http servers). config.address still describes the bound address, so address()/url/errors are unchanged, and a server started on an inherited descriptor no longer unlinks the unix socket file the primary owns. close()/closeAllConnections() release the handle; the handle's owner is the server so disconnect() closes it. exclusive (and reusePort, which implies it) still binds in the worker, as does a worker without a channel, and Windows keeps binding per worker: several processes accepting on copies of one listening socket block each other in accept() there and Bun.serve cannot take round-robin handoffs. --- packages/bun-uws/src/App.h | 6 + packages/bun-uws/src/HttpContext.h | 21 ++ src/js/node/_http_server.ts | 188 +++++++++++---- src/runtime/server/ServerConfig.rs | 47 ++++ src/runtime/server/mod.rs | 50 +++- src/uws_sys/App.rs | 37 ++- src/uws_sys/libuwsockets.cpp | 25 ++ test/js/node/cluster.test.ts | 359 +++++++++++++++++++++++++++++ 8 files changed, 683 insertions(+), 50 deletions(-) diff --git a/packages/bun-uws/src/App.h b/packages/bun-uws/src/App.h index eb5ee2ba132d..c5ad83c270b1 100644 --- a/packages/bun-uws/src/App.h +++ b/packages/bun-uws/src/App.h @@ -756,6 +756,12 @@ struct TemplatedApp { return std::move(*this); } + /* already-bound descriptor, options, callback */ + TemplatedApp &&listen_fd(LIBUS_SOCKET_DESCRIPTOR fd, int options, MoveOnlyFunction &&handler) { + handler(httpContext ? trackListenSocket(httpContext->listen_fd(sslCtxOrNull(), fd, options)) : nullptr); + return std::move(*this); + } + void setOnSocketClosed(HttpContextData::OnSocketClosedCallback onClose) { httpContext->getSocketContextData()->onSocketClosed = onClose; } diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index 6f7cf7441cca..fec88c92ff07 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -35,6 +35,7 @@ #include #include #include +#include extern "C" void Bun__NodeHTTP__onReadsResumable(int ssl, struct us_socket_t *s); @@ -1070,6 +1071,26 @@ struct HttpContext { return socket; } + + /* Accept on a descriptor that is already bound (node:cluster's shared listen handle: + * the primary bound it, every worker listen()s and accepts on its own copy). The fd + * is owned by the returned listen socket; on failure the caller still owns it. */ + us_listen_socket_t *listen_fd(struct ssl_ctx_st *sslCtx, LIBUS_SOCKET_DESCRIPTOR fd, int options) { + int error = 0; + auto* socket = us_socket_group_listen_fd(&group, socketKind(), sslCtx, fd, 512, options | LIBUS_LISTEN_DEFER_ACCEPT, socketExtSize(), &error); + if (socket) { + // we dont depend on libuv ref for keeping it alive + us_socket_unref(&socket->s); + return socket; + } +#ifndef _WIN32 + /* Bun.serve's listen-failure path reads errno (see us_socket_group_listen). */ + if (error) { + errno = error; + } +#endif + return nullptr; + } }; } diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 89b6db00bacd..3e38e5af6c1c 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -1,7 +1,7 @@ // Hardcoded module "node:_http_server" const EventEmitter: typeof import("node:events").EventEmitter = require("node:events"); const { Stream } = require("node:stream"); -const { Socket: NetSocket } = require("node:net"); +const { Socket: NetSocket, isIP } = require("node:net"); const { _checkInvalidHeaderChar: checkInvalidHeaderChar, chunkExpression, @@ -18,7 +18,15 @@ const { validateFunction, validateOneOf, } = require("internal/validators"); -const { ConnResetException, hasObserver, startPerf, stopPerf, kInternalSendOptions } = require("internal/shared"); +const { + ConnResetException, + ExceptionWithHostPort, + hasObserver, + startPerf, + stopPerf, + kClusterOwner, + kInternalSendOptions, +} = require("internal/shared"); const kServerResponseStatistics = Symbol("ServerResponseStatistics"); const { isPrimary } = require("internal/cluster/isPrimary"); @@ -93,6 +101,10 @@ const { const kConnectionsCheckingInterval = Symbol("http.server.connectionsCheckingInterval"); const kTrackedConnections = Symbol("http.server.trackedConnections"); const kHttpAllowHalfOpen = Symbol("http.server.httpAllowHalfOpen"); +// Cluster worker state: the primary's shared listen handle this server accepts on, and the +// counter that tells a late primary reply that close() or a newer listen() superseded it. +const kClusterHandle = Symbol("http.server.clusterHandle"); +const kClusterListeningId = Symbol("http.server.clusterListeningId"); // node.http trace events ('http.server.request' b/e). The agent module is // only created on the first request, and emission is gated per-request on the @@ -504,6 +516,7 @@ Server.prototype.closeAllConnections = function () { clearInterval(this[kConnectionsCheckingInterval]); this.listening = false; + releaseClusterHandle(this); server.stop(true); }; @@ -525,6 +538,8 @@ Server.prototype.closeIdleConnections = function () { Server.prototype.close = function (optionalCallback?) { const server = this[serverSymbol]; + // Invalidates a cluster listen() whose primary reply has not arrived yet (see listenInCluster). + this[kClusterListeningId] = (this[kClusterListeningId] || 0) + 1; // Node.js's httpServerPreClose clears the connections-checking interval // even when the server was never listening. clearInterval(this[kConnectionsCheckingInterval]); @@ -537,6 +552,7 @@ Server.prototype.close = function (optionalCallback?) { this[serverSymbol] = undefined; if (typeof optionalCallback === "function") setCloseCallback(this, optionalCallback); this.listening = false; + releaseClusterHandle(this); server.closeIdleConnections(); server.stop(); return this; @@ -587,6 +603,8 @@ Server.prototype.listen = function () { const server = this; let port, host, onListen; let socketPath; + let exclusive = false; + let reusePort = false; let tls = this[tlsSymbol]; // This logic must align with: @@ -599,6 +617,9 @@ Server.prototype.listen = function () { port = arg0.port; host = arg0.host; socketPath = arg0.path; + // As in net.Server#listen, reusePort implies exclusive: the worker binds its own socket. + reusePort = arg0.reusePort === true; + exclusive = !!arg0.exclusive || reusePort; const otherTLS = arg0.tls; if (otherTLS && $isObject(otherTLS)) { @@ -644,45 +665,26 @@ Server.prototype.listen = function () { if (cluster === undefined) cluster = require("node:cluster"); - // const serverQuery = { - // // address: address, - // port: port, - // addressType: 4, - // // fd: fd, - // // flags, - // // backlog, - // // ...options, - // }; - // cluster._getServer(server, serverQuery, function listenOnPrimaryHandle(err, handle) { - // // err = checkBindError(err, port, handle); - // // if (err) { - // // throw new ExceptionWithHostPort(err, "bind", address, port); - // // } - // if (err) { - // throw err; - // } - // server[kRealListen](port, host, socketPath, onListen); - // }); - - server.once("listening", () => { - // No channel (NODE_UNIQUE_ID inherited by a plain child, or already disconnected): nothing to notify. - if (!process.connected) return; - cluster.worker.state = "listening"; - const address = server.address(); - const isObjectAddress = address !== null && typeof address === "object"; - const boundHost = host && isObjectAddress ? address : null; - const message = { - cmd: "NODE_CLUSTER", - act: "listening", - port: socketPath ? -1 : (isObjectAddress && address.port) || port, - data: null, - address: socketPath ?? (boundHost && boundHost.address) ?? null, - addressType: socketPath ? -1 : boundHost && boundHost.family === "IPv6" ? 6 : 4, - }; - process.send(message, undefined, kInternalSendOptions); - }); + // The worker binds its own socket when node would (exclusive, or reusePort which implies + // it), when there is no primary to ask (NODE_UNIQUE_ID inherited by a plain child, or + // already disconnected), when the arguments are nothing the primary could bind (Bun.serve + // reports them), and on Windows: processes accepting on copies of one listening socket + // block each other in accept() there, and the primary's alternative of handing over each + // accepted connection is something Bun.serve's native accept loop cannot take. + if ( + exclusive || + !process.connected || + process.platform === "win32" || + !isBindableByPrimary(port, host, socketPath) + ) { + notifyPrimaryWhenListening(server, port, host, socketPath); + // Only an exclusive listen gets to opt out of SO_REUSEPORT: the other cases still need + // it for the workers to end up on one port at all. + server[kRealListen](tls, port, host, socketPath, exclusive ? reusePort : true, onListen); + return this; + } - server[kRealListen](tls, port, host, socketPath, true, onListen); + listenInCluster(server, tls, port, host, socketPath, onListen); } catch (err) { setTimeout(() => server.emit("error", err), 1); } @@ -690,7 +692,109 @@ Server.prototype.listen = function () { return this; }; -Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort, onListen) { +function isBindableByPrimary(port, host, socketPath) { + if (socketPath) return typeof socketPath === "string"; + if (host != null && typeof host !== "string") return false; + return typeof port === "number" && port === (port | 0) && port >= 0 && port <= 65535; +} + +// A worker that bound its own socket still reports it, so cluster.on("listening") fires. +function notifyPrimaryWhenListening(server, port, host, socketPath) { + server.once("listening", () => { + if (!process.connected) return; + cluster.worker.state = "listening"; + const address = server.address(); + const isObjectAddress = address !== null && typeof address === "object"; + const boundHost = host && isObjectAddress ? address : null; + const message = { + cmd: "NODE_CLUSTER", + act: "listening", + port: socketPath ? -1 : (isObjectAddress && address.port) || port, + data: null, + address: socketPath ?? (boundHost && boundHost.address) ?? null, + addressType: socketPath ? -1 : boundHost && boundHost.family === "IPv6" ? 6 : 4, + }; + process.send(message, undefined, kInternalSendOptions); + }); +} + +// Like net's listenInCluster: the primary binds the address once (or hands out the socket it +// already bound for this key, which is how every worker's listen(0) ends up on one port) and +// sends each worker the descriptor to accept on. Bun.serve accepts natively, so the server +// always asks for a shared handle (sharedOnly), the same way a TLS net.Server does. +function listenInCluster(server, tls, port, host, socketPath, onListen) { + const listeningId = (server[kClusterListeningId] = (server[kClusterListeningId] || 0) + 1); + const listenArgs = { server, listeningId, tls, port, host, socketPath, onListen }; + + // Same (address, port, addressType) tuple net.Server sends, so the primary keys the handle + // the way node does: pipes are port -1 / addressType -1, and _getServer adds the per-listen + // index that makes each worker's first listen(0) land on the same handle. + if (socketPath) { + queryPrimary(listenArgs, socketPath, -1, -1); + } else if (!host) { + queryPrimary(listenArgs, null, port, 4); + } else if (isIP(host) !== 0) { + queryPrimary(listenArgs, host, port, isIP(host)); + } else { + // The primary binds an address, not a name; node resolves it in the worker first. + require("node:dns").lookup(host, (err, ip, family) => { + if (listeningId !== server[kClusterListeningId]) return; + if (err) { + server.emit("error", err); + return; + } + queryPrimary(listenArgs, ip, port, family === 6 ? 6 : 4); + }); + } +} + +function queryPrimary({ server, listeningId, tls, port, host, socketPath, onListen }, address, queryPort, addressType) { + const serverQuery = { address, port: queryPort, addressType, fd: undefined, flags: 0, sharedOnly: true }; + cluster._getServer(server, serverQuery, function listenOnPrimaryHandle(err, handle, reply) { + if (listeningId !== server[kClusterListeningId]) { + // close() or another listen() came first; give the primary's handle straight back. + handle?.close(); + return; + } + const sharedFd = handle?.sharedFd; + if (!err && typeof sharedFd !== "number") { + // A primary that answers a sharedOnly query with a round-robin handle is not Bun's; + // connections handed over one at a time cannot be fed to Bun.serve. + handle?.close(); + err = process.binding("uv").UV_EINVAL; + } + if (err) { + const ex = new ExceptionWithHostPort(err, "bind", address, queryPort); + if (typeof reply?.bunHint === "string") ex.message += `\n note: ${reply.bunHint}`; + server.emit("error", ex); + return; + } + server[kClusterHandle] = handle; + // Lets worker.disconnect() close this server along with the worker's other servers. + handle[kClusterOwner] = server; + // Once adopted, the listener owns the descriptor and releasing the handle must not close + // it; a failed Bun.serve() never took it, so the release below closes it after all. + handle.adopted = true; + try { + server[kRealListen](tls, port, host, socketPath, false, onListen, sharedFd); + } catch (err) { + handle.adopted = false; + releaseClusterHandle(server); + setTimeout(() => server.emit("error", err), 1); + } + }); +} + +// The primary keeps the socket bound until every worker that asked for it has let go. +function releaseClusterHandle(server) { + const handle = server[kClusterHandle]; + if (handle === undefined) return; + server[kClusterHandle] = undefined; + handle[kClusterOwner] = null; + handle.close(); +} + +Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort, onListen, fd) { { const ResponseClass = this[optionsSymbol].ServerResponse || ServerResponse; const RequestClass = this[optionsSymbol].IncomingMessage || IncomingMessage; @@ -708,6 +812,8 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort hostname: host, unix: socketPath, reusePort, + // Set when the cluster primary bound the socket: accept on its descriptor instead of binding. + fd, // Bindings to be used for WS Server websocket: { open(ws) { diff --git a/src/runtime/server/ServerConfig.rs b/src/runtime/server/ServerConfig.rs index 1b83c3c08c0a..d37e02d90616 100644 --- a/src/runtime/server/ServerConfig.rs +++ b/src/runtime/server/ServerConfig.rs @@ -55,6 +55,12 @@ pub struct ServerConfig { pub(crate) websocket: Option, pub(crate) reuse_port: bool, + /// Accept on this already-bound descriptor instead of binding `address`. + /// Set by node:http in a cluster worker, which receives the primary's + /// shared listen socket; `address` still describes what was bound so + /// `address`/`url`/errors report it. Consumed by the listen socket on + /// success; a failed listen leaves it to the caller. + pub(crate) listen_fd: Option, pub(crate) id: Box<[u8]>, pub(crate) allow_hot: bool, pub(crate) ipv6_only: bool, @@ -88,6 +94,7 @@ impl Default for ServerConfig { on_node_http_request: JSValue::ZERO, websocket: None, reuse_port: false, + listen_fd: None, id: Box::default(), allow_hot: true, ipv6_only: false, @@ -274,6 +281,7 @@ impl ServerConfig { on_node_http_request: self.on_node_http_request, websocket: self.websocket.take(), reuse_port: self.reuse_port, + listen_fd: self.listen_fd, id: core::mem::take(&mut self.id), allow_hot: self.allow_hot, ipv6_only: self.ipv6_only, @@ -622,6 +630,23 @@ impl ServerConfig { // ─── from_js + JS-side parsing ─────────────────────────────────────────────── +/// A JS `fd` number as the OS descriptor `us_socket_group_listen_fd` accepts: +/// the POSIX int, or (as for `Bun.listen({ fd })`) the raw Windows SOCKET value. +fn listen_fd_from_number(number: f64) -> Option { + if !(number >= 0.0 && number.fract() == 0.0) { + return None; + } + #[cfg(not(windows))] + { + (number <= i32::MAX as f64).then(|| number as uws::LIBUS_SOCKET_DESCRIPTOR) + } + #[cfg(windows)] + { + // Exact in f64 up to 2^53; a SOCKET is a kernel handle, far below that. + (number < (1u64 << 53) as f64).then(|| number as u64 as uws::LIBUS_SOCKET_DESCRIPTOR) + } +} + fn validate_route_name(global: &JSGlobalObject, path: &[u8]) -> JsResult<()> { // Already validated by the caller debug_assert!(!path.is_empty() && path[0] == b'/'); @@ -1355,6 +1380,23 @@ impl ServerConfig { ))); } args.on_node_http_request = on_request_; + + // Only node:http servers may hand over a descriptor (see `listen_fd`); + // for a plain `Bun.serve()` the key is not an option and is left alone. + if let Some(fd) = arg.get(global, "fd")? { + let Some(fd) = fd.get_number().and_then(listen_fd_from_number) else { + return Err(global.throw_invalid_arguments(format_args!( + "Expected fd to be a non-negative integer" + ))); + }; + args.listen_fd = Some(fd); + // The descriptor is consumed by this listen; a `--hot` reload + // matching a previous server by address would silently drop it. + args.allow_hot = false; + } + } + if global.has_exception() { + return Err(JsError::Thrown); } if let Some(on_request_) = arg.get_truthy(global, "fetch")? { @@ -1469,6 +1511,11 @@ impl ServerConfig { "Cannot disable http1 with a unix socket — HTTP/3 over AF_UNIX is not supported", ))); } + if args.http3 && args.listen_fd.is_some() { + return Err(global.throw_invalid_arguments(format_args!( + "Cannot combine http3 with an inherited listen fd" + ))); + } // ---- base_uri / base_url normalization ---- if !args.base_uri.is_empty() { diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index 01077a8455a2..628815833f87 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1719,10 +1719,14 @@ impl NewServer { } self.notify_inspector_server_stopped(); - if let server_config::Address::Unix(path) = &self.config.address { - let bytes = path.as_bytes(); - if !bytes.is_empty() && bytes[0] != 0 { - let _ = bun_sys::unlink(path.as_zstr()); + // An inherited descriptor's socket file belongs to whoever bound it (the + // cluster primary unlinks it once the last worker lets go). + if self.config.listen_fd.is_none() { + if let server_config::Address::Unix(path) = &self.config.address { + let bytes = path.as_bytes(); + if !bytes.is_empty() && bytes[0] != 0 { + let _ = bun_sys::unlink(path.as_zstr()); + } } } @@ -1987,6 +1991,21 @@ impl NewServer { self.listener = None; let global = self.global_this(); + if self.config.listen_fd.is_some() { + // Nothing was bound here: the descriptor itself could not be + // listened on (HttpContext::listen_fd stores the cause in errno). + // EINVAL is what node reports for a descriptor it cannot serve on. + let errno = match bun_sys::get_errno(-1i32) { + bun_sys::E::SUCCESS => bun_sys::E::EINVAL, + e => e, + }; + let err = jsc::SystemError::from( + bun_sys::Error::from_code(errno, bun_sys::Tag::listen).to_system_error(), + ); + let _ = global.throw_value(err.to_error_instance(global)); + return; + } + let error_instance = match &self.config.address { server_config::Address::Tcp { port, @@ -2997,11 +3016,13 @@ impl NewServer { enum Addr { Tcp { port: u16, host: *const c_char }, Unix { ptr: *const u8, len: usize }, + Fd(uws_sys::LIBUS_SOCKET_DESCRIPTOR), } let (addr, http1, options) = { let cfg = &this_ref.get().config; - let addr = match &cfg.address { - server_config::Address::Tcp { port, hostname } => { + let addr = match (cfg.listen_fd, &cfg.address) { + (Some(fd), _) => Addr::Fd(fd), + (None, server_config::Address::Tcp { port, hostname }) => { let mut host: *const c_char = core::ptr::null(); if let Some(existing) = hostname.as_deref() { let bytes = existing.as_bytes(); @@ -3016,7 +3037,7 @@ impl NewServer { } Addr::Tcp { port: *port, host } } - server_config::Address::Unix(unix) => Addr::Unix { + (None, server_config::Address::Unix(unix)) => Addr::Unix { ptr: unix.as_ptr().cast(), len: unix.as_bytes().len(), }, @@ -3025,6 +3046,21 @@ impl NewServer { }; match addr { + Addr::Fd(fd) => { + // `ServerConfig::from_js` rejects http3 together with a descriptor, + // so there is no H3 listener to pair with it. + // SAFETY: app is a live uws handle owned by this server. No + // `&*this` is live across this call; the trampoline's + // `&mut *this` is the sole borrow while it runs. + unsafe { + (*app).listen_fd( + Some(trampoline::on_listen::), + this.cast::(), + fd, + options, + ); + } + } Addr::Tcp { port, host } => { // With `{port: 0, http3: true}` we bind TCP:0 (kernel picks N), // then must bind UDP:N for QUIC so Alt-Svc works. UDP:N may diff --git a/src/uws_sys/App.rs b/src/uws_sys/App.rs index c051b5aeb846..392f6721aab8 100644 --- a/src/uws_sys/App.rs +++ b/src/uws_sys/App.rs @@ -8,8 +8,8 @@ use bun_http_types::Method::Method; use crate::socket_context::BunSocketContextOptions; use crate::web_socket::c::uws_ws; use crate::{ - ListenSocket as UwsListenSocket, Opcode, Request, SendStatus, WebSocketBehavior, us_socket_t, - uws_res, + LIBUS_SOCKET_DESCRIPTOR, ListenSocket as UwsListenSocket, Opcode, Request, SendStatus, + WebSocketBehavior, us_socket_t, uws_res, }; // This file provides Rust bindings for the uWebSockets App class. @@ -287,6 +287,31 @@ impl App { } } + /// Accept on `fd`, a descriptor something else already bound (node:cluster's + /// shared listen handle). The listen socket takes ownership of `fd` when the + /// handler receives a non-null socket; otherwise the caller still owns it. + pub fn listen_fd( + &mut self, + handler: c::uws_listen_handler, + user_data: *mut c_void, + fd: LIBUS_SOCKET_DESCRIPTOR, + options: i32, + ) { + // Callers supply the C-ABI shim directly. + // SAFETY: self is a valid app; `fd` is a plain value and `handler`/`user_data` + // are only invoked synchronously inside this call. + unsafe { + c::uws_app_listen_fd( + Self::SSL_FLAG, + std::ptr::from_mut::(self).cast::(), + fd, + options, + handler, + user_data, + ) + } + } + pub fn listen_on_unix_socket( &mut self, handler: extern "C" fn(*mut UwsListenSocket, *const c_char, i32, *mut c_void), @@ -601,6 +626,14 @@ pub mod c { handler: uws_listen_handler, user_data: *mut c_void, ); + pub(crate) fn uws_app_listen_fd( + ssl: i32, + app: *mut uws_app_t, + fd: LIBUS_SOCKET_DESCRIPTOR, + options: i32, + handler: uws_listen_handler, + user_data: *mut c_void, + ); pub(crate) fn uws_num_subscribers( ssl: i32, app: *mut uws_app_t, diff --git a/src/uws_sys/libuwsockets.cpp b/src/uws_sys/libuwsockets.cpp index 89a62dcca2b1..f4b507525f0c 100644 --- a/src/uws_sys/libuwsockets.cpp +++ b/src/uws_sys/libuwsockets.cpp @@ -479,6 +479,31 @@ extern "C" } } + void uws_app_listen_fd(int ssl, uws_app_t *app, LIBUS_SOCKET_DESCRIPTOR fd, int32_t options, + uws_listen_handler handler, void *user_data) + { + if (ssl) + { + uWS::SSLApp *uwsApp = (uWS::SSLApp *)app; + uwsApp->listen_fd( + fd, options, + [handler, user_data](struct us_listen_socket_t *listen_socket) + { + handler((struct us_listen_socket_t *)listen_socket, user_data); + }); + } + else + { + uWS::App *uwsApp = (uWS::App *)app; + uwsApp->listen_fd( + fd, options, + [handler, user_data](struct us_listen_socket_t *listen_socket) + { + handler((struct us_listen_socket_t *)listen_socket, user_data); + }); + } + } + /* callback, path to unix domain socket */ void uws_app_listen_domain(int ssl, uws_app_t *app, const char *domain, size_t pathlen, uws_listen_domain_handler handler, void *user_data) { diff --git a/test/js/node/cluster.test.ts b/test/js/node/cluster.test.ts index 64d101f15e8f..7c000d2a72dd 100644 --- a/test/js/node/cluster.test.ts +++ b/test/js/node/cluster.test.ts @@ -1276,3 +1276,362 @@ if (cluster.isPrimary) { }); expect(exitCode).toBe(0); }, 30_000); + +// node:http servers are backed by Bun.serve, whose accept loop is native, so in a worker they +// take the shared-handle path (the primary binds once and ships the descriptor) rather than +// round-robin. On Windows a worker still binds its own SO_REUSEADDR socket, because several +// processes accepting on copies of one listening socket block each other in accept() there. +const httpListenZeroFixture = /* js */ ` +const cluster = require("node:cluster"); +const fs = require("node:fs"); +const path = require("node:path"); +const kind = process.env.MODULE; +const mod = require("node:" + kind); +const serverOptions = + kind === "https" + ? { key: fs.readFileSync(path.join(__dirname, "key.pem")), cert: fs.readFileSync(path.join(__dirname, "cert.pem")) } + : {}; +const workerCount = 2; + +if (cluster.isPrimary) { + const fail = reason => { + console.log(JSON.stringify({ failed: reason })); + process.exit(1); + }; + cluster.on("exit", (worker, code, signal) => fail("worker " + worker.id + " exited early: " + (signal ?? code))); + // What cluster's 'listening' event reported for each worker, keyed by worker id. + const listeningPorts = {}; + cluster.on("listening", (worker, address) => (listeningPorts[worker.id] = address.port)); + // Resolves with the port the worker itself sees in server.address(). The worker's 'listening' + // notification to the primary precedes the message it sends from its listen callback. + const forkAndWaitForPort = env => + new Promise(resolve => { + const worker = cluster.fork(env); + worker.once("message", message => { + if (message.error) fail("worker " + worker.id + " listen error: " + message.error); + resolve({ id: worker.id, port: message.port }); + }); + }); + + (async () => { + const [exclusive, ...shared] = await Promise.all([ + forkAndWaitForPort({ EXCLUSIVE: "1" }), + ...Array.from({ length: workerCount }, () => forkAndWaitForPort()), + ]); + + const port = shared[0].port; + let served = 0; + for (let i = 0; i < 4; i++) { + const response = await fetch(kind + "://127.0.0.1:" + port + "/", { tls: { rejectUnauthorized: false } }); + if ((await response.text()).startsWith("worker ")) served++; + } + + console.log( + JSON.stringify({ + workers: shared.length, + distinctReportedPorts: new Set(shared.map(worker => worker.port)).size, + distinctListeningPorts: new Set(shared.map(worker => listeningPorts[worker.id])).size, + listeningMatchesReported: shared.every(worker => listeningPorts[worker.id] === worker.port), + exclusiveWorkerHasOwnPort: exclusive.port > 0 && exclusive.port !== port, + exclusiveListeningMatchesReported: listeningPorts[exclusive.id] === exclusive.port, + served, + }), + ); + cluster.removeAllListeners("exit"); + for (const id in cluster.workers) cluster.workers[id].process.kill(); + process.exit(0); + })().catch(error => fail(String(error))); +} else { + const server = mod.createServer(serverOptions, (req, res) => { + res.setHeader("Connection", "close"); + res.end("worker " + cluster.worker.id); + }); + server.on("error", error => process.send({ error: error.code })); + const report = () => process.send({ port: server.address().port }); + if (process.env.EXCLUSIVE) server.listen({ port: 0, exclusive: true }, report); + else server.listen(0, report); +} +`; + +test.skipIf(isWindows).each(["http", "https"])( + "%s workers all share the one port the primary picked for listen(0)", + async kind => { + using dir = tempDir("cluster-http-listen0", { + "cert.pem": tlsCerts.cert, + "key.pem": tlsCerts.key, + "fixture.js": httpListenZeroFixture, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.js"], + env: { ...bunEnv, MODULE: kind }, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ out: JSON.parse(stdout.trim().split("\n").pop()!), stderr }).toEqual({ + out: { + workers: 2, + distinctReportedPorts: 1, + distinctListeningPorts: 1, + listeningMatchesReported: true, + exclusiveWorkerHasOwnPort: true, + exclusiveListeningMatchesReported: true, + served: 4, + }, + stderr: expect.any(String), + }); + expect(exitCode).toBe(0); + }, + 30_000, +); + +test.skipIf(isWindows)( + "http workers listening on one unix path share the primary's socket", + async () => { + using dir = tempDir("cluster-http-unix", { + "fixture.js": /* js */ ` +const cluster = require("node:cluster"); +const http = require("node:http"); +const fs = require("node:fs"); +const path = require("node:path"); +const SOCK = path.join(__dirname, "shared.sock"); + +if (cluster.isPrimary) { + const fail = reason => { + console.log(JSON.stringify({ failed: reason })); + process.exit(1); + }; + cluster.on("exit", (worker, code, signal) => fail("worker " + worker.id + " exited early: " + (signal ?? code))); + const listening = []; + cluster.on("listening", (_worker, address) => listening.push(address)); + const nextMessage = worker => + new Promise(resolve => + worker.once("message", message => { + if (message.error) fail("worker " + worker.id + ": " + message.error); + resolve(message); + }), + ); + const get = async () => { + const response = await fetch("http://localhost/", { unix: SOCK }); + return response.text(); + }; + const workers = [cluster.fork(), cluster.fork()]; + + (async () => { + const addresses = (await Promise.all(workers.map(nextMessage))).map(message => message.address); + const bodies = await Promise.all([get(), get(), get(), get()]); + + workers[0].send("close"); + await nextMessage(workers[0]); + const existsAfterFirstClose = fs.existsSync(SOCK); + const bodyAfterFirstClose = await get(); + + workers[1].send("close"); + await nextMessage(workers[1]); + // The worker's release (act: close) is sent before its reply, so the primary has handled it. + const existsAfterLastClose = fs.existsSync(SOCK); + + console.log( + JSON.stringify({ + addresses: addresses.map(address => address === SOCK), + listening: listening.map(address => ({ ...address, address: address.address === SOCK })), + served: bodies.filter(body => body.startsWith("worker ")).length, + existsAfterFirstClose, + servedAfterFirstClose: bodyAfterFirstClose.startsWith("worker "), + existsAfterLastClose, + }), + ); + cluster.removeAllListeners("exit"); + for (const worker of workers) worker.process.kill(); + process.exit(0); + })().catch(error => fail(String(error))); +} else { + const server = http.createServer((req, res) => { + res.setHeader("Connection", "close"); + res.end("worker " + cluster.worker.id); + }); + server.on("error", error => process.send({ error: error.code })); + process.on("message", () => server.close(() => process.send({ closed: true }))); + server.listen(SOCK, () => process.send({ address: server.address() })); +} +`, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.js"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ out: JSON.parse(stdout.trim().split("\n").pop()!), stderr }).toEqual({ + out: { + addresses: [true, true], + listening: [ + { address: true, addressType: -1, port: -1 }, + { address: true, addressType: -1, port: -1 }, + ], + served: 4, + existsAfterFirstClose: true, + servedAfterFirstClose: true, + existsAfterLastClose: false, + }, + stderr: expect.any(String), + }); + expect(exitCode).toBe(0); + }, + 30_000, +); + +test.skipIf(isWindows)( + "http worker reports a port the primary cannot bind the way node does", + async () => { + using dir = tempDir("cluster-http-bind-error", { + "fixture.js": /* js */ ` +const cluster = require("node:cluster"); +const http = require("node:http"); +const net = require("node:net"); + +if (cluster.isPrimary) { + const blocker = net.createServer(); + blocker.listen(0, "127.0.0.1", () => { + const port = blocker.address().port; + const worker = cluster.fork({ BUSY_PORT: String(port) }); + worker.on("message", error => { + console.log(JSON.stringify({ ...error, portMatches: error.port === port, message: error.message.replace(":" + port, ":PORT") })); + worker.process.kill(); + blocker.close(); + }); + }); +} else { + const server = http.createServer(); + server.on("error", error => + process.send({ code: error.code, syscall: error.syscall, address: error.address, port: error.port, message: error.message }), + ); + server.listen(+process.env.BUSY_PORT, "127.0.0.1"); +} +`, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.js"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ out: JSON.parse(stdout.trim()), stderr }).toEqual({ + out: { + code: "EADDRINUSE", + syscall: "bind", + address: "127.0.0.1", + port: expect.any(Number), + portMatches: true, + message: "bind EADDRINUSE 127.0.0.1:PORT", + }, + stderr: expect.any(String), + }); + expect(exitCode).toBe(0); + }, + 30_000, +); + +test.skipIf(isWindows)( + "http worker's close() releases the shared port so it can listen on it again", + async () => { + using dir = tempDir("cluster-http-relisten", { + "fixture.js": /* js */ ` +const cluster = require("node:cluster"); +const http = require("node:http"); + +if (cluster.isPrimary) { + const worker = cluster.fork(); + worker.on("message", result => { + console.log(JSON.stringify(result)); + worker.process.kill(); + }); +} else { + const first = http.createServer(); + first.listen(0, "127.0.0.1", () => { + const port = first.address().port; + first.close(); + const second = http.createServer(); + second.on("error", error => process.send({ relisten: error.code })); + second.listen(port, "127.0.0.1", () => process.send({ relisten: "ok", samePort: second.address().port === port })); + }); +} +`, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.js"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ out: JSON.parse(stdout.trim()), stderr }).toEqual({ + out: { relisten: "ok", samePort: true }, + stderr: expect.any(String), + }); + expect(exitCode).toBe(0); + }, + 30_000, +); + +test.skipIf(isWindows)( + "worker.disconnect() closes the worker's http server before the channel goes away", + async () => { + using dir = tempDir("cluster-http-disconnect", { + "fixture.js": /* js */ ` +const cluster = require("node:cluster"); +const http = require("node:http"); + +if (cluster.isPrimary) { + const worker = cluster.fork(); + let port; + worker.on("message", message => (port = message.port)); + const exited = new Promise(resolve => worker.on("exit", (code, signal) => resolve({ code, signal }))); + worker.on("disconnect", async () => { + const connectAfterDisconnect = await new Promise(resolve => { + http + .get({ host: "127.0.0.1", port }, response => { + response.resume(); + resolve("served " + response.statusCode); + }) + .on("error", error => resolve(error.code)); + }); + // A worker whose server is still up never exits on its own. + const killed = connectAfterDisconnect !== "ECONNREFUSED"; + if (killed) worker.process.kill(); + console.log(JSON.stringify({ connectAfterDisconnect, killed, exit: await exited })); + }); +} else { + const server = http.createServer((req, res) => res.end("still here")); + server.on("close", () => console.log("worker: server closed")); + server.listen(0, "127.0.0.1", () => { + process.send({ port: server.address().port }); + cluster.worker.disconnect(); + }); +} +`, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.js"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + const lines = stdout.trim().split("\n"); + expect({ workerLines: lines.slice(0, -1), out: JSON.parse(lines.at(-1)!), stderr }).toEqual({ + workerLines: ["worker: server closed"], + out: { connectAfterDisconnect: "ECONNREFUSED", killed: false, exit: { code: 0, signal: null } }, + stderr: expect.any(String), + }); + expect(exitCode).toBe(0); + }, + 30_000, +); From 36154f05e1acf61a134e9f9fa9e0495eb91b7f89 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 13 Aug 2026 09:54:19 +0000 Subject: [PATCH 2/5] node:http cluster listen: unbracket IPv6 hosts, fix late-failure teardown, update mixed-kind hint The primary binds addresses, so strip the brackets Bun.serve() accepts around an IPv6 literal before isIP and the query; listen(0, "[::1]") in a worker otherwise went to dns.lookup and failed with ENOTFOUND. A listen() that throws after Bun.serve() returned (while wiring up the listener) used to close the shared descriptor underneath the live listener; tear it down with close() instead, on every listen path, so a failed listen never leaves a listener running either. The hint the primary attaches when a shared-only query meets a round-robin handle now describes http/https/tls vs net servers, since http servers reach it too. --- src/js/internal/cluster/primary.ts | 5 +- src/js/node/_http_server.ts | 41 +++++-- test/js/node/cluster.test.ts | 175 ++++++++++++++++++++++++++++- 3 files changed, 207 insertions(+), 14 deletions(-) diff --git a/src/js/internal/cluster/primary.ts b/src/js/internal/cluster/primary.ts index d4840fccbcbf..5ac3fa36c680 100644 --- a/src/js/internal/cluster/primary.ts +++ b/src/js/internal/cluster/primary.ts @@ -231,8 +231,9 @@ function queryServer(worker, message) { if (cachedHandle && !cachedHandle.has(worker)) handle = cachedHandle; const kSharedOnlyHint = - "TLS and non-TLS cluster workers cannot share the same address:port under SCHED_RR " + - "(Bun's TLS accept is native and cannot adopt round-robin connection fds)"; + "tls, https and http servers in cluster workers accept natively, so they can only share a listening socket, " + + "while under SCHED_RR the primary hands a net server its connections one by one; " + + "the two kinds cannot share one address:port (use different ports, or cluster.schedulingPolicy = cluster.SCHED_NONE)"; if (handle !== undefined && message.sharedOnly === true && handle instanceof RoundRobinHandle) { send(worker, { errno: UV_EINVAL, key, ack: message.seq, data: handle.data, bunHint: kSharedOnlyHint }, null); return; diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 3e38e5af6c1c..22ee2a7f83da 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -655,6 +655,7 @@ Server.prototype.listen = function () { onListen = lastArg; } + const previousBunServer = server[serverSymbol]; try { // listenInCluster @@ -686,12 +687,19 @@ Server.prototype.listen = function () { listenInCluster(server, tls, port, host, socketPath, onListen); } catch (err) { - setTimeout(() => server.emit("error", err), 1); + emitListenError(server, err, previousBunServer); } return this; }; +// A listen() that failed after Bun.serve() had already returned (while wiring the listener up) +// must not leave that listener running behind the 'error' it reports. +function emitListenError(server, err, previousBunServer) { + if (server[serverSymbol] !== previousBunServer) server.close(); + setTimeout(() => server.emit("error", err), 1); +} + function isBindableByPrimary(port, host, socketPath) { if (socketPath) return typeof socketPath === "string"; if (host != null && typeof host !== "string") return false; @@ -731,13 +739,21 @@ function listenInCluster(server, tls, port, host, socketPath, onListen) { // index that makes each worker's first listen(0) land on the same handle. if (socketPath) { queryPrimary(listenArgs, socketPath, -1, -1); - } else if (!host) { + return; + } + if (!host) { queryPrimary(listenArgs, null, port, 4); - } else if (isIP(host) !== 0) { - queryPrimary(listenArgs, host, port, isIP(host)); + return; + } + // Bun.serve() (the other way a worker binds) takes an IPv6 literal in brackets; the primary's + // bind, like isIP, wants it bare. `host` itself still goes to Bun.serve for reporting. + const address = host.length > 2 && host.charCodeAt(0) === 0x5b /* [ */ ? host.slice(1, -1) : host; + const addressType = isIP(address); + if (addressType !== 0) { + queryPrimary(listenArgs, address, port, addressType); } else { // The primary binds an address, not a name; node resolves it in the worker first. - require("node:dns").lookup(host, (err, ip, family) => { + require("node:dns").lookup(address, (err, ip, family) => { if (listeningId !== server[kClusterListeningId]) return; if (err) { server.emit("error", err); @@ -772,15 +788,20 @@ function queryPrimary({ server, listeningId, tls, port, host, socketPath, onList server[kClusterHandle] = handle; // Lets worker.disconnect() close this server along with the worker's other servers. handle[kClusterOwner] = server; - // Once adopted, the listener owns the descriptor and releasing the handle must not close - // it; a failed Bun.serve() never took it, so the release below closes it after all. + // Once Bun.serve() returns, its listener owns the descriptor and releasing the handle must + // only notify the primary; until then the descriptor is ours to close. handle.adopted = true; + const previousBunServer = server[serverSymbol]; try { server[kRealListen](tls, port, host, socketPath, false, onListen, sharedFd); } catch (err) { - handle.adopted = false; - releaseClusterHandle(server); - setTimeout(() => server.emit("error", err), 1); + if (server[serverSymbol] === previousBunServer) { + // Bun.serve() itself failed, so the descriptor is still ours. (Otherwise the listener + // has it, and the close() in emitListenError releases the handle along with it.) + handle.adopted = false; + releaseClusterHandle(server); + } + emitListenError(server, err, previousBunServer); } }); } diff --git a/test/js/node/cluster.test.ts b/test/js/node/cluster.test.ts index 7c000d2a72dd..14e64e99b836 100644 --- a/test/js/node/cluster.test.ts +++ b/test/js/node/cluster.test.ts @@ -203,7 +203,7 @@ if (cluster.isPrimary) { }); const { stdout } = await bunRun(joinP(dir, "main.ts"), bunEnv); expect(stdout).toContain("tls listen error code: EINVAL"); - expect(stdout).toContain("TLS and non-TLS cluster workers cannot share"); + expect(stdout).toContain("the two kinds cannot share one address:port"); }); test("cluster pipe listen error carries no port suffix", async () => { @@ -842,7 +842,7 @@ if (cluster.isPrimary) { }); const { stdout } = await bunRun(joinP(dir, "main.ts"), bunEnv); expect(stdout).toContain("net listen error code: EINVAL"); - expect(stdout).toContain("TLS and non-TLS cluster workers cannot share"); + expect(stdout).toContain("the two kinds cannot share one address:port"); }, 30_000); test.skipIf(isWindows)( @@ -1635,3 +1635,174 @@ if (cluster.isPrimary) { }, 30_000, ); + +test.skipIf(isWindows)( + "http worker's listen(0) on the key a net worker's round-robin handle owns fails with EINVAL", + async () => { + // Like the TLS variant above: a worker's first listen(0) on an address always gets index 0, + // so both workers ask for the same key and the http worker's shared-only query meets the + // round-robin handle the net worker got. The hint has to describe http now, not just TLS. + using dir = tempDir("cluster-http-mixed-kinds", { + "fixture.js": /* js */ ` +const cluster = require("node:cluster"); +const http = require("node:http"); +const net = require("node:net"); + +if (cluster.isPrimary) { + const netWorker = cluster.fork({ ROLE: "net" }); + cluster.once("listening", () => { + const httpWorker = cluster.fork({ ROLE: "http" }); + httpWorker.on("message", error => { + console.log(JSON.stringify(error)); + netWorker.process.kill(); + httpWorker.process.kill(); + }); + }); +} else if (process.env.ROLE === "net") { + net.createServer(() => {}).listen(0, "127.0.0.1"); +} else { + const server = http.createServer(); + server.on("error", error => + process.send({ + code: error.code, + syscall: error.syscall, + firstLine: error.message.split("\\n")[0], + namesHttp: error.message.includes("http servers"), + hasHint: error.message.includes("cannot share one address:port"), + }), + ); + server.on("listening", () => process.send({ unexpected: "listening" })); + server.listen(0, "127.0.0.1"); +} +`, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.js"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ out: JSON.parse(stdout.trim()), stderr }).toEqual({ + out: { code: "EINVAL", syscall: "bind", firstLine: "bind EINVAL 127.0.0.1", namesHttp: true, hasHint: true }, + stderr: expect.any(String), + }); + expect(exitCode).toBe(0); + }, + 30_000, +); + +// One fixture for the host spellings a worker has to translate before asking the primary: the +// primary binds addresses, so a name is resolved first and an IPv6 literal loses the brackets +// Bun.serve() accepts. HOST is what the worker passes to listen(); the primary prints what the +// 'listening' event reported next to what the worker's server.address() says it bound. +const httpHostFormFixture = /* js */ ` +const cluster = require("node:cluster"); +const http = require("node:http"); + +if (cluster.isPrimary) { + let listening; + cluster.on("listening", (_worker, address) => (listening = address)); + const worker = cluster.fork(); + worker.on("message", ({ error, bound }) => { + console.log(JSON.stringify({ error, listening: listening && { address: listening.address, addressType: listening.addressType }, bound })); + worker.process.kill(); + }); +} else { + const server = http.createServer(); + server.on("error", error => process.send({ error: error.code })); + server.listen(0, process.env.HOST, () => { + const { address, family, port } = server.address(); + process.send({ bound: { address, family, hasPort: port > 0 } }); + }); +} +`; + +async function runHttpHostFormFixture(host: string) { + using dir = tempDir("cluster-http-host-form", { "fixture.js": httpHostFormFixture }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.js"], + env: { ...bunEnv, HOST: host }, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect(stderr).toEqual(expect.any(String)); + expect(exitCode).toBe(0); + return JSON.parse(stdout.trim()); +} + +test.skipIf(isWindows || !isIPv6())( + "http worker listen(0, '[::1]') binds ::1 through the primary", + async () => { + expect(await runHttpHostFormFixture("[::1]")).toEqual({ + listening: { address: "::1", addressType: 6 }, + bound: { address: "::1", family: "IPv6", hasPort: true }, + }); + }, + 30_000, +); + +test.skipIf(isWindows)( + "http worker listen(0, 'localhost') resolves the name before asking the primary", + async () => { + const out = await runHttpHostFormFixture("localhost"); + expect(out).toEqual({ + listening: { address: out.bound.address, addressType: out.bound.family === "IPv6" ? 6 : 4 }, + bound: { address: expect.any(String), family: expect.stringMatching(/^IPv[46]$/), hasPort: true }, + }); + expect(net.isIP(out.bound.address)).toBeGreaterThan(0); + }, + 30_000, +); + +test.skipIf(isWindows)( + "http worker whose listen fails after Bun.serve() took the shared socket is left not running", + async () => { + using dir = tempDir("cluster-http-listen-late-failure", { + "fixture.js": /* js */ ` +const cluster = require("node:cluster"); +const http = require("node:http"); + +if (cluster.isPrimary) { + const worker = cluster.fork(); + worker.on("message", result => { + console.log(JSON.stringify(result)); + worker.process.kill(); + }); +} else { + const server = http.createServer(); + // requireHostHeader is pushed to the native listener right after Bun.serve() returns, so a + // throwing getter fails listen() at the point where the listener already owns the descriptor. + Object.defineProperty(server, "requireHostHeader", { + configurable: true, + get() { + throw new Error("boom from requireHostHeader"); + }, + }); + server.on("listening", () => process.send({ unexpected: "listening" })); + server.on("error", error => { + server.close(closeError => process.send({ listenError: error.message, closeError: closeError?.code ?? null })); + }); + server.listen(0, "127.0.0.1"); +} +`, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.js"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ out: JSON.parse(stdout.trim()), stderr }).toEqual({ + out: { listenError: "boom from requireHostHeader", closeError: "ERR_SERVER_NOT_RUNNING" }, + stderr: expect.any(String), + }); + expect(exitCode).toBe(0); + }, + 30_000, +); From 7c34434f2ff89595939a89d62b929c38a3f8fd57 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 13 Aug 2026 10:11:33 +0000 Subject: [PATCH 3/5] node:http cluster listen: shorten comments --- packages/bun-uws/src/HttpContext.h | 5 ++- src/js/node/_http_server.ts | 50 ++++++++++-------------------- src/runtime/server/ServerConfig.rs | 17 +++------- src/runtime/server/mod.rs | 10 ++---- src/uws_sys/App.rs | 8 ++--- 5 files changed, 29 insertions(+), 61 deletions(-) diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index fec88c92ff07..77e0daa02900 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -1072,9 +1072,8 @@ struct HttpContext { return socket; } - /* Accept on a descriptor that is already bound (node:cluster's shared listen handle: - * the primary bound it, every worker listen()s and accepts on its own copy). The fd - * is owned by the returned listen socket; on failure the caller still owns it. */ + /* Accept on an already-bound fd (a cluster primary's shared listen socket). The returned + * listen socket owns the fd; on failure (nullptr) the caller still does. */ us_listen_socket_t *listen_fd(struct ssl_ctx_st *sslCtx, LIBUS_SOCKET_DESCRIPTOR fd, int options) { int error = 0; auto* socket = us_socket_group_listen_fd(&group, socketKind(), sslCtx, fd, 512, options | LIBUS_LISTEN_DEFER_ACCEPT, socketExtSize(), &error); diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 22ee2a7f83da..aea719efb222 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -101,9 +101,9 @@ const { const kConnectionsCheckingInterval = Symbol("http.server.connectionsCheckingInterval"); const kTrackedConnections = Symbol("http.server.trackedConnections"); const kHttpAllowHalfOpen = Symbol("http.server.httpAllowHalfOpen"); -// Cluster worker state: the primary's shared listen handle this server accepts on, and the -// counter that tells a late primary reply that close() or a newer listen() superseded it. +// The cluster primary's shared listen handle this server accepts on (workers only). const kClusterHandle = Symbol("http.server.clusterHandle"); +// Bumped by listen() and close(); a primary reply carrying an older value is stale. const kClusterListeningId = Symbol("http.server.clusterListeningId"); // node.http trace events ('http.server.request' b/e). The agent module is @@ -666,21 +666,12 @@ Server.prototype.listen = function () { if (cluster === undefined) cluster = require("node:cluster"); - // The worker binds its own socket when node would (exclusive, or reusePort which implies - // it), when there is no primary to ask (NODE_UNIQUE_ID inherited by a plain child, or - // already disconnected), when the arguments are nothing the primary could bind (Bun.serve - // reports them), and on Windows: processes accepting on copies of one listening socket - // block each other in accept() there, and the primary's alternative of handing over each - // accepted connection is something Bun.serve's native accept loop cannot take. - if ( - exclusive || - !process.connected || - process.platform === "win32" || - !isBindableByPrimary(port, host, socketPath) - ) { + // Windows: processes accepting on copies of one listening socket block each other in accept(). + const bindsInWorker = + exclusive || !process.connected || process.platform === "win32" || !isBindableByPrimary(port, host, socketPath); + if (bindsInWorker) { notifyPrimaryWhenListening(server, port, host, socketPath); - // Only an exclusive listen gets to opt out of SO_REUSEPORT: the other cases still need - // it for the workers to end up on one port at all. + // Non-exclusive workers still share a fixed port through SO_REUSEPORT. server[kRealListen](tls, port, host, socketPath, exclusive ? reusePort : true, onListen); return this; } @@ -693,8 +684,7 @@ Server.prototype.listen = function () { return this; }; -// A listen() that failed after Bun.serve() had already returned (while wiring the listener up) -// must not leave that listener running behind the 'error' it reports. +// A listener that Bun.serve() created before the failure must not stay up behind the 'error'. function emitListenError(server, err, previousBunServer) { if (server[serverSymbol] !== previousBunServer) server.close(); setTimeout(() => server.emit("error", err), 1); @@ -726,17 +716,13 @@ function notifyPrimaryWhenListening(server, port, host, socketPath) { }); } -// Like net's listenInCluster: the primary binds the address once (or hands out the socket it -// already bound for this key, which is how every worker's listen(0) ends up on one port) and -// sends each worker the descriptor to accept on. Bun.serve accepts natively, so the server -// always asks for a shared handle (sharedOnly), the same way a TLS net.Server does. +// net.Server's listenInCluster, minus round-robin: Bun.serve accepts natively, so like a TLS +// net.Server it always asks for the shared handle (the primary's bound socket) of its key. function listenInCluster(server, tls, port, host, socketPath, onListen) { const listeningId = (server[kClusterListeningId] = (server[kClusterListeningId] || 0) + 1); const listenArgs = { server, listeningId, tls, port, host, socketPath, onListen }; - // Same (address, port, addressType) tuple net.Server sends, so the primary keys the handle - // the way node does: pipes are port -1 / addressType -1, and _getServer adds the per-listen - // index that makes each worker's first listen(0) land on the same handle. + // The (address, port, addressType) tuple must match net.Server's, since it keys the primary's handles. if (socketPath) { queryPrimary(listenArgs, socketPath, -1, -1); return; @@ -745,14 +731,13 @@ function listenInCluster(server, tls, port, host, socketPath, onListen) { queryPrimary(listenArgs, null, port, 4); return; } - // Bun.serve() (the other way a worker binds) takes an IPv6 literal in brackets; the primary's - // bind, like isIP, wants it bare. `host` itself still goes to Bun.serve for reporting. + // Bun.serve() accepts "[::1]"; the primary's bind (and isIP) want it bare. const address = host.length > 2 && host.charCodeAt(0) === 0x5b /* [ */ ? host.slice(1, -1) : host; const addressType = isIP(address); if (addressType !== 0) { queryPrimary(listenArgs, address, port, addressType); } else { - // The primary binds an address, not a name; node resolves it in the worker first. + // The primary binds addresses, not names (node resolves in the worker too). require("node:dns").lookup(address, (err, ip, family) => { if (listeningId !== server[kClusterListeningId]) return; if (err) { @@ -774,8 +759,7 @@ function queryPrimary({ server, listeningId, tls, port, host, socketPath, onList } const sharedFd = handle?.sharedFd; if (!err && typeof sharedFd !== "number") { - // A primary that answers a sharedOnly query with a round-robin handle is not Bun's; - // connections handed over one at a time cannot be fed to Bun.serve. + // A round-robin handle despite sharedOnly (a foreign primary): Bun.serve cannot use it. handle?.close(); err = process.binding("uv").UV_EINVAL; } @@ -788,16 +772,14 @@ function queryPrimary({ server, listeningId, tls, port, host, socketPath, onList server[kClusterHandle] = handle; // Lets worker.disconnect() close this server along with the worker's other servers. handle[kClusterOwner] = server; - // Once Bun.serve() returns, its listener owns the descriptor and releasing the handle must - // only notify the primary; until then the descriptor is ours to close. + // adopted: the handle's close() leaves the descriptor to the listener and only notifies the primary. handle.adopted = true; const previousBunServer = server[serverSymbol]; try { server[kRealListen](tls, port, host, socketPath, false, onListen, sharedFd); } catch (err) { if (server[serverSymbol] === previousBunServer) { - // Bun.serve() itself failed, so the descriptor is still ours. (Otherwise the listener - // has it, and the close() in emitListenError releases the handle along with it.) + // Bun.serve() itself failed: no listener took the descriptor, so close it here. handle.adopted = false; releaseClusterHandle(server); } diff --git a/src/runtime/server/ServerConfig.rs b/src/runtime/server/ServerConfig.rs index d37e02d90616..f5dd110c40fd 100644 --- a/src/runtime/server/ServerConfig.rs +++ b/src/runtime/server/ServerConfig.rs @@ -55,11 +55,8 @@ pub struct ServerConfig { pub(crate) websocket: Option, pub(crate) reuse_port: bool, - /// Accept on this already-bound descriptor instead of binding `address`. - /// Set by node:http in a cluster worker, which receives the primary's - /// shared listen socket; `address` still describes what was bound so - /// `address`/`url`/errors report it. Consumed by the listen socket on - /// success; a failed listen leaves it to the caller. + /// Accept on this already-bound descriptor (node:http in a cluster worker) instead + /// of binding `address`, which still describes it for `address`/`url`/errors. pub(crate) listen_fd: Option, pub(crate) id: Box<[u8]>, pub(crate) allow_hot: bool, @@ -630,8 +627,7 @@ impl ServerConfig { // ─── from_js + JS-side parsing ─────────────────────────────────────────────── -/// A JS `fd` number as the OS descriptor `us_socket_group_listen_fd` accepts: -/// the POSIX int, or (as for `Bun.listen({ fd })`) the raw Windows SOCKET value. +/// POSIX fd, or (like `Bun.listen({ fd })`) the raw SOCKET value on Windows. fn listen_fd_from_number(number: f64) -> Option { if !(number >= 0.0 && number.fract() == 0.0) { return None; @@ -642,7 +638,6 @@ fn listen_fd_from_number(number: f64) -> Option { } #[cfg(windows)] { - // Exact in f64 up to 2^53; a SOCKET is a kernel handle, far below that. (number < (1u64 << 53) as f64).then(|| number as u64 as uws::LIBUS_SOCKET_DESCRIPTOR) } } @@ -1381,8 +1376,7 @@ impl ServerConfig { } args.on_node_http_request = on_request_; - // Only node:http servers may hand over a descriptor (see `listen_fd`); - // for a plain `Bun.serve()` the key is not an option and is left alone. + // `fd` is not a public Bun.serve() option: only node:http servers get to pass one. if let Some(fd) = arg.get(global, "fd")? { let Some(fd) = fd.get_number().and_then(listen_fd_from_number) else { return Err(global.throw_invalid_arguments(format_args!( @@ -1390,8 +1384,7 @@ impl ServerConfig { ))); }; args.listen_fd = Some(fd); - // The descriptor is consumed by this listen; a `--hot` reload - // matching a previous server by address would silently drop it. + // A `--hot` reload reusing the old server by address would drop the descriptor. args.allow_hot = false; } } diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index 628815833f87..073e0c77e8c7 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1719,8 +1719,7 @@ impl NewServer { } self.notify_inspector_server_stopped(); - // An inherited descriptor's socket file belongs to whoever bound it (the - // cluster primary unlinks it once the last worker lets go). + // A shared socket's file is unlinked by whoever bound it (the cluster primary). if self.config.listen_fd.is_none() { if let server_config::Address::Unix(path) = &self.config.address { let bytes = path.as_bytes(); @@ -1992,9 +1991,7 @@ impl NewServer { let global = self.global_this(); if self.config.listen_fd.is_some() { - // Nothing was bound here: the descriptor itself could not be - // listened on (HttpContext::listen_fd stores the cause in errno). - // EINVAL is what node reports for a descriptor it cannot serve on. + // HttpContext::listen_fd left the cause in errno; EINVAL is node's answer for an unusable fd. let errno = match bun_sys::get_errno(-1i32) { bun_sys::E::SUCCESS => bun_sys::E::EINVAL, e => e, @@ -3047,8 +3044,7 @@ impl NewServer { match addr { Addr::Fd(fd) => { - // `ServerConfig::from_js` rejects http3 together with a descriptor, - // so there is no H3 listener to pair with it. + // No H3 listener to pair with it: `from_js` rejects http3 together with a descriptor. // SAFETY: app is a live uws handle owned by this server. No // `&*this` is live across this call; the trampoline's // `&mut *this` is the sole borrow while it runs. diff --git a/src/uws_sys/App.rs b/src/uws_sys/App.rs index 392f6721aab8..3066a99d135b 100644 --- a/src/uws_sys/App.rs +++ b/src/uws_sys/App.rs @@ -287,9 +287,8 @@ impl App { } } - /// Accept on `fd`, a descriptor something else already bound (node:cluster's - /// shared listen handle). The listen socket takes ownership of `fd` when the - /// handler receives a non-null socket; otherwise the caller still owns it. + /// Accept on an already-bound `fd`; it is owned by the listen socket the handler + /// receives, or still by the caller when the handler receives null. pub fn listen_fd( &mut self, handler: c::uws_listen_handler, @@ -298,8 +297,7 @@ impl App { options: i32, ) { // Callers supply the C-ABI shim directly. - // SAFETY: self is a valid app; `fd` is a plain value and `handler`/`user_data` - // are only invoked synchronously inside this call. + // SAFETY: self is a valid app; the handler is only invoked synchronously inside this call. unsafe { c::uws_app_listen_fd( Self::SSL_FLAG, From 6255a6ae2bb3e38dc008e0d5d9b45a5d6e0c97e1 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 13 Aug 2026 10:23:21 +0000 Subject: [PATCH 4/5] node:http cluster listen: single-line comments --- packages/bun-uws/src/HttpContext.h | 3 +-- src/js/node/_http_server.ts | 3 +-- src/runtime/server/ServerConfig.rs | 3 +-- src/runtime/server/mod.rs | 1 + src/uws_sys/App.rs | 3 +-- 5 files changed, 5 insertions(+), 8 deletions(-) diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index 77e0daa02900..edd583eca595 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -1072,8 +1072,7 @@ struct HttpContext { return socket; } - /* Accept on an already-bound fd (a cluster primary's shared listen socket). The returned - * listen socket owns the fd; on failure (nullptr) the caller still does. */ + /* Accept on an already-bound fd; the returned listen socket owns it, on failure (nullptr) the caller still does. */ us_listen_socket_t *listen_fd(struct ssl_ctx_st *sslCtx, LIBUS_SOCKET_DESCRIPTOR fd, int options) { int error = 0; auto* socket = us_socket_group_listen_fd(&group, socketKind(), sslCtx, fd, 512, options | LIBUS_LISTEN_DEFER_ACCEPT, socketExtSize(), &error); diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index aea719efb222..26371786c1c8 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -716,8 +716,7 @@ function notifyPrimaryWhenListening(server, port, host, socketPath) { }); } -// net.Server's listenInCluster, minus round-robin: Bun.serve accepts natively, so like a TLS -// net.Server it always asks for the shared handle (the primary's bound socket) of its key. +// net.Server's listenInCluster, but always sharedOnly (like TLS there): Bun.serve accepts natively. function listenInCluster(server, tls, port, host, socketPath, onListen) { const listeningId = (server[kClusterListeningId] = (server[kClusterListeningId] || 0) + 1); const listenArgs = { server, listeningId, tls, port, host, socketPath, onListen }; diff --git a/src/runtime/server/ServerConfig.rs b/src/runtime/server/ServerConfig.rs index f5dd110c40fd..9165fd6d529e 100644 --- a/src/runtime/server/ServerConfig.rs +++ b/src/runtime/server/ServerConfig.rs @@ -55,8 +55,7 @@ pub struct ServerConfig { pub(crate) websocket: Option, pub(crate) reuse_port: bool, - /// Accept on this already-bound descriptor (node:http in a cluster worker) instead - /// of binding `address`, which still describes it for `address`/`url`/errors. + /// Accept on this already-bound descriptor instead of binding `address` (which still describes it). pub(crate) listen_fd: Option, pub(crate) id: Box<[u8]>, pub(crate) allow_hot: bool, diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index 073e0c77e8c7..b0b8c01bc26e 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -3045,6 +3045,7 @@ impl NewServer { match addr { Addr::Fd(fd) => { // No H3 listener to pair with it: `from_js` rejects http3 together with a descriptor. + // SAFETY: app is a live uws handle owned by this server. No // `&*this` is live across this call; the trampoline's // `&mut *this` is the sole borrow while it runs. diff --git a/src/uws_sys/App.rs b/src/uws_sys/App.rs index 3066a99d135b..09d1ac882bcb 100644 --- a/src/uws_sys/App.rs +++ b/src/uws_sys/App.rs @@ -287,8 +287,7 @@ impl App { } } - /// Accept on an already-bound `fd`; it is owned by the listen socket the handler - /// receives, or still by the caller when the handler receives null. + /// Accept on an already-bound `fd`, owned by the listen socket on success and still by the caller on null. pub fn listen_fd( &mut self, handler: c::uws_listen_handler, From 54dae5d14781d4f7ded7fec880d8698a825ea041 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 13 Aug 2026 11:11:54 +0000 Subject: [PATCH 5/5] node:http cluster listen: only unbracket a host that has both brackets --- src/js/node/_http_server.ts | 5 ++++- test/js/node/cluster.test.ts | 8 ++++++++ 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 26371786c1c8..83a411e987fb 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -731,7 +731,10 @@ function listenInCluster(server, tls, port, host, socketPath, onListen) { return; } // Bun.serve() accepts "[::1]"; the primary's bind (and isIP) want it bare. - const address = host.length > 2 && host.charCodeAt(0) === 0x5b /* [ */ ? host.slice(1, -1) : host; + const address = + host.length > 2 && host.charCodeAt(0) === 0x5b /* [ */ && host.charCodeAt(host.length - 1) === 0x5d /* ] */ + ? host.slice(1, -1) + : host; const addressType = isIP(address); if (addressType !== 0) { queryPrimary(listenArgs, address, port, addressType); diff --git a/test/js/node/cluster.test.ts b/test/js/node/cluster.test.ts index 14e64e99b836..a5376c2f698c 100644 --- a/test/js/node/cluster.test.ts +++ b/test/js/node/cluster.test.ts @@ -1745,6 +1745,14 @@ test.skipIf(isWindows || !isIPv6())( 30_000, ); +test.skipIf(isWindows)( + "http worker listen(0, '[::1') is a name lookup that fails, not a bind of ::", + async () => { + expect(await runHttpHostFormFixture("[::1")).toEqual({ error: "ENOTFOUND" }); + }, + 30_000, +); + test.skipIf(isWindows)( "http worker listen(0, 'localhost') resolves the name before asking the primary", async () => {