Skip to content
Open
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
6 changes: 6 additions & 0 deletions packages/bun-uws/src/App.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<void(us_listen_socket_t *)> &&handler) {
handler(httpContext ? trackListenSocket(httpContext->listen_fd(sslCtxOrNull(), fd, options)) : nullptr);
return std::move(*this);
}

void setOnSocketClosed(HttpContextData<SSL>::OnSocketClosedCallback onClose) {
httpContext->getSocketContextData()->onSocketClosed = onClose;
}
Expand Down
19 changes: 19 additions & 0 deletions packages/bun-uws/src/HttpContext.h
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
#include <span>
#include <array>
#include <mutex>
#include <cerrno>


extern "C" void Bun__NodeHTTP__onReadsResumable(int ssl, struct us_socket_t *s);
Expand Down Expand Up @@ -1070,6 +1071,24 @@ struct HttpContext {

return socket;
}

/* 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);
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;
}
};

}
5 changes: 3 additions & 2 deletions src/js/internal/cluster/primary.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
195 changes: 153 additions & 42 deletions src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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");
Expand Down Expand Up @@ -93,6 +101,10 @@ const {
const kConnectionsCheckingInterval = Symbol("http.server.connectionsCheckingInterval");
const kTrackedConnections = Symbol("http.server.trackedConnections");
const kHttpAllowHalfOpen = Symbol("http.server.httpAllowHalfOpen");
// 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
// only created on the first request, and emission is gated per-request on the
Expand Down Expand Up @@ -504,6 +516,7 @@ Server.prototype.closeAllConnections = function () {
clearInterval(this[kConnectionsCheckingInterval]);
this.listening = false;

releaseClusterHandle(this);
server.stop(true);
};

Expand All @@ -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]);
Expand All @@ -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;
Expand Down Expand Up @@ -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:
Expand All @@ -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)) {
Expand Down Expand Up @@ -634,6 +655,7 @@ Server.prototype.listen = function () {
onListen = lastArg;
}

const previousBunServer = server[serverSymbol];
try {
// listenInCluster

Expand All @@ -644,53 +666,140 @@ 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);
});
// 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);
// Non-exclusive workers still share a fixed port through SO_REUSEPORT.
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);
emitListenError(server, err, previousBunServer);
}

return this;
};

Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort, onListen) {
// 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);
}

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);
});
}

// 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 };

// 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;
}
if (!host) {
queryPrimary(listenArgs, null, port, 4);
return;
}
// Bun.serve() accepts "[::1]"; the primary's bind (and isIP) want it bare.
Comment thread
claude[bot] marked this conversation as resolved.
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);
} else {
// 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) {
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 round-robin handle despite sharedOnly (a foreign primary): Bun.serve cannot use it.
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;
// 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: no listener took the descriptor, so close it here.
handle.adopted = false;
releaseClusterHandle(server);
}
emitListenError(server, err, previousBunServer);
}
});
}

// 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;
Expand All @@ -708,6 +817,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) {
Expand Down
Loading
Loading