From 60b211c7bfd04b27b6bfe778df734d05523750e6 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 13:15:46 +0000 Subject: [PATCH 01/13] node:http: hand a pipelined Upgrade request to 'upgrade' and keep its writes behind the responses ahead An Upgrade request that arrives while an earlier response on the same connection is still in flight was dispatched to 'request' with req.upgrade === false. shouldUpgradeCallback now runs at dispatch for every Upgrade request. When it accepts, the connection enters tunnel mode at once and the socket goes to 'upgrade' at once, like Node. What the listener writes to the socket waits until the pipeline reaches this request, so the responses ahead keep their place on the wire. A pipelined CONNECT uses the same queue. The builtin ws adopts the socket through that queue as well, since the native upgrade takes the connection over, and the request keeps its WebSocket upgrade context when the listener ends the response ahead during the dispatch. --- src/js/internal/http.ts | 4 + src/js/node/_http_server.ts | 200 ++++++++++++++------ src/js/thirdparty/ws.js | 63 +++--- src/runtime/server/mod.rs | 9 +- test/js/node/http/node-http-with-ws.test.ts | 77 +++++++- test/js/node/http/node-http.test.ts | 193 +++++++++++++++++++ 6 files changed, 462 insertions(+), 84 deletions(-) diff --git a/src/js/internal/http.ts b/src/js/internal/http.ts index 3e2d45c0c0db..60c0d622526c 100644 --- a/src/js/internal/http.ts +++ b/src/js/internal/http.ts @@ -41,6 +41,9 @@ 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"); +// Method of a server socket: runs the callback once the socket handed to 'connect'/'upgrade' is +// the connection's current exchange (at once, or when the responses ahead of it have finished). +const kOnHandoffActive = Symbol("kOnHandoffActive"); const kCloseCallback = Symbol("closeCallback"); // node:_http_server registers its pipelined-response machinery here at module @@ -526,6 +529,7 @@ export { kHandoffResponse, kInternalSocketData, kNeedDrain, + kOnHandoffActive, kOutHeaders, kPendingCallbacks, kPerRequestCheckServerIdentity, diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index f073ba8d8dfe..bb2cb9c715df 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -33,6 +33,7 @@ const { serverSymbol, kHandle, kHandoffResponse, + kOnHandoffActive, kRealListen, tlsSymbol, optionsSymbol, @@ -736,6 +737,26 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort socketParser[kParserOnTimeout] = serverParserShimOnTimeout; const isPipelined = !!isPipelinedDispatch; + socket[kEnableStreaming](false); + + // The builtin ServerResponse consumes its options synchronously, so a + // reusable scratch object avoids one allocation per request. User + // subclasses (options.ServerResponse) might retain options, so they + // keep getting a fresh object. + let http_res; + if (ResponseClass === ServerResponse) { + scratchResponseOptions[kHandle] = handle; + scratchResponseOptions.highWaterMark = socket.writableHighWaterMark; + scratchResponseOptions[kRejectNonStandardBodyWrites] = server.rejectNonStandardBodyWrites; + http_res = new ResponseClass(http_req, scratchResponseOptions); + scratchResponseOptions[kHandle] = undefined; + } else { + http_res = new ResponseClass(http_req, { + [kHandle]: handle, + highWaterMark: socket.writableHighWaterMark, + [kRejectNonStandardBodyWrites]: server.rejectNonStandardBodyWrites, + }); + } // Pipelined or not, like Node.js: the native parser is in tunnel mode from this request on. if (method === "CONNECT") { // Handle CONNECT method for HTTP tunneling/proxy @@ -759,6 +780,9 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort http_req.upgrade = true; // Node frees the parser before handing the raw socket to 'connect'. releaseServerParserShim(socket, http_req); + // Behind responses that are still in flight, the listener gets the socket now (like + // Node) but what it writes waits for its turn in the pipeline. + if (isPipelined) queuePipelinedHandoff(server, socket, http_res, !!isAncientHTTP); try { server.emit("connect", http_req, socket, head); } catch (err) { @@ -779,26 +803,6 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort } return; } - socket[kEnableStreaming](false); - - // The builtin ServerResponse consumes its options synchronously, so a - // reusable scratch object avoids one allocation per request. User - // subclasses (options.ServerResponse) might retain options, so they - // keep getting a fresh object. - let http_res; - if (ResponseClass === ServerResponse) { - scratchResponseOptions[kHandle] = handle; - scratchResponseOptions.highWaterMark = socket.writableHighWaterMark; - scratchResponseOptions[kRejectNonStandardBodyWrites] = server.rejectNonStandardBodyWrites; - http_res = new ResponseClass(http_req, scratchResponseOptions); - scratchResponseOptions[kHandle] = undefined; - } else { - http_res = new ResponseClass(http_req, { - [kHandle]: handle, - highWaterMark: socket.writableHighWaterMark, - [kRejectNonStandardBodyWrites]: server.rejectNonStandardBodyWrites, - }); - } http_res._keepAliveTimeout = server.keepAliveTimeout; // Only stamp the symbol when the server actually set `uniqueHeaders`: // unconditionally adding it (even as undefined) forced a hidden-class @@ -883,13 +887,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) { is_upgrade = !!server.shouldUpgradeCallback(http_req); } // Like Node.js's parserOnIncoming: req.upgrade is true inside the @@ -901,14 +900,20 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // Node.js, this response is queued (res.socket === null) and its // writes are buffered until the in-flight response finishes and the // pipeline assigns it the socket (advanceResponsePipeline). - queuePipelinedResponse(socket, http_res, !!isAncientHTTP); - // 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. - if (socket._httpMessage == null && !socket[kPipelineKickScheduled]) { - socket[kPipelineKickScheduled] = true; - process.nextTick(advancePipelineIfIdleNT, server, socket); + if (is_upgrade) { + // The connection leaves HTTP here: nothing after this request head is + // parsed as a request. The listener gets the socket now, like Node, + // but what it writes waits for its turn in the pipeline, so the 101 + // follows the responses ahead of it on the wire. + socketHandle.upgradeToTunnel(hasBody, handle); + socket[kHandoffResponse] = handle; + socket[kEnableStreaming](true); + queuePipelinedHandoff(server, socket, http_res, !!isAncientHTTP); + emitUpgradeHandoff(server, socket, http_req, !hasBody && connectHead ? connectHead : kEmptyBuffer, hasBody); + return; } + queuePipelinedResponse(socket, http_res, !!isAncientHTTP); + kickPipelineIfIdle(server, socket); // Node's parserOnIncoming stops reading the connection once the bytes // queued on responses that do not own the socket yet reach the // socket's high water mark, so pipelined requests cannot flood it. @@ -976,28 +981,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)); - } 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(); + if (!emitUpgradeHandoff(server, socket, http_req, upgradeHead, hasBody)) { return; } // Like CONNECT: the connection is detached from the HTTP request @@ -1343,6 +1328,37 @@ function detachSocketListenersForHandoff(socket) { function resolveHandoffPromise(promise) { $resolvePromise(promise, undefined); } + +// Hand the raw socket to the 'upgrade' listeners, like Node.js's +// onParserExecuteCommon. Returns false when the socket was destroyed because +// shouldUpgradeCallback accepted the upgrade but no listener is installed. +function emitUpgradeHandoff(server, socket, req, upgradeHead, hasBody) { + detachSocketListenersForHandoff(socket); + // Node frees the parser before emitting 'upgrade' (socket.parser === null there). + releaseServerParserShim(socket, req); + // A read of a request ahead may have set the socket flowing. Flowing with no + // reader drops pushed bytes, so stop the flow (Node: onParserExecuteCommon). + if (socket.readableFlowing === true) socket.readableFlowing = null; + if (hasBody && !req.complete) { + socket[kUpgradeIncoming] = req; + req.once("end", clearUpgradeIncoming.bind(undefined, socket)); + } + let upgradeHandled; + try { + upgradeHandled = server.emit("upgrade", 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) { + socket.destroy(); + return false; + } + return true; +} const kSocketTimeoutTimer = Symbol("socketTimeoutTimer"); const kStreamingEnabled = Symbol("kStreamingEnabled"); // Scratch options object for the builtin ServerResponse (see the dispatcher). @@ -1369,6 +1385,9 @@ 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"); +// On a server socket handed to 'connect'/'upgrade' behind responses still in flight: +// { write, final, ready } parked until the pipeline reaches the hand-off (undefined otherwise). +const kPendingHandoff = Symbol("kPendingHandoff"); const kStopParsingOnCloseListener = Symbol("kStopParsingOnCloseListener"); // Set when the dispatcher already detached a synchronously-finished response, // so the 'finish' listener does not detach/advance the pipeline a second time. @@ -1498,6 +1517,7 @@ function getNodeHTTPServerSocket() { [kHandle]; [kUpgradeIncoming] = undefined; [kHandoffResponse] = undefined; + [kPendingHandoff] = undefined; server: Server; _httpMessage; _secureEstablished = false; @@ -1753,6 +1773,11 @@ function getNodeHTTPServerSocket() { } _final(callback) { + const pending = this[kPendingHandoff]; + if (pending !== undefined) { + pending.final = callback; + return; + } const handle = this[kHandle]; if (!handle) { callback(); @@ -1762,6 +1787,15 @@ function getNodeHTTPServerSocket() { callback(); } + [kOnHandoffActive](callback) { + const pending = this[kPendingHandoff]; + if (pending !== undefined) { + pending.ready.push(callback); + } else { + callback(); + } + } + get localAddress() { return this[kHandle]?.localAddress?.address; } @@ -1942,6 +1976,12 @@ function getNodeHTTPServerSocket() { } _write(_chunk, _encoding, _callback) { + const pending = this[kPendingHandoff]; + if (pending !== undefined) { + // Writable issues one _write at a time: holding its callback holds everything after it. + pending.write = { chunk: _chunk, encoding: _encoding, callback: _callback }; + return; + } const handle = this[kHandle]; let err; this._unrefTimer(); @@ -1978,7 +2018,9 @@ function getNodeHTTPServerSocket() { } get [kInternalSocketData]() { - return this[kHandle]?.response; + // After a 'connect'/'upgrade' hand-off the response of that request, which behind a + // pipeline is not the socket's current response yet. + return this[kHandoffResponse] ?? this[kHandle]?.response; } } as unknown as typeof import("node:net").Socket; Object.defineProperty(NodeHTTPServerSocket, "name", { value: "Socket" }); @@ -2532,10 +2574,46 @@ function queuePipelinedResponse(socket, res, isAncient) { ended: false, isAncient, socket, + // A CONNECT/Upgrade hand-off: its turn in the pipeline releases what the + // listener wrote to the socket instead of replaying a response. + handoff: false, }; (socket[kPipelinedResponses] ??= []).push(res); } +// 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. +function kickPipelineIfIdle(server, socket) { + if (socket._httpMessage == null && !socket[kPipelineKickScheduled]) { + socket[kPipelineKickScheduled] = true; + process.nextTick(advancePipelineIfIdleNT, server, socket); + } +} + +// A CONNECT/Upgrade behind responses still in flight. The listener gets the socket at +// once, like Node, but the socket parks what it writes (_write/_final) and what must +// run once it is the connection's current exchange (kOnHandoffActive) until the +// pipeline reaches this entry, so the responses ahead keep their place on the wire. +function queuePipelinedHandoff(server, socket, res, isAncient) { + queuePipelinedResponse(socket, res, isAncient); + res[kPipelinedQueuedState].handoff = true; + socket[kPendingHandoff] = { write: undefined, final: undefined, ready: [] }; + kickPipelineIfIdle(server, socket); +} + +function activatePipelinedHandoff(socket) { + const pending = socket[kPendingHandoff]; + socket[kPendingHandoff] = undefined; + if (pending === undefined) return; + const write = pending.write; + if (write !== undefined) socket._write(write.chunk, write.encoding, write.callback); + const final = pending.final; + if (final !== undefined) socket._final(final); + const ready = pending.ready; + for (let i = 0; i < ready.length; i++) ready[i](); +} + // When the connection dies with pipelined responses still queued behind the // in-flight one, abort them and their requests, like Node.js's socketOnClose // (abortIncoming). Runs from the native socket's close path and from the @@ -2547,6 +2625,8 @@ function abortQueuedPipelinedResponses(socket) { socket[kPipelinedResponses] = undefined; for (let i = 0; i < pipelinedLength; i++) { const queuedRes = pipelined[i]; + // The request was handed to 'connect'/'upgrade': Node's socketOnClose is gone there. + if (queuedRes[kPipelinedQueuedState]?.handoff) continue; const queuedReq = queuedRes.req; if (queuedReq && !queuedReq.destroyed) { queuedReq[kHandle] = undefined; @@ -2571,7 +2651,8 @@ function advanceResponsePipeline(server, socket) { // the pipeline is mutually exclusive with closing the socket - the queued // responses are aborted by the socket close path instead of being replayed // onto a half-closed connection. - if (!socket || socket.writableEnded || socket.destroyed) { + // (A socket.end() from a 'connect'/'upgrade' listener is parked, not ended yet.) + if (!socket || socket.destroyed || (socket.writableEnded && socket[kPendingHandoff] === undefined)) { return; } const queue = socket[kPipelinedResponses]; @@ -2611,6 +2692,13 @@ function advanceResponsePipeline(server, socket) { return; } + if (queued.handoff) { + // The responses ahead have finished and this request's native response + // is now the socket's current one: the tunnel owns the connection. + activatePipelinedHandoff(socket); + return; + } + if (res.assignSocket === ServerResponse.prototype.assignSocket) { assignSocketInternal(res, socket); } else { diff --git a/src/js/thirdparty/ws.js b/src/js/thirdparty/ws.js index 301d18f6abb7..36c9ae957278 100644 --- a/src/js/thirdparty/ws.js +++ b/src/js/thirdparty/ws.js @@ -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 = { @@ -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); @@ -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); } } @@ -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]; @@ -1562,26 +1564,37 @@ 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 = () => { + if (socket.destroyed) return; + 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); + }; + // node:http hands a socket to 'upgrade' while earlier responses on the connection can still + // be in flight. The native upgrade takes the connection over, so it waits for them. + const onHandoffActive = socket[require("internal/http").kOnHandoffActive]; + if (onHandoffActive !== undefined) { + onHandoffActive.$call(socket, upgrade); } else { - abortHandshake(socket, 500); + upgrade(); } } /** @@ -1631,7 +1644,7 @@ class WebSocketServer extends EventEmitter { } if (!this.shouldHandle(req)) { - abortHandshake(socket, 400); + abortHandshake(socket, 400, undefined, undefined, req); return; } @@ -1665,7 +1678,7 @@ 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); @@ -1673,7 +1686,7 @@ class WebSocketServer extends EventEmitter { 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); diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index da89a189c0f4..b7a861d5ad31 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1498,8 +1498,13 @@ impl NewServer { let nhr_flags = nhr.flags.get(); if !nhr_flags.contains(NhrFlags::UPGRADED) { if let Some(raw) = nhr.raw_response.get() { - if !nhr_flags.contains(NhrFlags::REQUEST_HAS_COMPLETED) - && raw.state().is_response_pending() + // A raw 'upgrade'/'connect' handoff keeps the WebSocket + // upgrade context: behind a pipeline the builtin ws adopts + // the tunnel only once the responses ahead are complete, + // and the listener may have ended them during this dispatch. + if nhr_flags.contains(NhrFlags::TUNNELED) + || (!nhr_flags.contains(NhrFlags::REQUEST_HAS_COMPLETED) + && raw.state().is_response_pending()) { nhr.set_on_aborted_handler(); } diff --git a/test/js/node/http/node-http-with-ws.test.ts b/test/js/node/http/node-http-with-ws.test.ts index 0ab97a4b1d26..8f4cb0ba80c6 100644 --- a/test/js/node/http/node-http-with-ws.test.ts +++ b/test/js/node/http/node-http-with-ws.test.ts @@ -3,7 +3,7 @@ import { bunEnv, bunExe, tls as options } from "harness"; import http from "http"; import https from "https"; import { once } from "node:events"; -import type { AddressInfo } from "node:net"; +import net, { type AddressInfo } from "node:net"; import tls from "tls"; import { WebSocketServer, type WebSocket as WsWebSocket } from "ws"; @@ -162,3 +162,78 @@ describe.concurrent("request handlers run to completion before the callbacks the expect(order).toEqual(["rest of handler", "nextTick", "microtask"]); }); }); + +test.concurrent("a WebSocket upgrade pipelined behind a pending response completes the handshake", async () => { + // The 'upgrade' listener gets the socket once the response ahead has finished, so the + // WebSocketServer answers through the upgrade request's own response and not through the + // response that was still in flight. + const events: string[] = []; + let first: http.ServerResponse | undefined; + const { promise: finished, resolve: onFinished, reject: onFailure } = Promise.withResolvers(); + await using server = http.createServer((req, res) => { + events.push(`request ${req.url}`); + if (req.headers.upgrade !== undefined) { + onFailure(new Error(`dispatched as a request: ${events}`)); + return; + } + first = res; + res.write("first"); + }); + server.on("clientError", onFailure); + const wsServer = new WebSocketServer({ noServer: true }); + server.on("upgrade", (req, socket, head) => { + events.push(`upgrade ${req.url} upgrade=${req.upgrade}`); + wsServer.handleUpgrade(req, socket, head, ws => { + events.push("connection"); + ws.on("message", data => { + events.push(`message ${data}`); + ws.close(); + }); + ws.on("close", onFinished); + }); + // The handshake completes only once the response ahead has finished. + events.push("first end()"); + first!.end(); + }); + server.shouldUpgradeCallback = req => { + events.push(`shouldUpgradeCallback ${req.url}`); + return true; + }; + await once(server.listen(0, "127.0.0.1"), "listening"); + + const port = (server.address() as AddressInfo).port; + const client = net.connect(port, "127.0.0.1"); + try { + let received = Buffer.alloc(0); + client.on("data", chunk => { + received = Buffer.concat([received, chunk]); + if (received.includes("\r\n\r\n", received.indexOf("HTTP/1.1 101"))) { + // A masked text frame "hi" (RFC 6455 5.7), sent once the handshake is complete. + client.write(Buffer.from([0x81, 0x82, 0x01, 0x02, 0x03, 0x04, 0x68 ^ 0x01, 0x69 ^ 0x02])); + client.removeAllListeners("data"); + client.resume(); + } + }); + client.on("error", onFailure); + client.write( + `GET /first HTTP/1.1\r\nHost: localhost:${port}\r\n\r\n` + + `GET /ws HTTP/1.1\r\nHost: localhost:${port}\r\nConnection: Upgrade\r\nUpgrade: websocket\r\nSec-WebSocket-Version: 13\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\n\r\n`, + ); + await finished; + expect(events).toEqual([ + "request /first", + "shouldUpgradeCallback /ws", + "upgrade /ws upgrade=true", + "first end()", + "connection", + "message hi", + ]); + const text = received.toString("latin1"); + // The response ahead is complete before the 101 starts. + expect(text.indexOf("5\r\nfirst\r\n0\r\n\r\n")).toBeGreaterThan(0); + expect(text.indexOf("HTTP/1.1 101")).toBeGreaterThan(text.indexOf("5\r\nfirst\r\n0\r\n\r\n")); + } finally { + client.destroy(); + wsServer.close(); + } +}); diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index 71ad7e6435d6..ccbcf1bdfdd5 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -4668,3 +4668,196 @@ it("connectionListener pauses reads when queued pipelined responses back up", as clientSide.destroy(); serverSide.destroy(); }); + +describe("Upgrade pipelined behind a pending response", () => { + const get = (path: string) => `GET ${path} HTTP/1.1\r\nHost: example.com\r\n\r\n`; + const upgradeRequest = (path: string, extra = "") => + `GET ${path} HTTP/1.1\r\nHost: example.com\r\nConnection: Upgrade\r\nUpgrade: foo\r\n${extra}\r\n`; + const SWITCHING = "HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nUpgrade: foo\r\n\r\n"; + // Drops the header block of each 200 response. The bodies, the 101 and the tunnel bytes stay. + const withoutResponseHeads = (received: string) => + received.replace(/HTTP\/1\.1 200 OK\r\n(?:[^\r\n]+\r\n)*\r\n/g, ""); + + // The client writes `written` in one write. `respond` answers the requests ahead of the + // Upgrade and leaves the first response open. The 'upgrade' listener writes the 101 and then + // ends that response, like the script in the issue. The tunnel answers once, when the client + // ends its side. + async function pipelinedUpgrade(options: { + written: string; + events: string[]; + respond: (req: http.IncomingMessage, res: http.ServerResponse) => void; + }) { + const events: string[] = []; + let first: http.ServerResponse | undefined; + const { promise: handedOff, resolve: onHandoff, reject: onFailure } = Promise.withResolvers(); + await using server = http.createServer((req, res) => { + events.push(`request ${req.method} ${req.url} upgrade=${req.upgrade}`); + if (req.headers.upgrade !== undefined) { + // An Upgrade gets here only when it was not dispatched as 'upgrade'. Both responses + // end, so the events below are compared and nothing waits. + res.end("not a tunnel"); + first!.end(); + onHandoff(); + return; + } + first ??= res; + options.respond(req, res); + }); + server.on("clientError", onFailure); + server.shouldUpgradeCallback = req => { + events.push(`shouldUpgradeCallback ${req.url}`); + return true; + }; + server.on("upgrade", (req, socket, head) => { + events.push(`upgrade ${req.url} upgrade=${req.upgrade} listeners=${socket.listenerCount("error")}`); + socket.write(SWITCHING); + first!.end(); + let body = ""; + req.on("data", chunk => (body += chunk)); + let tunneled = head.toString(); + socket.on("data", chunk => (tunneled += chunk)); + socket.on("end", () => socket.end(`body:${body};tunneled:${tunneled}`)); + onHandoff(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + const received: Buffer[] = []; + client.on("data", chunk => received.push(chunk)); + client.on("error", onFailure); + client.on("close", () => onFailure(new Error(`closed before the handoff: ${events}`))); + client.write(options.written); + await handedOff; + expect(events).toEqual(options.events); + // The tunnel is not an idle HTTP connection, also when the response ahead of it is complete. + server.closeIdleConnections(); + client.end("later"); + await once(client, "close"); + return Buffer.concat(received).toString("latin1"); + } finally { + client.destroy(); + } + } + + test("should emit 'upgrade' at once and write the 101 after the response ahead", async () => { + const received = await pipelinedUpgrade({ + written: get("/first") + upgradeRequest("/second") + "head,", + events: [ + "request GET /first upgrade=false", + "shouldUpgradeCallback /second", + "upgrade /second upgrade=true listeners=0", + ], + respond: (req, res) => void res.write("first"), + }); + // The 101 follows the response ahead on the wire. Node v26 writes it first. + expect(withoutResponseHeads(received)).toBe(`5\r\nfirst\r\n0\r\n\r\n${SWITCHING}body:;tunneled:head,later`); + }); + + test("should keep the order of the responses that are queued ahead of it", async () => { + const received = await pipelinedUpgrade({ + written: get("/first") + get("/second") + upgradeRequest("/third"), + events: [ + "request GET /first upgrade=false", + "request GET /second upgrade=false", + "shouldUpgradeCallback /third", + "upgrade /third upgrade=true listeners=0", + ], + respond: (req, res) => void (req.url === "/first" ? res.write("first") : res.end("second")), + }); + expect(withoutResponseHeads(received)).toBe(`5\r\nfirst\r\n0\r\n\r\nsecond${SWITCHING}body:;tunneled:later`); + }); + + test("should deliver the body of the Upgrade request through req and tunnel what follows", async () => { + const body = "0123456789"; + const received = await pipelinedUpgrade({ + written: get("/first") + upgradeRequest("/second", `Content-Length: ${body.length}\r\n`) + body + "head,", + events: [ + "request GET /first upgrade=false", + "shouldUpgradeCallback /second", + "upgrade /second upgrade=true listeners=0", + ], + respond: (req, res) => void res.write("first"), + }); + expect(withoutResponseHeads(received)).toBe(`5\r\nfirst\r\n0\r\n\r\n${SWITCHING}body:${body};tunneled:head,later`); + }); + + test("should hold socket.end() from the listener until the response ahead has finished", async () => { + // The listener answers and ends the tunnel before the response ahead is complete. + const events: string[] = []; + let first: http.ServerResponse | undefined; + const { promise: handedOff, resolve: onHandoff, reject: onFailure } = Promise.withResolvers(); + await using server = http.createServer((req, res) => { + events.push(`request ${req.url}`); + if (req.headers.upgrade !== undefined) return void onFailure(new Error(`dispatched as a request: ${events}`)); + first = res; + res.write("first"); + }); + server.on("clientError", onFailure); + server.on("upgrade", (req, socket, head) => { + events.push(`upgrade ${req.url} head=${head}`); + socket.end(`${SWITCHING}bye`); + onHandoff(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + const received: Buffer[] = []; + client.on("data", chunk => received.push(chunk)); + client.on("error", onFailure); + client.write(get("/first") + upgradeRequest("/second") + "head"); + await handedOff; + expect(events).toEqual(["request /first", "upgrade /second head=head"]); + first!.end("-done"); + await once(client, "end"); + expect(withoutResponseHeads(Buffer.concat(received).toString("latin1"))).toBe( + `5\r\nfirst\r\n5\r\n-done\r\n0\r\n\r\n${SWITCHING}bye`, + ); + } finally { + client.destroy(); + } + }); + + test("should go to 'request' when shouldUpgradeCallback declines", async () => { + const events: string[] = []; + let first: http.ServerResponse | undefined; + await using server = http.createServer((req, res) => { + events.push(`request ${req.url} upgrade=${req.upgrade}`); + if (req.url === "/first") { + first = res; + res.write("first"); + } else { + res.end("second"); + first!.end(); + } + }); + server.shouldUpgradeCallback = req => { + events.push(`shouldUpgradeCallback ${req.url}`); + return false; + }; + server.on("upgrade", () => events.push("upgrade")); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + const received: Buffer[] = []; + client.on("data", chunk => received.push(chunk)); + client.write( + get("/first") + + "GET /second HTTP/1.1\r\nHost: example.com\r\nConnection: Upgrade, close\r\nUpgrade: foo\r\n\r\n", + ); + await once(client, "close"); + expect(events).toEqual([ + "request /first upgrade=false", + "shouldUpgradeCallback /second", + "request /second upgrade=false", + ]); + expect( + Buffer.concat(received) + .toString("latin1") + .match(/HTTP\/1\.1 200 OK/g), + ).toHaveLength(2); + } finally { + client.destroy(); + } + }); +}); From 347bf43a31e91b9d4d2783d233df713e44ac9f9e Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 13:39:28 +0000 Subject: [PATCH 02/13] ci: retrigger From ffbd1ca36a6ed9c9bba535fd668431ef5f78e1d3 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 13:40:31 +0000 Subject: [PATCH 03/13] node:http: shorten the kOnHandoffActive comment --- src/js/internal/http.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/src/js/internal/http.ts b/src/js/internal/http.ts index 60c0d622526c..11ed51cc456f 100644 --- a/src/js/internal/http.ts +++ b/src/js/internal/http.ts @@ -41,8 +41,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"); -// Method of a server socket: runs the callback once the socket handed to 'connect'/'upgrade' is -// the connection's current exchange (at once, or when the responses ahead of it have finished). +// Server socket method: run a callback once the 'connect'/'upgrade' hand-off owns the connection. const kOnHandoffActive = Symbol("kOnHandoffActive"); const kCloseCallback = Symbol("closeCallback"); From 057d0891a64eb4a2ac28efe8b18e0d1312ac207d Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 13:43:33 +0000 Subject: [PATCH 04/13] node:http: shorten the comments of the pipelined hand-off --- src/js/node/_http_server.ts | 36 ++++++++++-------------------------- src/js/thirdparty/ws.js | 3 +-- src/runtime/server/mod.rs | 6 ++---- 3 files changed, 13 insertions(+), 32 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index bb2cb9c715df..4ee2e3edd71c 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -780,8 +780,7 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort http_req.upgrade = true; // Node frees the parser before handing the raw socket to 'connect'. releaseServerParserShim(socket, http_req); - // Behind responses that are still in flight, the listener gets the socket now (like - // Node) but what it writes waits for its turn in the pipeline. + // Behind responses in flight, the listener's writes wait for their turn in the pipeline. if (isPipelined) queuePipelinedHandoff(server, socket, http_res, !!isAncientHTTP); try { server.emit("connect", http_req, socket, head); @@ -901,10 +900,8 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // writes are buffered until the in-flight response finishes and the // pipeline assigns it the socket (advanceResponsePipeline). if (is_upgrade) { - // The connection leaves HTTP here: nothing after this request head is - // parsed as a request. The listener gets the socket now, like Node, - // but what it writes waits for its turn in the pipeline, so the 101 - // follows the responses ahead of it on the wire. + // The listener gets the socket now, like Node. Its writes wait for their turn in the + // pipeline, so the 101 follows the responses ahead of it on the wire. socketHandle.upgradeToTunnel(hasBody, handle); socket[kHandoffResponse] = handle; socket[kEnableStreaming](true); @@ -1329,15 +1326,12 @@ function resolveHandoffPromise(promise) { $resolvePromise(promise, undefined); } -// Hand the raw socket to the 'upgrade' listeners, like Node.js's -// onParserExecuteCommon. Returns false when the socket was destroyed because -// shouldUpgradeCallback accepted the upgrade but no listener is installed. +// Node's onParserExecuteCommon. False: no 'upgrade' listener, so the socket was destroyed. function emitUpgradeHandoff(server, socket, req, upgradeHead, hasBody) { detachSocketListenersForHandoff(socket); // Node frees the parser before emitting 'upgrade' (socket.parser === null there). releaseServerParserShim(socket, req); - // A read of a request ahead may have set the socket flowing. Flowing with no - // reader drops pushed bytes, so stop the flow (Node: onParserExecuteCommon). + // A flowing socket with no reader drops pushed bytes (Node does the same reset). if (socket.readableFlowing === true) socket.readableFlowing = null; if (hasBody && !req.complete) { socket[kUpgradeIncoming] = req; @@ -1385,8 +1379,7 @@ 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"); -// On a server socket handed to 'connect'/'upgrade' behind responses still in flight: -// { write, final, ready } parked until the pipeline reaches the hand-off (undefined otherwise). +// { write, final, ready } parked on a handed-off socket until the pipeline reaches it. const kPendingHandoff = Symbol("kPendingHandoff"); const kStopParsingOnCloseListener = Symbol("kStopParsingOnCloseListener"); // Set when the dispatcher already detached a synchronously-finished response, @@ -2018,8 +2011,7 @@ function getNodeHTTPServerSocket() { } get [kInternalSocketData]() { - // After a 'connect'/'upgrade' hand-off the response of that request, which behind a - // pipeline is not the socket's current response yet. + // After a hand-off: that request's response, not the one still in flight ahead of it. return this[kHandoffResponse] ?? this[kHandle]?.response; } } as unknown as typeof import("node:net").Socket; @@ -2574,16 +2566,13 @@ function queuePipelinedResponse(socket, res, isAncient) { ended: false, isAncient, socket, - // A CONNECT/Upgrade hand-off: its turn in the pipeline releases what the - // listener wrote to the socket instead of replaying a response. + // A CONNECT/Upgrade hand-off: its turn releases the listener's parked writes. handoff: false, }; (socket[kPipelinedResponses] ??= []).push(res); } -// 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. +// Nothing in flight advances the queue when the previous response finished but still flushes. function kickPipelineIfIdle(server, socket) { if (socket._httpMessage == null && !socket[kPipelineKickScheduled]) { socket[kPipelineKickScheduled] = true; @@ -2591,10 +2580,7 @@ function kickPipelineIfIdle(server, socket) { } } -// A CONNECT/Upgrade behind responses still in flight. The listener gets the socket at -// once, like Node, but the socket parks what it writes (_write/_final) and what must -// run once it is the connection's current exchange (kOnHandoffActive) until the -// pipeline reaches this entry, so the responses ahead keep their place on the wire. +// The socket parks its writes and kOnHandoffActive callbacks until the pipeline reaches it. function queuePipelinedHandoff(server, socket, res, isAncient) { queuePipelinedResponse(socket, res, isAncient); res[kPipelinedQueuedState].handoff = true; @@ -2693,8 +2679,6 @@ function advanceResponsePipeline(server, socket) { } if (queued.handoff) { - // The responses ahead have finished and this request's native response - // is now the socket's current one: the tunnel owns the connection. activatePipelinedHandoff(socket); return; } diff --git a/src/js/thirdparty/ws.js b/src/js/thirdparty/ws.js index 36c9ae957278..90312f10e8eb 100644 --- a/src/js/thirdparty/ws.js +++ b/src/js/thirdparty/ws.js @@ -1588,8 +1588,7 @@ class WebSocketServer extends EventEmitter { abortHandshake(socket, 500, undefined, undefined, request); } }; - // node:http hands a socket to 'upgrade' while earlier responses on the connection can still - // be in flight. The native upgrade takes the connection over, so it waits for them. + // 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); diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index b7a861d5ad31..53ad6e6ac65d 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1498,10 +1498,8 @@ impl NewServer { let nhr_flags = nhr.flags.get(); if !nhr_flags.contains(NhrFlags::UPGRADED) { if let Some(raw) = nhr.raw_response.get() { - // A raw 'upgrade'/'connect' handoff keeps the WebSocket - // upgrade context: behind a pipeline the builtin ws adopts - // the tunnel only once the responses ahead are complete, - // and the listener may have ended them during this dispatch. + // A tunnel keeps its WebSocket upgrade context: ws adopts + // it only once the responses ahead of it are complete. if nhr_flags.contains(NhrFlags::TUNNELED) || (!nhr_flags.contains(NhrFlags::REQUEST_HAS_COMPLETED) && raw.state().is_response_pending()) From 1799b18724aa4870e087feb2cd6d010b23784828 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 13:46:53 +0000 Subject: [PATCH 05/13] node:http: one-line comments for the pipelined hand-off --- src/js/node/_http_server.ts | 3 +-- src/runtime/server/mod.rs | 3 +-- 2 files changed, 2 insertions(+), 4 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 4ee2e3edd71c..90ee409e2aae 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -900,8 +900,7 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // writes are buffered until the in-flight response finishes and the // pipeline assigns it the socket (advanceResponsePipeline). if (is_upgrade) { - // The listener gets the socket now, like Node. Its writes wait for their turn in the - // pipeline, so the 101 follows the responses ahead of it on the wire. + // The listener gets the socket now, like Node; its writes follow the responses ahead. socketHandle.upgradeToTunnel(hasBody, handle); socket[kHandoffResponse] = handle; socket[kEnableStreaming](true); diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index 53ad6e6ac65d..ec470070d88c 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1498,8 +1498,7 @@ impl NewServer { let nhr_flags = nhr.flags.get(); if !nhr_flags.contains(NhrFlags::UPGRADED) { if let Some(raw) = nhr.raw_response.get() { - // A tunnel keeps its WebSocket upgrade context: ws adopts - // it only once the responses ahead of it are complete. + // A tunnel keeps its upgrade context: ws adopts it later. if nhr_flags.contains(NhrFlags::TUNNELED) || (!nhr_flags.contains(NhrFlags::REQUEST_HAS_COMPLETED) && raw.state().is_response_pending()) From 74feb6340200b22695bc30075a0ca3107934d89f Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 16:53:44 +0000 Subject: [PATCH 06/13] node:http: order a tunnel's writes behind the HTTP output ahead of it A CONNECT or Upgrade socket wrote straight to the connection. Behind a response whose tail was still in the uWS send buffer or in a zero-copy write, those bytes landed inside that response, and a socket.end() lost its FIN. The tunnel's bytes now wait in its stream buffer while the connection owes HTTP output, and HttpContext::onWritable flushes them once the responses ahead are done. Also from the review of the pipelined hand-off: uncork the socket at the hand-off, fail parked write callbacks when the socket is destroyed, let pause/resume target the handed-off request's response, recheck the WebSocketServer state in the deferred ws upgrade, and build no ServerResponse for a CONNECT that is not pipelined. --- packages/bun-uws/src/HttpContext.h | 14 ++- src/js/node/_http_server.ts | 83 +++++++++++------ src/js/thirdparty/ws.js | 4 +- .../bindings/node/JSNodeHTTPServerSocket.cpp | 29 +++++- .../bindings/node/JSNodeHTTPServerSocket.h | 6 ++ .../node/JSNodeHTTPServerSocketPrototype.cpp | 12 ++- src/runtime/socket/uws_jsc.rs | 8 ++ test/js/node/http/node-http-connect.test.ts | 41 +++++++++ test/js/node/http/node-http.test.ts | 92 +++++++++++++++++++ 9 files changed, 250 insertions(+), 39 deletions(-) diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index 6978d65dbe63..732eab22953a 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -820,10 +820,6 @@ struct HttpContext { auto *httpContextData = getSocketContextDataS(s); - - if (httpResponseData->isConnectRequest && httpResponseData->socketData && httpContextData->onSocketDrain) { - httpContextData->onSocketDrain(httpResponseData->socketData, SSL, (struct us_socket_t *) s); - } /* Ask the developer to write data and return success (true) or failure (false), OR skip sending anything and return success (true). */ if (httpResponseData->onWritable) { /* We are now writable, so hang timeout again, the user does not have to do anything so we should hang until end or tryEnd rearms timeout */ @@ -862,6 +858,16 @@ struct HttpContext { /* Drain any socket buffer, this might empty our backpressure and thus finish the request */ asyncSocket->flush(); + /* node:http compat: a CONNECT/Upgrade tunnel's own writes (JSNodeHTTPServerSocket's + * stream buffer) go out once the responses ahead of it owe nothing more. */ + if constexpr (IsNodeHttp) { + bool tunnel = httpResponseData->isConnectRequest || (httpResponseData->state & HttpResponseData::HTTP_NODE_TUNNEL_AFTER_BODY); + if (tunnel && httpResponseData->socketData && httpContextData->onSocketDrain + && httpResponseData->onWritable == nullptr && asyncSocket->getBufferedAmount() == 0) { + httpContextData->onSocketDrain(httpResponseData->socketData, SSL, (struct us_socket_t *) s); + } + } + /* node:http compat: reads were paused while pipelined responses were * queued and stayed paused because the socket still had outgoing * backpressure when the queue drained; now that it has flushed, read diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 90ee409e2aae..3b9163b9ace9 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -737,26 +737,6 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort socketParser[kParserOnTimeout] = serverParserShimOnTimeout; const isPipelined = !!isPipelinedDispatch; - socket[kEnableStreaming](false); - - // The builtin ServerResponse consumes its options synchronously, so a - // reusable scratch object avoids one allocation per request. User - // subclasses (options.ServerResponse) might retain options, so they - // keep getting a fresh object. - let http_res; - if (ResponseClass === ServerResponse) { - scratchResponseOptions[kHandle] = handle; - scratchResponseOptions.highWaterMark = socket.writableHighWaterMark; - scratchResponseOptions[kRejectNonStandardBodyWrites] = server.rejectNonStandardBodyWrites; - http_res = new ResponseClass(http_req, scratchResponseOptions); - scratchResponseOptions[kHandle] = undefined; - } else { - http_res = new ResponseClass(http_req, { - [kHandle]: handle, - highWaterMark: socket.writableHighWaterMark, - [kRejectNonStandardBodyWrites]: server.rejectNonStandardBodyWrites, - }); - } // Pipelined or not, like Node.js: the native parser is in tunnel mode from this request on. if (method === "CONNECT") { // Handle CONNECT method for HTTP tunneling/proxy @@ -770,6 +750,15 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // readable side without tearing the tunnel down (allowHalfOpen). socketHandle.upgradeToTunnel(false, handle); socket[kHandoffResponse] = handle; + // Behind responses in flight, the listener's writes wait for their turn in the pipeline. + if (isPipelined) { + queuePipelinedHandoff( + server, + socket, + newServerResponse(ResponseClass, server, http_req, handle, socket), + !!isAncientHTTP, + ); + } // The parser is detached: the socket is handed over with only // net.Socket's 'end' listener left, like Node.js. detachSocketListenersForHandoff(socket); @@ -780,8 +769,6 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort http_req.upgrade = true; // Node frees the parser before handing the raw socket to 'connect'. releaseServerParserShim(socket, http_req); - // Behind responses in flight, the listener's writes wait for their turn in the pipeline. - if (isPipelined) queuePipelinedHandoff(server, socket, http_res, !!isAncientHTTP); try { server.emit("connect", http_req, socket, head); } catch (err) { @@ -802,6 +789,13 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort } return; } + socket[kEnableStreaming](false); + + // The builtin ServerResponse consumes its options synchronously, so a + // reusable scratch object avoids one allocation per request. User + // subclasses (options.ServerResponse) might retain options, so they + // keep getting a fresh object. + const http_res = newServerResponse(ResponseClass, server, http_req, handle, socket); http_res._keepAliveTimeout = server.keepAliveTimeout; // Only stamp the symbol when the server actually set `uniqueHeaders`: // unconditionally adding it (even as undefined) forced a hidden-class @@ -1320,6 +1314,14 @@ function detachSocketListenersForHandoff(socket) { socket.removeListener("error", socketOnError); socket.removeListener("timeout", onNodeHTTPServerSocketTimeout); socket.on("end", onReadableStreamEnd); + // The dispatcher corked the socket after the requests ahead; the listener owns it now. + const writableState = socket._writableState; + if (writableState?.corked) { + writableState.corked = 1; + socket.uncork(); + } + // Node's onParserExecuteCommon: the listener's 'data' handler starts the flow. + socket.readableFlowing = null; } function resolveHandoffPromise(promise) { $resolvePromise(promise, undefined); @@ -1330,8 +1332,6 @@ function emitUpgradeHandoff(server, socket, req, upgradeHead, hasBody) { detachSocketListenersForHandoff(socket); // Node frees the parser before emitting 'upgrade' (socket.parser === null there). releaseServerParserShim(socket, req); - // A flowing socket with no reader drops pushed bytes (Node does the same reset). - if (socket.readableFlowing === true) socket.readableFlowing = null; if (hasBody && !req.complete) { socket[kUpgradeIncoming] = req; req.once("end", clearUpgradeIncoming.bind(undefined, socket)); @@ -1360,6 +1360,21 @@ const scratchResponseOptions = { highWaterMark: 0, [kRejectNonStandardBodyWrites]: false, }; +function newServerResponse(ResponseClass, server, req, handle, socket) { + if (ResponseClass === ServerResponse) { + scratchResponseOptions[kHandle] = handle; + scratchResponseOptions.highWaterMark = socket.writableHighWaterMark; + scratchResponseOptions[kRejectNonStandardBodyWrites] = server.rejectNonStandardBodyWrites; + const res = new ResponseClass(req, scratchResponseOptions); + scratchResponseOptions[kHandle] = undefined; + return res; + } + return new ResponseClass(req, { + [kHandle]: handle, + highWaterMark: socket.writableHighWaterMark, + [kRejectNonStandardBodyWrites]: server.rejectNonStandardBodyWrites, + }); +} // Per-socket cached bound abort handler (the socket outlives its requests). const kBoundOnAbort = Symbol("kBoundOnAbort"); const kKeepAliveTimeoutSet = Symbol("keepAliveTimeoutSet"); @@ -1745,6 +1760,13 @@ function getNodeHTTPServerSocket() { } _destroy(err, callback) { + const pending = this[kPendingHandoff]; + if (pending !== undefined) { + this[kPendingHandoff] = undefined; + const reason = err ?? $ERR_STREAM_DESTROYED("write"); + pending.write?.callback(reason); + pending.final?.(reason); + } const handle = this[kHandle]; if (!handle) { if ($isCallable(callback)) callback(err); @@ -1806,7 +1828,8 @@ function getNodeHTTPServerSocket() { #resumeSocket() { const handle = this[kHandle]; - const response = handle?.response; + // Behind a pipeline, the handed-off request's response is not the socket's current one yet. + 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 @@ -1996,7 +2019,7 @@ function getNodeHTTPServerSocket() { pause() { const handle = this[kHandle]; - const response = handle?.response; + const response = this[kHandoffResponse] ?? handle?.response; if (response) { response.pause(); } @@ -2636,8 +2659,12 @@ function advanceResponsePipeline(server, socket) { // the pipeline is mutually exclusive with closing the socket - the queued // responses are aborted by the socket close path instead of being replayed // onto a half-closed connection. - // (A socket.end() from a 'connect'/'upgrade' listener is parked, not ended yet.) - if (!socket || socket.destroyed || (socket.writableEnded && socket[kPendingHandoff] === undefined)) { + if (!socket || socket.destroyed) { + return; + } + if (socket.writableEnded) { + // The end is parked behind the responses ahead; let a hand-off's writes go out before the FIN. + activatePipelinedHandoff(socket); return; } const queue = socket[kPipelinedResponses]; diff --git a/src/js/thirdparty/ws.js b/src/js/thirdparty/ws.js index 90312f10e8eb..911c3a3108ad 100644 --- a/src/js/thirdparty/ws.js +++ b/src/js/thirdparty/ws.js @@ -1565,7 +1565,9 @@ class WebSocketServer extends EventEmitter { this.emit("headers", headers, request); const upgrade = () => { - if (socket.destroyed) return; + // 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], diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp index c830813cc4b9..c0b737d9fb77 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp @@ -20,7 +20,7 @@ extern "C" void Bun__NodeHTTPResponse_onClose(void* zigResponse, JSC::EncodedJSV extern "C" void us_socket_free_stream_buffer(us_socket_stream_buffer_t* streamBuffer); extern "C" uint64_t uws_res_get_remote_address_info(void* res, const char** dest, int* port, bool* is_ipv6); extern "C" uint64_t uws_res_get_local_address_info(void* res, const char** dest, int* port, bool* is_ipv6); -extern "C" EncodedJSValue us_socket_buffered_js_write(void* socket, bool is_ssl, bool ended, us_socket_stream_buffer_t* streamBuffer, JSC::JSGlobalObject* globalObject, JSC::EncodedJSValue data, JSC::EncodedJSValue encoding); +extern "C" EncodedJSValue us_socket_buffered_js_write(void* socket, bool is_ssl, bool ended, bool hold, us_socket_stream_buffer_t* streamBuffer, JSC::JSGlobalObject* globalObject, JSC::EncodedJSValue data, JSC::EncodedJSValue encoding); extern "C" int us_socket_is_ssl_handshake_finished(struct us_socket_t* s); extern "C" int us_socket_ssl_handshake_callback_has_fired(struct us_socket_t* s); @@ -274,6 +274,28 @@ static bool deferShutdownUntilResponseDrains(us_socket_t* socket) return true; } +template +static bool tunnelOwesHttpOutputImpl(us_socket_t* socket) +{ + auto* httpResponseData = reinterpret_cast*>(us_socket_ext(socket)); + bool tunnel = httpResponseData->isConnectRequest || (httpResponseData->state & uWS::HttpResponseData::HTTP_NODE_TUNNEL_AFTER_BODY); + if (!tunnel) { + return false; + } + return reinterpret_cast*>(socket)->getBufferedAmount() > 0 || httpResponseData->onWritable != nullptr; +} + +bool JSNodeHTTPServerSocket::tunnelOwesHttpOutput() +{ + if (!socket || upgraded || us_socket_is_closed(socket)) { + return false; + } + if (is_ssl) { + return tunnelOwesHttpOutputImpl(socket); + } + return tunnelOwesHttpOutputImpl(socket); +} + bool JSNodeHTTPServerSocket::shutdownAfterResponseDrains() { if (!socket || upgraded || us_socket_is_closed(socket) || us_socket_is_shut_down(socket)) { @@ -731,10 +753,11 @@ void JSNodeHTTPServerSocket::onDrain() } auto bufferedSize = this->streamBuffer.bufferedSize(); - if (bufferedSize > 0) { + /* Also a socket.end() that waited for the HTTP output ahead of it: nothing left to write, FIN now. */ + if (bufferedSize > 0 || (this->ended && this->socket && !us_socket_is_shut_down(this->socket))) { auto* globalObject = defaultGlobalObject(this->globalObject()); auto scope = DECLARE_TOP_EXCEPTION_SCOPE(globalObject->vm()); - us_socket_buffered_js_write(this->socket, this->is_ssl, this->ended, &this->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); + us_socket_buffered_js_write(this->socket, this->is_ssl, this->ended, false, &this->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); if (auto* exception = scope.exception()) { (void)scope.tryClearException(); globalObject->reportUncaughtExceptionAtEventLoop(globalObject, exception); diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h index df304c66e63d..e69cb9c577c6 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h @@ -108,6 +108,12 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject { * truncate the response. Returns true after handing the close to uWS. */ bool shutdownAfterResponseDrains(); + /* A CONNECT/Upgrade tunnel whose connection still owes HTTP output: the + * send buffer of the responses ahead is not empty, or one of them is in + * the middle of a zero-copy write. The tunnel's bytes wait behind them + * (HttpContext::onWritable flushes the stream buffer once they are out). */ + bool tunnelOwesHttpOutput(); + /* Switch the connection into CONNECT-style tunnel mode after an accepted * Upgrade: subsequent bytes bypass the HTTP parser and stream to the * ondata callback as opaque data. With afterBody, the switch is deferred diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp index f8a4c450720a..d6b280db3a93 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp @@ -9,7 +9,7 @@ #include #include -extern "C" EncodedJSValue us_socket_buffered_js_write(void* socket, bool is_ssl, bool ended, us_socket_stream_buffer_t* streamBuffer, JSC::JSGlobalObject* globalObject, JSC::EncodedJSValue data, JSC::EncodedJSValue encoding); +extern "C" EncodedJSValue us_socket_buffered_js_write(void* socket, bool is_ssl, bool ended, bool hold, us_socket_stream_buffer_t* streamBuffer, JSC::JSGlobalObject* globalObject, JSC::EncodedJSValue data, JSC::EncodedJSValue encoding); extern "C" uint64_t uws_res_get_remote_address_info(void* res, const char** dest, int* port, bool* is_ipv6); extern "C" uint64_t uws_res_get_local_address_info(void* res, const char** dest, int* port, bool* is_ipv6); extern "C" void us_socket_resume(us_socket_t*); @@ -213,7 +213,9 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketWrite, (JSC::JSGlobalObje return JSValue::encode(JSC::jsNumber(0)); } - return us_socket_buffered_js_write(thisObject->socket, thisObject->is_ssl, thisObject->ended, &thisObject->streamBuffer, globalObject, JSValue::encode(callFrame->argument(0)), JSValue::encode(callFrame->argument(1))); + // Behind the HTTP output of the responses ahead, the bytes wait in the stream buffer. + bool hold = thisObject->tunnelOwesHttpOutput(); + return us_socket_buffered_js_write(thisObject->socket, thisObject->is_ssl, thisObject->ended, hold, &thisObject->streamBuffer, globalObject, JSValue::encode(callFrame->argument(0)), JSValue::encode(callFrame->argument(1))); } JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketEnd, (JSC::JSGlobalObject * globalObject, JSC::CallFrame* callFrame)) @@ -227,6 +229,10 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketEnd, (JSC::JSGlobalObject } thisObject->ended = true; + // A tunnel's FIN follows the HTTP output ahead of it (JSNodeHTTPServerSocket::onDrain). + if (thisObject->tunnelOwesHttpOutput()) { + return JSValue::encode(JSC::jsUndefined()); + } // The response's buffered body must reach the kernel before the FIN; uWS // performs the shutdown after its send buffer drains. if (thisObject->shutdownAfterResponseDrains()) { @@ -240,7 +246,7 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketEnd, (JSC::JSGlobalObject if (thisObject->socket && !thisObject->upgraded) { us_socket_pause(thisObject->socket); } - auto result = us_socket_buffered_js_write(thisObject->socket, thisObject->is_ssl, thisObject->ended, &thisObject->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); + auto result = us_socket_buffered_js_write(thisObject->socket, thisObject->is_ssl, thisObject->ended, false, &thisObject->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); // Undo the pause above after the shutdown so the unread body drains // and kqueue's one-shot EVFILT_WRITE (which delivers EV_EOF on // SHUT_WR) is not deleted by a W -> R|W -> R step. diff --git a/src/runtime/socket/uws_jsc.rs b/src/runtime/socket/uws_jsc.rs index b23a41ac95f8..d99350bc645a 100644 --- a/src/runtime/socket/uws_jsc.rs +++ b/src/runtime/socket/uws_jsc.rs @@ -134,6 +134,8 @@ unsafe extern "C" fn us_socket_buffered_js_write( // kept for ABI parity with the C++ caller; TLS is now per-socket _ssl: bool, ended: bool, + // Only append to the stream buffer: the socket still owes bytes that go first. + hold: bool, buffer: *mut us_socket_stream_buffer_t, global_object: &JSGlobalObject, data: JSValue, @@ -201,6 +203,12 @@ unsafe extern "C" fn us_socket_buffered_js_write( // single `&mut` does not alias the re-entrant write path documented at // the top of this fn (raw `socket` is still kept for that reason). let socket_ref = us_socket_t::opaque_mut(socket); + if hold { + if !data_slice.is_empty() { + stream_buffer.write(data_slice); + } + break 'body JSValue::FALSE; + } if stream_buffer.is_not_empty() { let to_flush = stream_buffer.slice(); let to_flush_len = to_flush.len(); diff --git a/test/js/node/http/node-http-connect.test.ts b/test/js/node/http/node-http-connect.test.ts index f242bb980c31..def60405b5a0 100644 --- a/test/js/node/http/node-http-connect.test.ts +++ b/test/js/node/http/node-http-connect.test.ts @@ -946,6 +946,47 @@ describe("CONNECT pipelined behind a pending response", () => { } }); + test("should write what the 'connect' listener writes at once after the response ahead", async () => { + // The listener answers synchronously and keeps the tunnel open. Its bytes wait for the + // response ahead, which ends later, and are not cut into it. + let first: http.ServerResponse | undefined; + const { promise: handedOff, resolve: onHandoff, reject: onFailure } = Promise.withResolvers(); + await using server = http.createServer((req, res) => { + if (req.method === "CONNECT") return void onFailure(new Error("dispatched as a request")); + first = res; + res.write("first"); + }); + server.on("clientError", onFailure); + server.on("connect", (req, socket) => { + socket.write(ESTABLISHED); + socket.on("data", chunk => socket.write(`echo:${chunk}`)); + onHandoff(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = net.connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + let received = ""; + const { promise: gotEstablished, resolve: onEstablished } = Promise.withResolvers(); + const { promise: gotEcho, resolve: onEcho } = Promise.withResolvers(); + client.setEncoding("latin1"); + client.on("data", chunk => { + received += chunk; + if (received.endsWith(ESTABLISHED)) onEstablished(); + if (received.endsWith("echo:ping")) onEcho(); + }); + client.on("error", onFailure); + client.write(get("/first") + CONNECT); + await handedOff; + first!.end("-done"); + await gotEstablished; + expect(withoutResponseHeads(received)).toBe(`5\r\nfirst\r\n5\r\n-done\r\n0\r\n\r\n${ESTABLISHED}`); + client.write("ping"); + await gotEcho; + } finally { + client.destroy(); + } + }); + test("should close the connection when the server has no 'connect' listener", async () => { const events: string[] = []; let first: http.ServerResponse | undefined; diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index ccbcf1bdfdd5..670dc285eea7 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -4818,6 +4818,98 @@ describe("Upgrade pipelined behind a pending response", () => { } }); + test("should write the 101 of a listener that keeps the tunnel open, after the response ahead", async () => { + // The listener writes the 101 and waits. Nothing ends or uncorks the socket for it. + let first: http.ServerResponse | undefined; + const { promise: handedOff, resolve: onHandoff, reject: onFailure } = Promise.withResolvers(); + await using server = http.createServer((req, res) => { + if (req.headers.upgrade !== undefined) return void onFailure(new Error("dispatched as a request")); + first = res; + res.write("first"); + }); + server.on("clientError", onFailure); + server.on("upgrade", (req, socket) => { + socket.write(SWITCHING); + socket.on("data", chunk => socket.write(`echo:${chunk}`)); + onHandoff(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + let received = ""; + const { promise: got101, resolve: on101 } = Promise.withResolvers(); + const { promise: gotEcho, resolve: onEcho } = Promise.withResolvers(); + client.setEncoding("latin1"); + client.on("data", chunk => { + received += chunk; + if (received.endsWith(SWITCHING)) on101(); + if (received.endsWith("echo:ping")) onEcho(); + }); + client.on("error", onFailure); + client.write(get("/first") + upgradeRequest("/second")); + await handedOff; + first!.end("-done"); + // The client waits for the 101, which waits for the response ahead. + await got101; + expect(withoutResponseHeads(received)).toBe(`5\r\nfirst\r\n5\r\n-done\r\n0\r\n\r\n${SWITCHING}`); + client.write("ping"); + await gotEcho; + } finally { + client.destroy(); + } + }); + + test("should write the 101 after a response ahead that is larger than the socket buffers", async () => { + // The response ahead has finished for JS, but most of its body is still unsent when the + // pipeline reaches the hand-off. The tunnel's bytes and its FIN follow that body. + const total = 4 * 1024 * 1024; + await using server = http.createServer((req, res) => { + res.writeHead(200, { "Content-Length": total }); + res.end(Buffer.alloc(total, "a")); + }); + server.on("upgrade", (req, socket) => socket.end(`${SWITCHING}bye`)); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + const received: Buffer[] = []; + client.on("data", chunk => received.push(chunk)); + // The client reads only once the server has everything queued. + client.pause(); + client.write(get("/first") + upgradeRequest("/second")); + await once(server, "upgrade"); + client.resume(); + await once(client, "end"); + const text = Buffer.concat(received).toString("latin1"); + const bodyStart = text.indexOf("\r\n\r\n") + 4; + expect(text.slice(bodyStart + total)).toBe(`${SWITCHING}bye`); + expect(text.slice(bodyStart, bodyStart + total)).not.toContain("HTTP/1.1 101"); + } finally { + client.destroy(); + } + }); + + test("should fail the parked writes of the listener when the client disconnects early", async () => { + const { promise: handedOff, resolve: onHandoff, reject: onFailure } = Promise.withResolvers(); + const { promise: writeDone, resolve: onWriteDone } = Promise.withResolvers(); + const { promise: endDone, resolve: onEndDone } = Promise.withResolvers(); + await using server = http.createServer((req, res) => { + if (req.headers.upgrade !== undefined) return void onFailure(new Error("dispatched as a request")); + res.write("first"); + }); + server.on("upgrade", (req, socket) => { + socket.write(SWITCHING, onWriteDone); + socket.end("bye", onEndDone); + onHandoff(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + client.write(get("/first") + upgradeRequest("/second")); + await handedOff; + client.destroy(); + expect((await writeDone)?.code).toBe("ERR_STREAM_DESTROYED"); + expect((await endDone)?.code).toBe("ERR_STREAM_DESTROYED"); + }); + test("should go to 'request' when shouldUpgradeCallback declines", async () => { const events: string[] = []; let first: http.ServerResponse | undefined; From f342554b7f6a00320be95c27bacb0a6d1bdd2183 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 16:54:25 +0000 Subject: [PATCH 07/13] node:http: one-line comments for the tunnel write order --- packages/bun-uws/src/HttpContext.h | 3 +-- src/jsc/bindings/node/JSNodeHTTPServerSocket.h | 5 +---- 2 files changed, 2 insertions(+), 6 deletions(-) diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index 732eab22953a..cfeb6276d27a 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -858,8 +858,7 @@ struct HttpContext { /* Drain any socket buffer, this might empty our backpressure and thus finish the request */ asyncSocket->flush(); - /* node:http compat: a CONNECT/Upgrade tunnel's own writes (JSNodeHTTPServerSocket's - * stream buffer) go out once the responses ahead of it owe nothing more. */ + /* node:http compat: a tunnel's own writes go out once the responses ahead owe nothing more. */ if constexpr (IsNodeHttp) { bool tunnel = httpResponseData->isConnectRequest || (httpResponseData->state & HttpResponseData::HTTP_NODE_TUNNEL_AFTER_BODY); if (tunnel && httpResponseData->socketData && httpContextData->onSocketDrain diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h index e69cb9c577c6..2387319fd37f 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h @@ -108,10 +108,7 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject { * truncate the response. Returns true after handing the close to uWS. */ bool shutdownAfterResponseDrains(); - /* A CONNECT/Upgrade tunnel whose connection still owes HTTP output: the - * send buffer of the responses ahead is not empty, or one of them is in - * the middle of a zero-copy write. The tunnel's bytes wait behind them - * (HttpContext::onWritable flushes the stream buffer once they are out). */ + /* A tunnel whose responses ahead still have unsent bytes: its writes wait (HttpContext::onWritable). */ bool tunnelOwesHttpOutput(); /* Switch the connection into CONNECT-style tunnel mode after an accepted From af15405f787053e36de8ddc9094412adba740a08 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 17:34:42 +0000 Subject: [PATCH 08/13] node:http: finish a tunnel's end with its FIN and keep the hand-off in queue order From the review: the ws adoption runs on a fresh turn, since a native callback on the stack may not outlive the socket it replaces. A socket whose FIN waits for the bytes ahead of it completes 'finish' once the FIN is out. A listener's socket.end() keeps its place behind every queued response, and a response that closes the connection releases the hand-off before the FIN. Failed parked callbacks run on the next tick. --- src/js/node/_http_server.ts | 32 ++++++++++------ .../bindings/node/JSNodeHTTPServerSocket.cpp | 20 +++++++++- .../bindings/node/JSNodeHTTPServerSocket.h | 2 + .../node/JSNodeHTTPServerSocketPrototype.cpp | 30 +++++---------- test/js/node/http/node-http.test.ts | 38 +++++++++++++++++++ 5 files changed, 89 insertions(+), 33 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 3b9163b9ace9..f2dbd9b25947 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -1763,9 +1763,7 @@ function getNodeHTTPServerSocket() { const pending = this[kPendingHandoff]; if (pending !== undefined) { this[kPendingHandoff] = undefined; - const reason = err ?? $ERR_STREAM_DESTROYED("write"); - pending.write?.callback(reason); - pending.final?.(reason); + process.nextTick(failParkedHandoff, pending, err ?? $ERR_STREAM_DESTROYED("write")); } const handle = this[kHandle]; if (!handle) { @@ -1797,7 +1795,11 @@ function getNodeHTTPServerSocket() { callback(); return; } - handle.end(); + // A tunnel's FIN waits for the bytes ahead of it; 'finish' waits with it (native drain). + if (handle.end() === false && handle.ondrain) { + this.#pendingCallback = callback; + return; + } callback(); } @@ -2454,6 +2456,8 @@ function emitResponseFinish() { // is eventually closed. function onResponseFinishHandleSocket(server, socket, res) { if (res[kMustCloseConnection]) { + // A hand-off waiting behind this response: its bytes go out before the FIN. + if (socket?.[kPendingHandoff] !== undefined) activatePipelinedHandoff(socket); socket?.end(); return; } @@ -2618,10 +2622,20 @@ function activatePipelinedHandoff(socket) { if (write !== undefined) socket._write(write.chunk, write.encoding, write.callback); const final = pending.final; if (final !== undefined) socket._final(final); - const ready = pending.ready; + // On a fresh turn: a ws adoption replaces the native socket, which no native callback + // on the stack (the response ahead's drain) may outlive. + if (pending.ready.length !== 0) setImmediate(runHandoffReady, pending.ready); +} + +function runHandoffReady(ready) { for (let i = 0; i < ready.length; i++) ready[i](); } +function failParkedHandoff(pending, reason) { + pending.write?.callback(reason); + pending.final?.(reason); +} + // When the connection dies with pipelined responses still queued behind the // in-flight one, abort them and their requests, like Node.js's socketOnClose // (abortIncoming). Runs from the native socket's close path and from the @@ -2659,12 +2673,8 @@ function advanceResponsePipeline(server, socket) { // the pipeline is mutually exclusive with closing the socket - the queued // responses are aborted by the socket close path instead of being replayed // onto a half-closed connection. - if (!socket || socket.destroyed) { - return; - } - if (socket.writableEnded) { - // The end is parked behind the responses ahead; let a hand-off's writes go out before the FIN. - activatePipelinedHandoff(socket); + // (A socket.end() from a 'connect'/'upgrade' listener is parked, not ended yet.) + if (!socket || socket.destroyed || (socket.writableEnded && socket[kPendingHandoff] === undefined)) { return; } const queue = socket[kPipelinedResponses]; diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp index c0b737d9fb77..8123b548672a 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp @@ -296,6 +296,24 @@ bool JSNodeHTTPServerSocket::tunnelOwesHttpOutput() return tunnelOwesHttpOutputImpl(socket); } +void JSNodeHTTPServerSocket::flushAndShutdown(JSC::JSGlobalObject* globalObject) +{ + // onNodeHTTPRequest no longer pauses at dispatch; pause here so the + // shutdown+resume below still cycles kqueue's EVFILT_READ (delete then + // re-add), without which macOS 26 does not deliver the peer's close. + bool cycle = this->ended && this->socket && !this->upgraded; + if (cycle) { + us_socket_pause(this->socket); + } + us_socket_buffered_js_write(this->socket, this->is_ssl, this->ended, false, &this->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); + // Undo the pause above after the shutdown so the unread body drains + // and kqueue's one-shot EVFILT_WRITE (which delivers EV_EOF on + // SHUT_WR) is not deleted by a W -> R|W -> R step. + if (cycle && this->socket) { + us_socket_resume(this->socket); + } +} + bool JSNodeHTTPServerSocket::shutdownAfterResponseDrains() { if (!socket || upgraded || us_socket_is_closed(socket) || us_socket_is_shut_down(socket)) { @@ -757,7 +775,7 @@ void JSNodeHTTPServerSocket::onDrain() if (bufferedSize > 0 || (this->ended && this->socket && !us_socket_is_shut_down(this->socket))) { auto* globalObject = defaultGlobalObject(this->globalObject()); auto scope = DECLARE_TOP_EXCEPTION_SCOPE(globalObject->vm()); - us_socket_buffered_js_write(this->socket, this->is_ssl, this->ended, false, &this->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); + this->flushAndShutdown(globalObject); if (auto* exception = scope.exception()) { (void)scope.tryClearException(); globalObject->reportUncaughtExceptionAtEventLoop(globalObject, exception); diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h index 2387319fd37f..b7930600a031 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h @@ -110,6 +110,8 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject { /* A tunnel whose responses ahead still have unsent bytes: its writes wait (HttpContext::onWritable). */ bool tunnelOwesHttpOutput(); + /* Flush the stream buffer and, once it is empty and ended, shut the write side down. */ + void flushAndShutdown(JSC::JSGlobalObject* globalObject); /* Switch the connection into CONNECT-style tunnel mode after an accepted * Upgrade: subsequent bytes bypass the HTTP parser and stream to the diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp index d6b280db3a93..83128b295f33 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp @@ -218,44 +218,32 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketWrite, (JSC::JSGlobalObje return us_socket_buffered_js_write(thisObject->socket, thisObject->is_ssl, thisObject->ended, hold, &thisObject->streamBuffer, globalObject, JSValue::encode(callFrame->argument(0)), JSValue::encode(callFrame->argument(1))); } +// Returns false when the FIN waits for bytes ahead of it (JSNodeHTTPServerSocket::onDrain sends it). JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketEnd, (JSC::JSGlobalObject * globalObject, JSC::CallFrame* callFrame)) { auto* thisObject = dynamicDowncast(callFrame->thisValue()); if (!thisObject) [[unlikely]] { - return JSValue::encode(JSC::jsUndefined()); + return JSValue::encode(JSC::jsBoolean(true)); } if (thisObject->isClosed()) { - return JSValue::encode(JSC::jsUndefined()); + return JSValue::encode(JSC::jsBoolean(true)); } thisObject->ended = true; // A tunnel's FIN follows the HTTP output ahead of it (JSNodeHTTPServerSocket::onDrain). if (thisObject->tunnelOwesHttpOutput()) { - return JSValue::encode(JSC::jsUndefined()); + return JSValue::encode(JSC::jsBoolean(false)); } // The response's buffered body must reach the kernel before the FIN; uWS // performs the shutdown after its send buffer drains. if (thisObject->shutdownAfterResponseDrains()) { - return JSValue::encode(JSC::jsUndefined()); + return JSValue::encode(JSC::jsBoolean(false)); } - auto bufferedSize = thisObject->streamBuffer.bufferedSize(); - if (bufferedSize == 0) { - // onNodeHTTPRequest no longer pauses at dispatch; pause here so the - // shutdown+resume below still cycles kqueue's EVFILT_READ (delete then - // re-add), without which macOS 26 does not deliver the peer's close. - if (thisObject->socket && !thisObject->upgraded) { - us_socket_pause(thisObject->socket); - } - auto result = us_socket_buffered_js_write(thisObject->socket, thisObject->is_ssl, thisObject->ended, false, &thisObject->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); - // Undo the pause above after the shutdown so the unread body drains - // and kqueue's one-shot EVFILT_WRITE (which delivers EV_EOF on - // SHUT_WR) is not deleted by a W -> R|W -> R step. - if (thisObject->socket && !thisObject->upgraded) { - us_socket_resume(thisObject->socket); - } - return result; + if (thisObject->streamBuffer.bufferedSize() == 0) { + thisObject->flushAndShutdown(globalObject); + return JSValue::encode(JSC::jsBoolean(true)); } - return JSValue::encode(JSC::jsUndefined()); + return JSValue::encode(JSC::jsBoolean(false)); } // Implementation of custom getters diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index 670dc285eea7..b0dba46b45d5 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -4910,6 +4910,44 @@ describe("Upgrade pipelined behind a pending response", () => { expect((await endDone)?.code).toBe("ERR_STREAM_DESTROYED"); }); + test("should send the listener's socket.end() after every response queued ahead of it", async () => { + // Two requests ahead: the first ends later, the second at once. The listener declines with + // socket.end(), like the builtin ws does for a bad handshake. + const DECLINED = "HTTP/1.1 400 Bad Request\r\nConnection: close\r\n\r\n"; + let first: http.ServerResponse | undefined; + const { promise: handedOff, resolve: onHandoff, reject: onFailure } = Promise.withResolvers(); + await using server = http.createServer((req, res) => { + if (req.headers.upgrade !== undefined) return void onFailure(new Error("dispatched as a request")); + if (req.url === "/first") { + first = res; + res.write("first"); + } else { + res.end("second"); + } + }); + server.on("clientError", onFailure); + server.on("upgrade", (req, socket) => { + socket.end(DECLINED); + onHandoff(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + const received: Buffer[] = []; + client.on("data", chunk => received.push(chunk)); + client.on("error", onFailure); + client.write(get("/first") + get("/second") + upgradeRequest("/third")); + await handedOff; + first!.end("-done"); + await once(client, "end"); + expect(withoutResponseHeads(Buffer.concat(received).toString("latin1"))).toBe( + `5\r\nfirst\r\n5\r\n-done\r\n0\r\n\r\nsecond${DECLINED}`, + ); + } finally { + client.destroy(); + } + }); + test("should go to 'request' when shouldUpgradeCallback declines", async () => { const events: string[] = []; let first: http.ServerResponse | undefined; From beba49156d253c7ce68a9fd5539239afde193ed6 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 17:36:51 +0000 Subject: [PATCH 09/13] node:http: one-line comments in the hand-off and the shutdown helper --- src/js/node/_http_server.ts | 3 +-- src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp | 8 ++------ 2 files changed, 3 insertions(+), 8 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index f2dbd9b25947..4dd0468400ed 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -2622,8 +2622,7 @@ function activatePipelinedHandoff(socket) { if (write !== undefined) socket._write(write.chunk, write.encoding, write.callback); const final = pending.final; if (final !== undefined) socket._final(final); - // On a fresh turn: a ws adoption replaces the native socket, which no native callback - // on the stack (the response ahead's drain) may outlive. + // A fresh turn: a ws adoption must not run inside a native callback of the response ahead. if (pending.ready.length !== 0) setImmediate(runHandoffReady, pending.ready); } diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp index 8123b548672a..6a4f3c817375 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp @@ -298,17 +298,13 @@ bool JSNodeHTTPServerSocket::tunnelOwesHttpOutput() void JSNodeHTTPServerSocket::flushAndShutdown(JSC::JSGlobalObject* globalObject) { - // onNodeHTTPRequest no longer pauses at dispatch; pause here so the - // shutdown+resume below still cycles kqueue's EVFILT_READ (delete then - // re-add), without which macOS 26 does not deliver the peer's close. + // Pause, shut down, resume: cycles kqueue's EVFILT_READ, or macOS 26 does not deliver the peer's close. bool cycle = this->ended && this->socket && !this->upgraded; if (cycle) { us_socket_pause(this->socket); } us_socket_buffered_js_write(this->socket, this->is_ssl, this->ended, false, &this->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); - // Undo the pause above after the shutdown so the unread body drains - // and kqueue's one-shot EVFILT_WRITE (which delivers EV_EOF on - // SHUT_WR) is not deleted by a W -> R|W -> R step. + // Resume after the shutdown: a W -> R|W -> R step would delete kqueue's one-shot EVFILT_WRITE. if (cycle && this->socket) { us_socket_resume(this->socket); } From 8b0b1d973348c7c38ac6c80cd5342445a95ee84e Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 18:26:21 +0000 Subject: [PATCH 10/13] node:http: let a queued tunnel end the connection after a Connection: close response ahead A response with Connection: close closed the connection natively at its end, before the hand-off queued behind it could write. The native close gate now leaves that to the tunnel's own socket.end(). A destroyed socket also fails the write callback that waits for a native drain. --- packages/bun-uws/src/HttpResponseData.h | 4 ++++ src/js/node/_http_server.ts | 14 +++++++---- test/js/node/http/node-http.test.ts | 31 +++++++++++++++++++++++++ 3 files changed, 44 insertions(+), 5 deletions(-) diff --git a/packages/bun-uws/src/HttpResponseData.h b/packages/bun-uws/src/HttpResponseData.h index e1d258cc9067..328e15b6fd8a 100644 --- a/packages/bun-uws/src/HttpResponseData.h +++ b/packages/bun-uws/src/HttpResponseData.h @@ -233,6 +233,10 @@ struct HttpResponseData : AsyncSocketData, HttpParser { /* Whether the connection should be torn down once the in-flight response (if * any) has completed and all buffered outgoing data has been flushed. */ bool shouldCloseConnection() const { + /* node:http: a tunnel queued behind this response ends the connection itself, after its parked bytes. */ + if ((isConnectRequest || (state & HTTP_NODE_TUNNEL_AFTER_BODY)) && nodeHttpQueuedPipelinedCount > 0) { + return false; + } return (state & HTTP_CONNECTION_CLOSE) || ((state & HTTP_NODE_RECEIVED_FIN) && nodeHttpQueuedPipelinedCount == 0) || ((state & HTTP_CLOSE_WHEN_IDLE) && this->isIdle); diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 4dd0468400ed..0732498f284f 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -1760,10 +1760,13 @@ function getNodeHTTPServerSocket() { } _destroy(err, callback) { + // Writes that wait for a native drain that never comes now. const pending = this[kPendingHandoff]; - if (pending !== undefined) { + const waiting = this.#pendingCallback; + if (pending !== undefined || waiting !== null) { this[kPendingHandoff] = undefined; - process.nextTick(failParkedHandoff, pending, err ?? $ERR_STREAM_DESTROYED("write")); + this.#pendingCallback = null; + process.nextTick(failParkedHandoff, pending, waiting, err ?? $ERR_STREAM_DESTROYED("write")); } const handle = this[kHandle]; if (!handle) { @@ -2630,9 +2633,10 @@ function runHandoffReady(ready) { for (let i = 0; i < ready.length; i++) ready[i](); } -function failParkedHandoff(pending, reason) { - pending.write?.callback(reason); - pending.final?.(reason); +function failParkedHandoff(pending, waiting, reason) { + pending?.write?.callback(reason); + pending?.final?.(reason); + waiting?.(reason); } // When the connection dies with pipelined responses still queued behind the diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index b0dba46b45d5..781018543b9a 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -4948,6 +4948,37 @@ describe("Upgrade pipelined behind a pending response", () => { } }); + test("should send the listener's reply before the FIN of a response ahead that closes the connection", async () => { + const { promise: handedOff, resolve: onHandoff, reject: onFailure } = Promise.withResolvers(); + let first: http.ServerResponse | undefined; + await using server = http.createServer((req, res) => { + if (req.headers.upgrade !== undefined) return void onFailure(new Error("dispatched as a request")); + first = res; + res.setHeader("Connection", "close"); + res.write("first"); + }); + server.on("upgrade", (req, socket) => { + socket.write(SWITCHING); + onHandoff(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + const received: Buffer[] = []; + client.on("data", chunk => received.push(chunk)); + client.on("error", onFailure); + client.write(get("/first") + upgradeRequest("/second")); + await handedOff; + first!.end("-done"); + await once(client, "end"); + expect(withoutResponseHeads(Buffer.concat(received).toString("latin1"))).toBe( + `5\r\nfirst\r\n5\r\n-done\r\n0\r\n\r\n${SWITCHING}`, + ); + } finally { + client.destroy(); + } + }); + test("should go to 'request' when shouldUpgradeCallback declines", async () => { const events: string[] = []; let first: http.ServerResponse | undefined; From 55a5d3d7e343add2283e16e73fc850fab00be6fd Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 19:36:54 +0000 Subject: [PATCH 11/13] node:http: a queued tunnel holds no reads and closes after a Connection: close response ahead An Upgrade with a body is a tunnel before its body ends; a queued one no longer keeps reads paused, so the listener can read the body while the response ahead drains. After a response that closes the connection, the hand-off socket is destroyed once its FIN is out, like Node's destroySoon(). --- src/js/node/_http_server.ts | 12 ++++- .../bindings/node/JSNodeHTTPServerSocket.cpp | 5 ++- test/js/node/http/node-http.test.ts | 44 +++++++++++++++++++ 3 files changed, 57 insertions(+), 4 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 0732498f284f..046e4e731aae 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -2459,8 +2459,12 @@ function emitResponseFinish() { // is eventually closed. function onResponseFinishHandleSocket(server, socket, res) { if (res[kMustCloseConnection]) { - // A hand-off waiting behind this response: its bytes go out before the FIN. - if (socket?.[kPendingHandoff] !== undefined) activatePipelinedHandoff(socket); + if (socket?.[kPendingHandoff] !== undefined) { + // A hand-off waiting behind this response: its bytes go out before the FIN, and the + // connection closes after it, like Node's destroySoon() (native leaves that to this end). + activatePipelinedHandoff(socket); + socket.once("finish", destroyHandoffSocketNT); + } socket?.end(); return; } @@ -2633,6 +2637,10 @@ function runHandoffReady(ready) { for (let i = 0; i < ready.length; i++) ready[i](); } +function destroyHandoffSocketNT(this: NodeHTTPServerSocket) { + this.destroy(); +} + function failParkedHandoff(pending, waiting, reason) { pending?.write?.callback(reason); pending?.final?.(reason); diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp index 6a4f3c817375..de17299cacb0 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp @@ -462,11 +462,12 @@ void JSNodeHTTPServerSocket::appendPipelinedResponse(JSC::VM& vm, WebCore::JSNod m_pipelinedResponses.last().set(vm, this, response); } -/* A pipelined CONNECT stays queued so that the connection never counts as idle. No request follows it, so it holds no reads. */ +/* A pipelined CONNECT/Upgrade stays queued so that the connection never counts as idle. No request follows it, so it holds no reads. */ template static bool queuedResponsesHoldReads(uWS::NodeHttpResponseData* httpResponseData) { - return httpResponseData->nodeHttpQueuedPipelinedCount > 0 && !httpResponseData->isConnectRequest; + bool tunnel = httpResponseData->isConnectRequest || (httpResponseData->state & uWS::HttpResponseData::HTTP_NODE_TUNNEL_AFTER_BODY); + return httpResponseData->nodeHttpQueuedPipelinedCount > 0 && !tunnel; } /* node:http flood prevention, resume half. Parked pipelined requests (HttpParser::nodeHttpPausedSpill) diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index 781018543b9a..89f726f28586 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -4979,6 +4979,50 @@ describe("Upgrade pipelined behind a pending response", () => { } }); + test("should read the body of the Upgrade request while the response ahead has not drained", async () => { + // The response ahead fills the socket buffers, which pauses reads for queued responses. + // A queued hand-off must not hold them: the listener waits for the body before it ends + // the response ahead. + const total = 4 * 1024 * 1024; + const body = "0123456789"; + let first: http.ServerResponse | undefined; + const { promise: gotBody, resolve: onBody, reject: onFailure } = Promise.withResolvers(); + await using server = http.createServer((req, res) => { + if (req.headers.upgrade !== undefined) return void onFailure(new Error("dispatched as a request")); + first = res; + res.writeHead(200, { "Content-Length": total }); + res.write(Buffer.alloc(total, "a")); + }); + server.on("clientError", onFailure); + server.on("upgrade", (req, socket) => { + let read = ""; + req.on("data", chunk => (read += chunk)); + req.on("end", () => { + onBody(read); + socket.end(SWITCHING); + first!.end(); + }); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const client = connect((server.address() as AddressInfo).port, "127.0.0.1"); + try { + const received: Buffer[] = []; + client.on("data", chunk => received.push(chunk)); + client.on("error", onFailure); + client.pause(); + client.write(get("/first") + upgradeRequest("/second", `Content-Length: ${body.length}\r\n`)); + await once(server, "upgrade"); + // The body comes in a packet of its own, while the client still does not read. + client.write(body); + expect(await gotBody).toBe(body); + client.resume(); + await once(client, "end"); + expect(Buffer.concat(received).toString("latin1").endsWith(SWITCHING)).toBe(true); + } finally { + client.destroy(); + } + }); + test("should go to 'request' when shouldUpgradeCallback declines", async () => { const events: string[] = []; let first: http.ServerResponse | undefined; From a74f3cc4e20dbfe7555c6c6039a859d67aa7d2ee Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 19:42:44 +0000 Subject: [PATCH 12/13] node:http: one-line comment in the close path of the hand-off --- src/js/node/_http_server.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 046e4e731aae..bd6a7e02d3e7 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -2460,8 +2460,7 @@ function emitResponseFinish() { function onResponseFinishHandleSocket(server, socket, res) { if (res[kMustCloseConnection]) { if (socket?.[kPendingHandoff] !== undefined) { - // A hand-off waiting behind this response: its bytes go out before the FIN, and the - // connection closes after it, like Node's destroySoon() (native leaves that to this end). + // A hand-off behind this response writes before the FIN; then destroySoon(), like Node. activatePipelinedHandoff(socket); socket.once("finish", destroyHandoffSocketNT); } From dc2445606c4f090b4863b843d74df9369d3d1dab Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 20:51:37 +0000 Subject: [PATCH 13/13] node:http: make the reads-gate test fill the uWS send buffer --- test/js/node/http/node-http.test.ts | 22 +++++++++++++--------- 1 file changed, 13 insertions(+), 9 deletions(-) diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index 89f726f28586..9f3491eb623f 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -4980,18 +4980,18 @@ describe("Upgrade pipelined behind a pending response", () => { }); test("should read the body of the Upgrade request while the response ahead has not drained", async () => { - // The response ahead fills the socket buffers, which pauses reads for queued responses. - // A queued hand-off must not hold them: the listener waits for the body before it ends - // the response ahead. - const total = 4 * 1024 * 1024; + // Small writes past the kernel buffers fill the uWS send buffer, which pauses the reads of + // the connection when the Upgrade is parsed. Once the buffer drains, a queued hand-off must + // not keep them paused: the listener waits for the body before it ends the response ahead. + const chunk = Buffer.alloc(8 * 1024, "a"); + const count = 1024; const body = "0123456789"; let first: http.ServerResponse | undefined; const { promise: gotBody, resolve: onBody, reject: onFailure } = Promise.withResolvers(); await using server = http.createServer((req, res) => { if (req.headers.upgrade !== undefined) return void onFailure(new Error("dispatched as a request")); first = res; - res.writeHead(200, { "Content-Length": total }); - res.write(Buffer.alloc(total, "a")); + for (let i = 0; i < count; i++) res.write(chunk); }); server.on("clientError", onFailure); server.on("upgrade", (req, socket) => { @@ -5012,12 +5012,16 @@ describe("Upgrade pipelined behind a pending response", () => { client.pause(); client.write(get("/first") + upgradeRequest("/second", `Content-Length: ${body.length}\r\n`)); await once(server, "upgrade"); - // The body comes in a packet of its own, while the client still does not read. + // The body comes in a packet of its own. The server reads it once its buffer has drained. client.write(body); - expect(await gotBody).toBe(body); client.resume(); + expect(await gotBody).toBe(body); await once(client, "end"); - expect(Buffer.concat(received).toString("latin1").endsWith(SWITCHING)).toBe(true); + const text = withoutResponseHeads(Buffer.concat(received).toString("latin1")); + const bodyEnd = text.indexOf("0\r\n\r\n"); + expect(text.slice(bodyEnd)).toBe(`0\r\n\r\n${SWITCHING}`); + expect(text.slice(0, bodyEnd).replace(/2000\r\na+\r\n/g, "")).toBe(""); + expect(received.reduce((n, b) => n + b.length, 0)).toBeGreaterThan(chunk.length * count); } finally { client.destroy(); }