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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -92,3 +92,4 @@
- Make `process.title` and `--title` visible to OS process tools on Linux and macOS while preserving startup argv and worker-local assignments. Ports [oven-sh/bun#44318](https://github.com/oven-sh/bun/pull/44318). Thanks @tnrich!
- Propagate synchronous `process.emit()` listener exceptions to the caller, preserving once-listener removal and stopping dispatch before later listeners.
- Poll macOS file watches on their owning JavaScript event loop so kqueue coalesces bursts like Node, preserving directory watches, inode replacement, and unref behavior.
- Emit queued HTTP upgrades immediately and defer built-in WebSocket adoption until earlier responses drain. Adapts [oven-sh/bun#43441](https://github.com/oven-sh/bun/pull/43441), preserving the fork's socket backpressure and pause behavior. Thanks @robobun!
2 changes: 2 additions & 0 deletions docs/runtime/nodejs-compat.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,8 @@ On macOS, `fs.copyFile` and `fs.cp` accept regular-file `/dev/fd` sources withou

Destroying an unfinished request or response closes its connection, including asynchronous destruction after `writeHead()`. Bun does not complete the response as an empty success.

Upgrade requests emit `upgrade` while an earlier response on the connection is still pending. Built-in WebSocket adoption waits for that response to drain.

🟢 Fully implemented. `http.Server` does not extend `net.Server`. Bun ignores `listen(handle)` and the `fd`, `ipv6Only` and `signal` options of `listen()`. `keepAlive`/`keepAliveInitialDelay` on the server are no-ops.

### [`node:https`](https://nodejs.org/api/https.html)
Expand Down
2 changes: 2 additions & 0 deletions src/js/internal/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ const kPendingCallbacks = Symbol("pendingCallbacks");
const kRequest = Symbol("request");
// Set on a server socket at the 'connect'/'upgrade' handoff: the native response of that request.
const kHandoffResponse = Symbol("kHandoffResponse");
const kOnHandoffActive = Symbol("kOnHandoffActive");
const kCloseCallback = Symbol("closeCallback");

// node:_http_server registers its pipelined-response machinery here at module
Expand Down Expand Up @@ -545,6 +546,7 @@ export {
kHandoffResponse,
kInternalSocketData,
kNeedDrain,
kOnHandoffActive,
kOnReadParsed,
kOutHeaders,
kPendingCallbacks,
Expand Down
92 changes: 58 additions & 34 deletions src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ const {
kHandle,
kOnReadParsed,
kHandoffResponse,
kOnHandoffActive,
kRealListen,
tlsSymbol,
optionsSymbol,
Expand Down Expand Up @@ -1112,13 +1113,8 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
// token; the server then consults shouldUpgradeCallback (default: an
// 'upgrade' listener is installed) and otherwise dispatches the
// request normally.
// Not when pipelined: the builtin ws answers through the socket's current response, the one in flight.
let is_upgrade = false;
if (
!isPipelined &&
(dispatchBits & DISPATCH_HAS_UPGRADE) !== 0 &&
(dispatchBits & DISPATCH_CONN_UPGRADE) !== 0
) {
if ((dispatchBits & DISPATCH_HAS_UPGRADE) !== 0 && (dispatchBits & DISPATCH_CONN_UPGRADE) !== 0) {
// Like Node.js, shouldUpgradeCallback sees req.upgrade === true.
http_req.upgrade = true;
is_upgrade = !!server.shouldUpgradeCallback(http_req);
Expand All @@ -1134,6 +1130,16 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
// pipeline assigns it the socket (advanceResponsePipeline).
socket[kPipelinedResponses]?.pop(); // the turn that the native handle held
queuePipelinedResponse(socket, http_res, !!isAncientHTTP);
if (is_upgrade) {
http_res[kPipelinedQueuedState].handoff = true;
socket[kPendingHandoff] = [];
socketHandle.upgradeToTunnel(hasBody, handle);
socket[kHandoffResponse] = handle;
socket[kEnableStreaming](true);
kickPipelineIfIdle(server, socket);
emitUpgradeHandoff(server, socket, http_req, !hasBody && connectHead ? connectHead : kEmptyBuffer, hasBody);
return;
}
// A pipelined dispatch can arrive after the previous response finished and detached
// (bytes still flushing keep it pending), leaving nothing in flight to advance the
// queue. Kick the pipeline once this dispatch settles.
Expand Down Expand Up @@ -1230,32 +1236,8 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
socketHandle.upgradeToTunnel(hasBody, handle);
socket[kHandoffResponse] = handle;
socket[kEnableStreaming](true);
detachSocketListenersForHandoff(socket);
// Node frees the parser before emitting 'upgrade' (socket.parser === null there).
releaseServerParserShim(socket, http_req);
if (hasBody) {
socket[kUpgradeIncoming] = http_req;
http_req.once("end", clearUpgradeIncoming.bind(undefined, socket));
} else {
http_req.complete = true;
}
const upgradeHead = !hasBody && connectHead ? connectHead : kEmptyBuffer;
let upgradeHandled;
try {
upgradeHandled = server.emit("upgrade", http_req, socket, upgradeHead);
} catch (err) {
// A throwing 'upgrade' listener surfaces as an uncaught
// exception, like Node.js (the emit happens outside any JS try
// frame there).
process.nextTick(rethrowUncaught, err);
upgradeHandled = true;
}
if (!upgradeHandled) {
// shouldUpgradeCallback accepted the upgrade but no 'upgrade'
// listener is installed: Node.js destroys the socket.
socket.destroy();
return;
}
if (!emitUpgradeHandoff(server, socket, http_req, upgradeHead, hasBody)) return;
// Like CONNECT: the connection is detached from the HTTP request
// machinery; hold the native callback open until the raw socket
// closes.
Expand Down Expand Up @@ -1654,6 +1636,26 @@ function detachSocketListenersForHandoff(socket) {
function resolveHandoffPromise(promise) {
$resolvePromise(promise, undefined);
}
// Emit immediately: an upgrade listener may finish the response ahead of it.
function emitUpgradeHandoff(server, socket, req, head, hasBody) {
detachSocketListenersForHandoff(socket);
releaseServerParserShim(socket, req);
if (hasBody && !req.complete) {
socket[kUpgradeIncoming] = req;
req.once("end", clearUpgradeIncoming.bind(undefined, socket));
} else if (!hasBody) {
req.complete = true;
}
let handled;
try {
handled = server.emit("upgrade", req, socket, head);
} catch (err) {
process.nextTick(rethrowUncaught, err);
handled = true;
}
if (!handled) socket.destroy();
return handled;
}
const kSocketTimeoutTimer = Symbol("socketTimeoutTimer");
const kStreamingEnabled = Symbol("kStreamingEnabled");
// destroySoon() was called: native lets the read that is being parsed reach the request, then closes right behind the FIN.
Expand Down Expand Up @@ -1682,6 +1684,8 @@ const kPipelinedQueuedState = Symbol("kPipelinedQueuedState");
// responses. Reads are paused while it is at or above the high water mark.
const kOutgoingData = Symbol("kOutgoingData");
const kReplayingPipelinedOps = Symbol("kReplayingPipelinedOps");
// Native WebSocket adoption must wait until responses ahead have drained.
const kPendingHandoff = Symbol("kPendingHandoff");
const kStopParsingOnCloseListener = Symbol("kStopParsingOnCloseListener");

// https://github.com/nodejs/node/blob/v26.3.0/lib/_http_server.js (socketOnError)
Expand Down Expand Up @@ -1827,6 +1831,7 @@ function getNodeHTTPServerSocket() {
[kDispatcherCorkDepth] = 0;
[kOnReadParsed] = undefined;
[kHandoffResponse] = undefined;
[kPendingHandoff]: (() => void)[] | undefined = undefined;
[kDestroySoon] = false;
[kHandedOff] = false;
server: Server;
Expand Down Expand Up @@ -2138,6 +2143,7 @@ function getNodeHTTPServerSocket() {
}

_destroy(err, callback) {
this[kPendingHandoff] = undefined;
// Match net.Socket._destroy: #onClose clears this too, but the native close is async and an already-due timer would fire first.
const timer = this[kSocketTimeoutTimer];
if (timer) {
Expand Down Expand Up @@ -2191,6 +2197,12 @@ function getNodeHTTPServerSocket() {
super.destroySoon();
}

[kOnHandoffActive](callback) {
const pending = this[kPendingHandoff];
if (pending !== undefined) pending.push(callback);
else callback();
}

get localAddress() {
return this[kHandle]?.localAddress?.address;
}
Expand All @@ -2211,7 +2223,7 @@ function getNodeHTTPServerSocket() {
const handle = this[kHandle];
// A tunnel reads again: response.resume() below does nothing for it.
if (readStart && this[kStreamingEnabled]) handle?.readStart();
const response = handle?.response;
const response = this[kHandoffResponse] ?? handle?.response;
const upgradeIncoming = this[kUpgradeIncoming];
if (upgradeIncoming) {
// Upgrade with a body: reading the raw socket resumes the request so its
Expand Down Expand Up @@ -2409,7 +2421,7 @@ function getNodeHTTPServerSocket() {
}

get [kInternalSocketData]() {
return this[kHandle]?.response;
return this[kHandoffResponse] ?? this[kHandle]?.response;
}
} as unknown as typeof import("node:net").Socket;
Object.defineProperty(NodeHTTPServerSocket, "name", { value: "Socket" });
Expand Down Expand Up @@ -3021,7 +3033,7 @@ function abortQueuedPipelinedResponses(socket, error = $ERR_STREAM_DESTROYED("wr
for (let i = 0; i < pipelinedLength; i++) {
const queuedRes = pipelined[i];
// A turn that the native handle still holds: the native close path notifies that one.
if (queuedRes[kPipelinedQueuedState] === undefined) continue;
if (queuedRes[kPipelinedQueuedState] === undefined || queuedRes[kPipelinedQueuedState].handoff) continue;
const queuedReq = queuedRes.req;
failQueuedPipelinedWriteCallbacks(queuedRes[kPipelinedQueuedState], error);
if (queuedReq && !queuedReq.destroyed) {
Expand Down Expand Up @@ -3130,6 +3142,14 @@ function advanceResponsePipeline(server, socket) {
return;
}

if (queued.handoff) {
const ready = socket[kPendingHandoff];
socket[kPendingHandoff] = undefined;
// Adoption replaces the native socket; never run it inside the previous response's callback.
if (ready?.length) setImmediate(runHandoffReady, socket, ready);
return;
}

if (res.assignSocket === ServerResponse.prototype.assignSocket) {
assignSocketInternal(res, socket);
} else {
Expand Down Expand Up @@ -3189,6 +3209,10 @@ function advanceResponsePipeline(server, socket) {
}
}

function runHandoffReady(socket, ready) {
for (let i = 0; i < ready.length && !socket.destroyed; i++) ready[i]();
}

function markResponseEndedNT(res) {
res._ended = true;
}
Expand Down
64 changes: 39 additions & 25 deletions src/js/thirdparty/ws.js
Original file line number Diff line number Diff line change
Expand Up @@ -937,7 +937,7 @@ function socketOnError() {
this.destroy();
}

function abortHandshake(socket, code, message, headers) {
function abortHandshake(socket, code, message, headers, req) {
const { STATUS_CODES } = lazyHttp();
message = message || STATUS_CODES[code];
headers = {
Expand All @@ -948,7 +948,9 @@ function abortHandshake(socket, code, message, headers) {
};

// handleUpgrade() was called from a 'request' listener: answer through its ServerResponse.
const response = socket._httpMessage;
let response = socket._httpMessage;
// After an 'upgrade' hand-off behind a pipeline, that response belongs to a request ahead.
if (response && response.req && response.req !== req) response = undefined;
if (response) {
response.writeHead(code, headers);
response.write(message);
Expand Down Expand Up @@ -978,7 +980,7 @@ function abortHandshakeOrEmitwsClientError(server, req, socket, code, message, h

server.emit("wsClientError", err, socket, req);
} else {
abortHandshake(socket, code, message, headers);
abortHandshake(socket, code, message, headers, req);
}
}

Expand Down Expand Up @@ -1544,7 +1546,7 @@ class WebSocketServer extends EventEmitter {
);
}

if (this._state > RUNNING) return abortHandshake(socket, 503);
if (this._state > RUNNING) return abortHandshake(socket, 503, undefined, undefined, request);

const server = socket.server[kBunInternals];

Expand All @@ -1562,26 +1564,38 @@ class WebSocketServer extends EventEmitter {
const headers = ["HTTP/1.1 101 Switching Protocols", "Upgrade: websocket", "Connection: Upgrade"];
this.emit("headers", headers, request);

if (
server.upgrade(req, {
data: ws[kBunInternals],
headers: protocol ? { "sec-websocket-protocol": protocol } : undefined,
})
) {
const clients = this.clients;
if (clients) {
clients.add(ws);
ws.on("close", () => {
clients.delete(ws);

if (this._shouldEmitClose && !clients.size) {
process.nextTick(wsEmitClose, this);
}
});
const upgrade = () => {
// The checks above, again: the responses ahead may have taken a while.
if (!socket.readable || !socket.writable) return socket.destroy();
if (this._state > RUNNING) return abortHandshake(socket, 503, undefined, undefined, request);
if (
server.upgrade(req, {
data: ws[kBunInternals],
headers: protocol ? { "sec-websocket-protocol": protocol } : undefined,
})
) {
const clients = this.clients;
if (clients) {
clients.add(ws);
ws.on("close", () => {
clients.delete(ws);

if (this._shouldEmitClose && !clients.size) {
process.nextTick(wsEmitClose, this);
}
});
}
cb(ws, request);
} else {
abortHandshake(socket, 500, undefined, undefined, request);
}
cb(ws, request);
};
// The native upgrade takes the connection over: not while responses ahead are in flight.
const onHandoffActive = socket[require("internal/http").kOnHandoffActive];
if (onHandoffActive !== undefined) {
onHandoffActive.$call(socket, upgrade);
} else {
abortHandshake(socket, 500);
upgrade();
}
}
/**
Expand Down Expand Up @@ -1631,7 +1645,7 @@ class WebSocketServer extends EventEmitter {
}

if (!this.shouldHandle(req)) {
abortHandshake(socket, 400);
abortHandshake(socket, 400, undefined, undefined, req);
return;
}

Expand Down Expand Up @@ -1665,15 +1679,15 @@ class WebSocketServer extends EventEmitter {
if (this.options.verifyClient.length === 2) {
this.options.verifyClient(info, (verified, code, message, headers) => {
if (!verified) {
return abortHandshake(socket, code || 401, message, headers);
return abortHandshake(socket, code || 401, message, headers, req);
}

this.completeUpgrade(extensions, key, protocols, req, socket, head, cb);
});
return;
}

if (!this.options.verifyClient(info)) return abortHandshake(socket, 401);
if (!this.options.verifyClient(info)) return abortHandshake(socket, 401, undefined, undefined, req);
}

this.completeUpgrade(extensions, key, protocols, req, socket, head, cb);
Expand Down
5 changes: 3 additions & 2 deletions src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -581,11 +581,12 @@ static void endFloodPreventionPause(us_socket_t* socket, uWS::NodeHttpResponseDa
reinterpret_cast<uWS::HttpResponse<SSL>*>(socket)->resume();
}

/* A pipelined CONNECT stays queued so that the connection never counts as idle. No request follows it, so it holds no reads. */
/* A pipelined tunnel stays queued so that the connection never counts as idle. No request follows it, so it holds no reads. */
template<bool SSL>
static bool queuedResponsesHoldReads(uWS::NodeHttpResponseData<SSL>* httpResponseData)
{
return httpResponseData->nodeHttpQueuedPipelinedCount > 0 && !httpResponseData->isConnectRequest;
bool tunnel = httpResponseData->isConnectRequest || (httpResponseData->state & uWS::HttpResponseData<SSL>::HTTP_NODE_TUNNEL_AFTER_BODY);
return httpResponseData->nodeHttpQueuedPipelinedCount > 0 && !tunnel;
}

/* node:http flood prevention, resume half: unsent response bytes and queued responses hold the reads. */
Expand Down
6 changes: 4 additions & 2 deletions src/runtime/server/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1526,8 +1526,10 @@ impl<const SSL: bool, const DEBUG: bool> NewServer<SSL, DEBUG> {
let nhr_flags = nhr.flags.get();
if !nhr_flags.contains(NhrFlags::UPGRADED) {
if let Some(raw) = nhr.reader() {
if !nhr_flags.contains(NhrFlags::REQUEST_HAS_COMPLETED)
&& raw.state().is_response_pending()
// A queued WebSocket adoption still needs the upgrade context after dispatch.
if nhr_flags.contains(NhrFlags::TUNNELED)
|| (!nhr_flags.contains(NhrFlags::REQUEST_HAS_COMPLETED)
&& raw.state().is_response_pending())
{
nhr.set_on_aborted_handler();
}
Expand Down
Loading
Loading