From 39ba5346887ab35f8fb28da56a8afafa87a1bca5 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 23:45:16 +0000 Subject: [PATCH 1/5] node:http: leave the response ahead alone when a pipelined dispatch throws The native dispatch tail answers a throw by ending the uWS response of the connection. A pipelined request's response is queued, and the connection has one uWS response state, so the tail ended the response ahead of it: the client got part of that body and a close. The tail now leaves a queued response queued and marks it. When its turn comes, node:http sends it if the listener completed it. If not, the connection closes after the responses ahead of it: the socket ends as it does for a response that must close the connection, so bytes that are still buffered go out first. The turn of a destroyed queued response closes the connection the same way. destroy() there cut a large response ahead short. startPipelinedResponse() refuses a response that is not the next one in the native queue, so a dispatch that throws before node:http queues a response cannot give its turn to a later response. --- src/js/internal/http.ts | 1 + src/js/node/_http_server.ts | 53 +++++- .../bindings/node/JSNodeHTTPServerSocket.cpp | 8 +- .../bindings/node/JSNodeHTTPServerSocket.h | 3 +- src/runtime/server/NodeHTTPResponse.rs | 13 ++ src/runtime/server/mod.rs | 30 ++- .../http/node-http-pipelined-throw-fixture.js | 174 ++++++++++++++++++ test/js/node/http/node-http.test.ts | 97 ++++++++++ 8 files changed, 359 insertions(+), 20 deletions(-) create mode 100644 test/js/node/http/node-http-pipelined-throw-fixture.js diff --git a/src/js/internal/http.ts b/src/js/internal/http.ts index 9631dec34b71..e7ea3faca045 100644 --- a/src/js/internal/http.ts +++ b/src/js/internal/http.ts @@ -77,6 +77,7 @@ export const enum NodeHTTPResponseFlags { request_has_completed = 1 << 1, ended = 1 << 2, upgraded = 1 << 3, + dispatch_threw_while_queued = 1 << 9, closed_or_completed = socket_closed | request_has_completed, } diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 5e2c41acc30b..9a16bb99e335 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -2551,6 +2551,20 @@ function abortQueuedPipelinedResponses(socket) { } } +// The response at the head of the queue can never be sent, and an HTTP/1.1 +// connection cannot skip its turn, so the response that just finished was the +// last one. A native socket ends like one whose response must close the +// connection (onResponseFinishHandleSocket): uWS closes it once the bytes still +// buffered for that response have left, where destroy() would discard them. +// The close path then aborts what is queued. +function closeAfterLastSendableResponse(socket) { + if (NodeHTTPServerSocket && socket instanceof NodeHTTPServerSocket) { + socket.end(); + } else if (!socket.destroyed) { + socket.destroy(); + } +} + function advanceResponsePipeline(server, socket) { // The previous response on this connection closed it (Connection: close, // HTTP/1.0, maxRequestsPerSocket): like Node.js's resOnFinish, advancing @@ -2566,9 +2580,25 @@ function advanceResponsePipeline(server, socket) { } const res = queue.shift(); const queued = res[kPipelinedQueuedState]; + const handle = res[kHandle]; + + if ( + !queued.ended && + !res.destroyed && + handle && + (handle.flags & NodeHTTPResponseFlags.dispatch_threw_while_queued) !== 0 + ) { + // The dispatch of this request threw and nothing ended the response since. + // For the connection's current response the native dispatch tail answers + // and closes at once; a queued one gets the same at its turn. It goes back + // in the queue so the close path aborts it, and its request, with the rest. + queue.unshift(res); + closeAfterLastSendableResponse(socket); + return; + } + res[kPipelinedQueuedState] = undefined; releasePipelineOutgoingData(socket, queued.bytes); - const handle = res[kHandle]; if (res.destroyed || !handle) { // The queued response was destroyed before it could be sent; the @@ -2576,9 +2606,7 @@ function advanceResponsePipeline(server, socket) { // Deliberate divergence from Node v26, which assigns the destroyed // message and wedges the connection until requestTimeout: an HTTP/1.1 // connection cannot skip a response slot, so reset it instead. - if (!socket.destroyed) { - socket.destroy(); - } + closeAfterLastSendableResponse(socket); return; } @@ -2589,11 +2617,20 @@ function advanceResponsePipeline(server, socket) { socket.destroyed || !socketHandle.startPipelinedResponse(handle, !!queued.isAncient, !requestShouldKeepAlive(res.req)) ) { - // The connection is already gone; the socket close path destroys queued - // responses, but make sure this (already dequeued) one is not skipped. - if (!res.destroyed) { - res.destroy(); + if (socket.destroyed) { + // The close path may have run already: do not skip this (dequeued) one. + if (!res.destroyed) { + res.destroy(); + } + return; } + // The connection is closing, or this response is not the next one + // natively: a dispatch ahead of it threw before it queued a response + // here, and that turn cannot be skipped. Back in the queue, the close + // path aborts it, and its request, with the rest. + res[kPipelinedQueuedState] = queued; + queue.unshift(res); + closeAfterLastSendableResponse(socket); return; } diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp index 921c90164c82..c45548a43d3e 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp @@ -586,7 +586,13 @@ bool JSNodeHTTPServerSocket::startPipelinedResponse(JSC::VM& vm, WebCore::JSNode bool hasMoreQueued = false; { Locker locker { m_pipelinedResponsesLock }; - m_pipelinedResponses.removeFirstMatching([&](auto& entry) { return entry.get() == response; }); + // Responses leave in request order. Another one first means JS never + // queued it (its dispatch threw before that), and its turn cannot be + // given to this one: the client would read this as the answer to that. + if (m_pipelinedResponses.isEmpty() || m_pipelinedResponses.first().get() != response) { + return false; + } + m_pipelinedResponses.removeAt(0); hasMoreQueued = !m_pipelinedResponses.isEmpty(); } diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h index bdd077af2db8..f15227f9e773 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h @@ -97,7 +97,8 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject { /* Make a previously queued pipelined response the connection's current * response: reset the per-response uWS state (the part the request handler * normally resets per parsed request) and, when the queue drained, resume - * socket reads. Returns false when the connection is already gone. */ + * socket reads. Returns false when the connection is already gone, or when + * the response is not the next one in the queue. */ bool startPipelinedResponse(JSC::VM& vm, WebCore::JSNodeHTTPResponse* response, bool isAncient, bool connectionClose); /* Stop parsing further HTTP requests on this connection (Node frees the * parser when 'close' is emitted on the socket). */ diff --git a/src/runtime/server/NodeHTTPResponse.rs b/src/runtime/server/NodeHTTPResponse.rs index 085b2f40f843..513fb7cbbd56 100644 --- a/src/runtime/server/NodeHTTPResponse.rs +++ b/src/runtime/server/NodeHTTPResponse.rs @@ -95,6 +95,10 @@ bitflags! { /// node:http handed this connection to a raw 'upgrade'/'connect' /// tunnel (JSNodeHTTPServerSocket::upgradeToTunnelMode). const TUNNELED = 1 << 8; + /// The dispatch of this request threw while the response was queued + /// behind another one. Nothing ended it natively: node:http closes the + /// connection at its turn unless JS ended it (advanceResponsePipeline). + const DISPATCH_THREW_WHILE_QUEUED = 1 << 9; } } @@ -442,6 +446,15 @@ impl NodeHTTPResponse { Bun__getNodeHTTPResponseThisValue(any_response_is_ssl(&raw), raw.socket().cast()) } + /// Pipelining: the connection has another current response, and this one + /// waits for its turn. Until then the state of `raw_response` (one per + /// connection) describes that other response. + pub(crate) fn is_queued_behind_current_response(&self) -> bool { + self.get_this_value() + .as_class_ref::() + .is_some_and(|current| !ptr::eq(current, self)) + } + fn get_server_socket_value(&self) -> JSValue { let flags = self.flags.get(); let Some(raw) = self.raw_response.get() else { diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index da89a189c0f4..cd722287ac26 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1464,7 +1464,10 @@ impl NewServer { // SAFETY: see `nhr` above. let nhr = unsafe { &*node_http_response }; let nhr_flags = nhr.flags.get(); - if !nhr_flags.contains(NhrFlags::UPGRADED) { + // A pipelined dispatch: the pending response that + // `raw_response` reports is the one ahead of this one. + let is_queued = nhr.is_queued_behind_current_response(); + if !nhr_flags.contains(NhrFlags::UPGRADED) && !is_queued { if let Some(raw) = nhr.raw_response.get() { if !nhr_flags.contains(NhrFlags::REQUEST_HAS_COMPLETED) && raw.state().is_response_pending() @@ -1478,15 +1481,22 @@ impl NewServer { } } } - // The handler threw before `res.end()`; we just ended (or - // will never end) the raw response above. Mark ENDED so - // `on_request_complete()` → `mark_request_as_done()` runs - // and releases the `IS_REQUEST_PENDING` ref (one of the - // initial 3). Without this the box leaks: the later - // `on_abort` socket-close path early-returns once - // `REQUEST_HAS_COMPLETED` is set and never balances it. - nhr.flags.set(nhr.flags.get() | NhrFlags::ENDED); - nhr.on_request_complete(); + if is_queued { + // Nothing was ended, so this response stays queued + // like any other one. + nhr.flags + .set(nhr.flags.get() | NhrFlags::DISPATCH_THREW_WHILE_QUEUED); + } else { + // The handler threw before `res.end()`; we just ended (or + // will never end) the raw response above. Mark ENDED so + // `on_request_complete()` → `mark_request_as_done()` runs + // and releases the `IS_REQUEST_PENDING` ref (one of the + // initial 3). Without this the box leaks: the later + // `on_abort` socket-close path early-returns once + // `REQUEST_HAS_COMPLETED` is set and never balances it. + nhr.flags.set(nhr.flags.get() | NhrFlags::ENDED); + nhr.on_request_complete(); + } } } HttpResult::Success | HttpResult::Pending => {} diff --git a/test/js/node/http/node-http-pipelined-throw-fixture.js b/test/js/node/http/node-http-pipelined-throw-fixture.js new file mode 100644 index 000000000000..912be6569c1c --- /dev/null +++ b/test/js/node/http/node-http-pipelined-throw-fixture.js @@ -0,0 +1,174 @@ +// The dispatch of a request throws while an earlier response on the connection is still pending. +// MODE selects the scenario. The only line of stdout is the result as JSON. +// Under Node.js (`MODE=request node `) the modes request, checkContinue, +// checkExpectation, ended, ended-later, large and large-destroyed print the same result. The +// others wait for a close that Node never makes. +const http = require("node:http"); +const net = require("node:net"); + +const mode = process.env.MODE; +const events = []; +process.on("uncaughtException", err => events.push(`uncaught: ${err.message}`)); + +const get = (path, headers = "") => `GET ${path} HTTP/1.1\r\nHost: example.com\r\n${headers}\r\n`; +const written = new Map([ + ["checkContinue", get("/first") + get("/second", "Expect: 100-continue\r\n")], + ["checkExpectation", get("/first") + get("/second", "Expect: something\r\n")], + ["unfinished", get("/first") + get("/second") + get("/third") + get("/fourth")], + ["not-pipelined", get("/first")], +]); +const pipelinedPair = get("/first") + get("/second"); + +// The "large" modes end /first with more bytes than a socket buffer takes at once, so most of +// them are still on their way out when the turn of /second comes. +const isLarge = mode.startsWith("large"); +const firstBodyLength = isLarge ? 8 * 1024 * 1024 : 10; + +let received = ""; +let firstBodyBytes = 0; +let clientClosed = false; +const serverSideCloses = []; + +function report(extra) { + const result = isLarge + ? { firstBodyBytes } + : { bodies: ["first-done", "second", "third", "fourth"].filter(body => received.includes(body)) }; + console.log(JSON.stringify({ events, ...result, closed: clientClosed, ...extra })); + process.exit(0); +} + +// "unfinished" waits for the client, and for each request and response behind /first, to close. +// server.close() calls back only when no request is pending, the one that threw included. +function reportUnfinished() { + if (clientClosed && serverSideCloses.length === 6) { + server.close(() => report({ serverSideCloses: serverSideCloses.sort() })); + } +} + +// The response to /first stays open until a later dispatch finishes it. +let first; +const finishFirst = () => first.end(isLarge ? Buffer.alloc(firstBodyLength - 5, "-") : "-done"); + +function fail(thrower, req) { + events.push(`${thrower} ${req.url}`); + throw new Error(`${thrower} threw`); +} + +// node:http queues a response only after it constructed it. +class ResponseThatThrows extends http.ServerResponse { + constructor(req, options) { + super(req, options); + if (req.url === "/second") { + setImmediate(finishFirst); + fail("constructor", req); + } + } +} + +const server = http.createServer(mode === "constructor" ? { ServerResponse: ResponseThatThrows } : {}, (req, res) => { + if (req.url === "/first") { + events.push(`request ${req.url}`); + if (mode === "not-pipelined") return void res.end("first-done"); + first = res; + res.writeHead(200, { "Content-Length": String(firstBodyLength) }); + res.write("first"); + return; + } + switch (mode) { + case "unfinished": + // /second and /fourth answer. /third throws and leaves its response open. + for (const emitter of [req, res]) { + emitter.on("close", () => { + serverSideCloses.push(req.url); + reportUnfinished(); + }); + } + if (req.url === "/third") fail("request", req); + res.end(req.url.slice(1)); + if (req.url === "/fourth") setImmediate(finishFirst); + return; + case "constructor": + events.push(`request ${req.url}`); + return void res.end(req.url.slice(1)); + case "large-destroyed": + // No throw: the queued response is destroyed, and its turn resets the connection too. + events.push(`request ${req.url}`); + setImmediate(finishFirst); + return void res.destroy(); + case "ended": + // The response is complete before the throw. + res.end("second"); + setImmediate(finishFirst); + break; + case "ended-later": + // The response is complete before its turn comes. + setImmediate(() => { + res.end("second"); + finishFirst(); + }); + break; + case "request": + case "large": + setImmediate(finishFirst); + break; + } + fail("request", req); +}); +for (const eventName of ["checkContinue", "checkExpectation"]) { + server.on(eventName, req => { + setImmediate(finishFirst); + fail(eventName, req); + }); +} + +server.listen(0, "127.0.0.1", () => { + const client = net.connect(server.address().port, "127.0.0.1"); + let wroteAgain = false; + let head = ""; + client.on("error", () => {}); + client.on("close", () => { + clientClosed = true; + if (mode === "unfinished") reportUnfinished(); + else report(); + }); + client.on("data", chunk => { + if (isLarge) { + // Count the body of /first: no other response can follow it. + if (head === undefined) { + firstBodyBytes += chunk.length; + } else { + head += chunk.toString("latin1"); + const headEnd = head.indexOf("\r\n\r\n"); + if (headEnd === -1) return; + firstBodyBytes = head.length - (headEnd + 4); + head = undefined; + } + if (firstBodyBytes === firstBodyLength) report(); + return; + } + received += chunk.toString("latin1"); + switch (mode) { + case "request": + case "checkContinue": + case "checkExpectation": + if (received.includes("first-done")) report(); + break; + case "ended": + case "ended-later": + if (received.includes("second")) report(); + break; + case "constructor": + // A response that took the turn of /second would show up here. + if (received.includes("third")) report(); + // fallthrough + case "not-pipelined": + // One more request, after the first response is complete. + if (received.includes("first-done") && !wroteAgain) { + wroteAgain = true; + client.write(get(mode === "constructor" ? "/third" : "/second")); + } + break; + } + }); + client.write(written.get(mode) ?? pipelinedPair); +}); diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index 71ad7e6435d6..b49c946b1dbc 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -2892,6 +2892,103 @@ it("pipelined responses buffered past the high water mark pause reads on the con } }); +// The native dispatch tail answers a throw by ending the connection's current response. For a +// pipelined request that was the response ahead of it: the client got "first" of a 10-byte body, +// then the close. The throw is an uncaught exception, so each scenario runs in a child process. +describe("a dispatch that throws while an earlier response on the connection is pending", () => { + async function run(mode: string) { + await using proc = Bun.spawn({ + cmd: [bunExe(), path.join(import.meta.dir, "node-http-pipelined-throw-fixture.js")], + env: { ...bunEnv, MODE: mode }, + stdout: "pipe", + stderr: "inherit", + }); + const [stdout, exitCode] = await Promise.all([proc.stdout.text(), proc.exited]); + return { result: stdout ? JSON.parse(stdout) : undefined, exitCode }; + } + + // node v26.3.0 gives the same result for the tests in these three loops. + for (const emitted of ["request", "checkContinue", "checkExpectation"]) { + it.concurrent(`'${emitted}' listener: the response ahead still completes`, async () => { + expect(await run(emitted)).toEqual({ + result: { + events: ["request /first", `${emitted} /second`, `uncaught: ${emitted} threw`], + bodies: ["first-done"], + closed: false, + }, + exitCode: 0, + }); + }); + } + // The response ahead has ended, and most of its 8 MB are still in the send buffer when the + // connection is reset at the turn of the response behind it. A destroyed queued response + // resets the connection in the same way, with no throw. + for (const [mode, events] of [ + ["large", ["request /first", "request /second", "uncaught: request threw"]], + ["large-destroyed", ["request /first", "request /second"]], + ] as const) { + it.concurrent(`a large response ahead arrives in full before the connection is reset (${mode})`, async () => { + expect(await run(mode)).toEqual({ + result: { events, firstBodyBytes: 8 * 1024 * 1024, closed: false }, + exitCode: 0, + }); + }); + } + for (const mode of ["ended", "ended-later"]) { + it.concurrent(`a response that is complete when its turn comes is sent (${mode})`, async () => { + expect(await run(mode)).toEqual({ + result: { + events: ["request /first", "request /second", "uncaught: request threw"], + bodies: ["first-done", "second"], + closed: false, + }, + exitCode: 0, + }); + }); + } + + // The tests below wait for a close. Node never makes it: it answers nothing and keeps the + // connection, which then waits forever. Bun answers a throw with a close. For a queued response + // that happens when its turn comes, after the responses ahead of it. The requests behind it are + // aborted with the connection. + it.concurrent("a response that is not complete resets the connection when its turn comes", async () => { + expect(await run("unfinished")).toEqual({ + result: { + events: ["request /first", "request /third", "uncaught: request threw"], + bodies: ["first-done", "second"], + closed: true, + serverSideCloses: ["/fourth", "/fourth", "/second", "/second", "/third", "/third"], + }, + exitCode: 0, + }); + }); + + // The throw comes before node:http queued a response, so nothing holds the turn of /second. + // A later response must not go out in its place. + it.concurrent("a throw from the ServerResponse constructor does not let a later response take the turn", async () => { + expect(await run("constructor")).toEqual({ + result: { + events: ["request /first", "constructor /second", "uncaught: constructor threw", "request /third"], + bodies: ["first-done"], + closed: true, + }, + exitCode: 0, + }); + }); + + // Unchanged: with no response ahead, the throwing request is the current one and is answered at once. + it.concurrent("a request that is not pipelined is still answered with a close", async () => { + expect(await run("not-pipelined")).toEqual({ + result: { + events: ["request /first", "request /second", "uncaught: request threw"], + bodies: ["first-done"], + closed: true, + }, + exitCode: 0, + }); + }); +}); + it("requireHostHeader still rejects Upgrade-carrying requests that dispatch as normal requests", async () => { // The native parser exempts Upgrade requests from the Host check so genuine // upgrades can reach the 'upgrade' event, but a request that falls through From a6139fda527537453817c31e5edebb54cce86fff Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 21:26:44 +0000 Subject: [PATCH 2/5] node:http: close a connection whose next pipelined response cannot be sent without waiting for the client socket.end() at such a turn sends only a FIN when nothing is buffered, and a client that never answers it kept the connection, its queued requests and the server's pending request count. closeWhenDrained() marks the connection and runs uWS's close gate: the connection closes now when nothing is buffered, and from the drain gate when bytes of the responses ahead are still on their way. A pipelined dispatch now holds its turn in the JS queue with its native handle until its response exists. A throw before that (the constructor of a user's IncomingMessage or ServerResponse subclass) then has an entry at its turn, and the connection closes there. This replaces the order check in startPipelinedResponse(), which closed such a connection only on the next request or on the keep-alive timeout. The code that the lint rule no-duplicate-conditional-property-access rejected goes away with it. --- src/js/node/_http_server.ts | 97 ++++++++++--------- .../bindings/node/JSNodeHTTPServerSocket.cpp | 33 +++++-- .../bindings/node/JSNodeHTTPServerSocket.h | 9 +- .../node/JSNodeHTTPServerSocketPrototype.cpp | 13 +++ .../http/node-http-pipelined-throw-fixture.js | 55 ++++++++--- test/js/node/http/node-http.test.ts | 32 +++--- 6 files changed, 154 insertions(+), 85 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 9a16bb99e335..266f639d957a 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -690,6 +690,21 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort socket = new (getNodeHTTPServerSocket())(server, socketHandle, !!tls); } + if (isPipelinedDispatch) { + // The native queue already holds this request's turn. Hold it here too, + // with the native handle, until the response exists and replaces it + // below: a throw on the way (the constructor of a user's IncomingMessage + // or ServerResponse subclass) must not let a later response take it. + (socket[kPipelinedResponses] ??= []).push(handle); + // 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); + } + } + // Like Node.js's resetSocketTimeout (parserOnIncoming): a new request // arriving on a kept-alive connection replaces the keep-alive idle // timeout with the server's regular per-socket timeout. @@ -889,15 +904,10 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // A previous response on this connection has not finished yet: like // 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). + // pipeline assigns it the socket (advanceResponsePipeline). It takes + // the turn that the native handle held since the top of this dispatch. + socket[kPipelinedResponses]?.pop(); 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); - } // 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. @@ -2533,6 +2543,9 @@ function abortQueuedPipelinedResponses(socket) { socket[kPipelinedResponses] = undefined; for (let i = 0; i < pipelinedLength; i++) { const queuedRes = pipelined[i]; + // A turn that the native handle still holds (its dispatch queued no + // response): nothing to abort here, the native close path notifies it. + if (queuedRes[kPipelinedQueuedState] === undefined) continue; const queuedReq = queuedRes.req; if (queuedReq && !queuedReq.destroyed) { queuedReq[kHandle] = undefined; @@ -2553,13 +2566,12 @@ function abortQueuedPipelinedResponses(socket) { // The response at the head of the queue can never be sent, and an HTTP/1.1 // connection cannot skip its turn, so the response that just finished was the -// last one. A native socket ends like one whose response must close the -// connection (onResponseFinishHandleSocket): uWS closes it once the bytes still -// buffered for that response have left, where destroy() would discard them. -// The close path then aborts what is queued. +// last one. A native socket closes once the bytes still buffered for that +// response have left: destroy() would discard them, and end() would wait for +// the client's FIN. The close path then aborts what is queued. function closeAfterLastSendableResponse(socket) { if (NodeHTTPServerSocket && socket instanceof NodeHTTPServerSocket) { - socket.end(); + socket[kHandle]?.closeWhenDrained(); } else if (!socket.destroyed) { socket.destroy(); } @@ -2578,38 +2590,40 @@ function advanceResponsePipeline(server, socket) { if (!queue || queue.length === 0) { return; } - const res = queue.shift(); + const res = queue[0]; const queued = res[kPipelinedQueuedState]; + if (queued === undefined) { + // Still the native handle that holds the turn of a pipelined dispatch: the + // dispatch queued no response. If it threw, none can come. + if ((res.flags & NodeHTTPResponseFlags.dispatch_threw_while_queued) !== 0) { + closeAfterLastSendableResponse(socket); + } + return; + } const handle = res[kHandle]; if ( - !queued.ended && - !res.destroyed && - handle && - (handle.flags & NodeHTTPResponseFlags.dispatch_threw_while_queued) !== 0 + res.destroyed || + !handle || + (!queued.ended && (handle.flags & NodeHTTPResponseFlags.dispatch_threw_while_queued) !== 0) ) { - // The dispatch of this request threw and nothing ended the response since. - // For the connection's current response the native dispatch tail answers - // and closes at once; a queued one gets the same at its turn. It goes back - // in the queue so the close path aborts it, and its request, with the rest. - queue.unshift(res); + // The queued response was destroyed before it could be sent, or the + // dispatch of its request threw and nothing ended it since (the native + // dispatch tail answers a throw at once only for the connection's current + // response). The connection cannot produce a response for this slot, so it + // is unusable. Deliberate divergence from Node v26, which assigns the + // message and wedges the connection until requestTimeout: an HTTP/1.1 + // connection cannot skip a response slot, so reset it instead. The entry + // stays queued: nothing behind it can start, and the close path aborts it, + // and its request, with the rest. closeAfterLastSendableResponse(socket); return; } + queue.shift(); res[kPipelinedQueuedState] = undefined; releasePipelineOutgoingData(socket, queued.bytes); - if (res.destroyed || !handle) { - // The queued response was destroyed before it could be sent; the - // connection cannot produce a response for this slot, so it is unusable. - // Deliberate divergence from Node v26, which assigns the destroyed - // message and wedges the connection until requestTimeout: an HTTP/1.1 - // connection cannot skip a response slot, so reset it instead. - closeAfterLastSendableResponse(socket); - return; - } - if (NodeHTTPServerSocket && socket instanceof NodeHTTPServerSocket) { const socketHandle = socket[kHandle]; if ( @@ -2617,20 +2631,11 @@ function advanceResponsePipeline(server, socket) { socket.destroyed || !socketHandle.startPipelinedResponse(handle, !!queued.isAncient, !requestShouldKeepAlive(res.req)) ) { - if (socket.destroyed) { - // The close path may have run already: do not skip this (dequeued) one. - if (!res.destroyed) { - res.destroy(); - } - return; + // The connection is already gone; the socket close path destroys queued + // responses, but make sure this (already dequeued) one is not skipped. + if (!res.destroyed) { + res.destroy(); } - // The connection is closing, or this response is not the next one - // natively: a dispatch ahead of it threw before it queued a response - // here, and that turn cannot be skipped. Back in the queue, the close - // path aborts it, and its request, with the rest. - res[kPipelinedQueuedState] = queued; - queue.unshift(res); - closeAfterLastSendableResponse(socket); return; } diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp index c45548a43d3e..73084656dc9d 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp @@ -280,6 +280,31 @@ bool JSNodeHTTPServerSocket::shutdownAfterResponseDrains() return deferShutdownUntilResponseDrains(socket); } +template +static void closeWhenDrainedImpl(us_socket_t* socket) +{ + auto* httpResponseData = reinterpret_cast*>(us_socket_ext(socket)); + /* What a response that closes the connection leaves behind: the close gates + * (here, or HttpContext::onWritable once the send buffer has flushed) + * shut down and close a connection marked like this with no response pending. */ + httpResponseData->state |= uWS::HttpResponseData::HTTP_CONNECTION_CLOSE; + /* A response that ended inside the read being parsed is still in the cork buffer. */ + reinterpret_cast*>(socket)->uncork(); + reinterpret_cast*>(socket)->closeIfDoneAndMarked(httpResponseData); +} + +void JSNodeHTTPServerSocket::closeWhenDrained() +{ + if (!socket || upgraded || us_socket_is_closed(socket)) { + return; + } + if (is_ssl) { + closeWhenDrainedImpl(socket); + } else { + closeWhenDrainedImpl(socket); + } +} + template static bool isRequestTimedOutImpl(us_socket_t* socket, uint64_t headersTimeoutMs, uint64_t requestTimeoutMs) { @@ -586,13 +611,7 @@ bool JSNodeHTTPServerSocket::startPipelinedResponse(JSC::VM& vm, WebCore::JSNode bool hasMoreQueued = false; { Locker locker { m_pipelinedResponsesLock }; - // Responses leave in request order. Another one first means JS never - // queued it (its dispatch threw before that), and its turn cannot be - // given to this one: the client would read this as the answer to that. - if (m_pipelinedResponses.isEmpty() || m_pipelinedResponses.first().get() != response) { - return false; - } - m_pipelinedResponses.removeAt(0); + m_pipelinedResponses.removeFirstMatching([&](auto& entry) { return entry.get() == response; }); hasMoreQueued = !m_pipelinedResponses.isEmpty(); } diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h index f15227f9e773..a3b18c2c6229 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h @@ -97,8 +97,7 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject { /* Make a previously queued pipelined response the connection's current * response: reset the per-response uWS state (the part the request handler * normally resets per parsed request) and, when the queue drained, resume - * socket reads. Returns false when the connection is already gone, or when - * the response is not the next one in the queue. */ + * socket reads. Returns false when the connection is already gone. */ bool startPipelinedResponse(JSC::VM& vm, WebCore::JSNodeHTTPResponse* response, bool isAncient, bool connectionClose); /* Stop parsing further HTTP requests on this connection (Node frees the * parser when 'close' is emitted on the socket). */ @@ -109,6 +108,12 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject { * truncate the response. Returns true after handing the close to uWS. */ bool shutdownAfterResponseDrains(); + /* node:http pipelining: the next queued response can never be sent, so the + * response that just ended was the last one. Close the connection once its + * bytes have left: now when none are buffered, otherwise from uWS's close + * gate. close() would discard them, and end() would wait for the peer's FIN. */ + void closeWhenDrained(); + /* 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 f85cd51c825e..952cd5f096c4 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp @@ -47,6 +47,7 @@ JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketSetResponseTrailers); JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketIsRequestTimedOut); JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketStartPipelinedResponse); JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketStopParsing); +JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketCloseWhenDrained); JSC_DECLARE_CUSTOM_GETTER(jsNodeHttpServerSocketGetterResponse); JSC_DECLARE_CUSTOM_GETTER(jsNodeHttpServerSocketGetterRemoteAddress); JSC_DECLARE_CUSTOM_GETTER(jsNodeHttpServerSocketGetterLocalAddress); @@ -83,6 +84,7 @@ static const JSC::HashTableValue JSNodeHTTPServerSocketPrototypeTableValues[] = { "isRequestTimedOut"_s, static_cast(JSC::PropertyAttribute::Function | JSC::PropertyAttribute::DontEnum), JSC::NoIntrinsic, { JSC::HashTableValue::NativeFunctionType, jsFunctionNodeHTTPServerSocketIsRequestTimedOut, 2 } }, { "startPipelinedResponse"_s, static_cast(JSC::PropertyAttribute::Function | JSC::PropertyAttribute::DontEnum), JSC::NoIntrinsic, { JSC::HashTableValue::NativeFunctionType, jsFunctionNodeHTTPServerSocketStartPipelinedResponse, 3 } }, { "stopParsing"_s, static_cast(JSC::PropertyAttribute::Function | JSC::PropertyAttribute::DontEnum), JSC::NoIntrinsic, { JSC::HashTableValue::NativeFunctionType, jsFunctionNodeHTTPServerSocketStopParsing, 0 } }, + { "closeWhenDrained"_s, static_cast(JSC::PropertyAttribute::Function | JSC::PropertyAttribute::DontEnum), JSC::NoIntrinsic, { JSC::HashTableValue::NativeFunctionType, jsFunctionNodeHTTPServerSocketCloseWhenDrained, 0 } }, { "secureEstablished"_s, static_cast(JSC::PropertyAttribute::CustomAccessor | JSC::PropertyAttribute::ReadOnly), JSC::NoIntrinsic, { JSC::HashTableValue::GetterSetterType, jsNodeHttpServerSocketGetterIsSecureEstablished, noOpSetter } }, { "servername"_s, static_cast(JSC::PropertyAttribute::CustomAccessor | JSC::PropertyAttribute::ReadOnly), JSC::NoIntrinsic, { JSC::HashTableValue::GetterSetterType, jsNodeHttpServerSocketGetterServername, noOpSetter } }, { "authorizationError"_s, static_cast(JSC::PropertyAttribute::CustomAccessor | JSC::PropertyAttribute::ReadOnly), JSC::NoIntrinsic, { JSC::HashTableValue::GetterSetterType, jsNodeHttpServerSocketGetterAuthorizationError, noOpSetter } }, @@ -203,6 +205,17 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketStopParsing, (JSC::JSGlob return JSValue::encode(JSC::jsUndefined()); } +// node:http pipelining: the response that just ended was the connection's last +// one (see JSNodeHTTPServerSocket::closeWhenDrained). +JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketCloseWhenDrained, (JSC::JSGlobalObject * globalObject, JSC::CallFrame* callFrame)) +{ + auto* thisObject = dynamicDowncast(callFrame->thisValue()); + if (thisObject) [[likely]] { + thisObject->closeWhenDrained(); + } + return JSValue::encode(JSC::jsUndefined()); +} + JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketWrite, (JSC::JSGlobalObject * globalObject, JSC::CallFrame* callFrame)) { auto* thisObject = dynamicDowncast(callFrame->thisValue()); diff --git a/test/js/node/http/node-http-pipelined-throw-fixture.js b/test/js/node/http/node-http-pipelined-throw-fixture.js index 912be6569c1c..74278acd2513 100644 --- a/test/js/node/http/node-http-pipelined-throw-fixture.js +++ b/test/js/node/http/node-http-pipelined-throw-fixture.js @@ -1,8 +1,7 @@ // The dispatch of a request throws while an earlier response on the connection is still pending. // MODE selects the scenario. The only line of stdout is the result as JSON. -// Under Node.js (`MODE=request node `) the modes request, checkContinue, -// checkExpectation, ended, ended-later, large and large-destroyed print the same result. The -// others wait for a close that Node never makes. +// Under Node.js (`MODE=request node `) every mode prints the same result, except +// unfinished and not-pipelined: they wait for a close that Node never makes. const http = require("node:http"); const net = require("node:net"); @@ -37,10 +36,13 @@ function report(extra) { process.exit(0); } -// "unfinished" waits for the client, and for each request and response behind /first, to close. -// server.close() calls back only when no request is pending, the one that threw included. +// "unfinished" waits for the end of the connection, and for each request and response behind +// /first to close. server.close() calls back only when no request is pending, the one that threw +// included. +let closingServer = false; function reportUnfinished() { - if (clientClosed && serverSideCloses.length === 6) { + if (clientClosed && serverSideCloses.length === 6 && !closingServer) { + closingServer = true; server.close(() => report({ serverSideCloses: serverSideCloses.sort() })); } } @@ -54,18 +56,34 @@ function fail(thrower, req) { throw new Error(`${thrower} threw`); } -// node:http queues a response only after it constructed it. +// node:http queues a response only after it constructed the request and the response. class ResponseThatThrows extends http.ServerResponse { constructor(req, options) { super(req, options); if (req.url === "/second") { setImmediate(finishFirst); - fail("constructor", req); + fail("ServerResponse", req); } } } +// The url is not known yet in this constructor: the second request is the one that throws. +let requests = 0; +class RequestThatThrows extends http.IncomingMessage { + constructor(...args) { + super(...args); + if (++requests === 2) { + setImmediate(finishFirst); + events.push("IncomingMessage /second"); + throw new Error("IncomingMessage threw"); + } + } +} +const serverOptions = new Map([ + ["constructor-response", { ServerResponse: ResponseThatThrows }], + ["constructor-request", { IncomingMessage: RequestThatThrows }], +]); -const server = http.createServer(mode === "constructor" ? { ServerResponse: ResponseThatThrows } : {}, (req, res) => { +const server = http.createServer(serverOptions.get(mode) ?? {}, (req, res) => { if (req.url === "/first") { events.push(`request ${req.url}`); if (mode === "not-pipelined") return void res.end("first-done"); @@ -87,7 +105,8 @@ const server = http.createServer(mode === "constructor" ? { ServerResponse: Resp res.end(req.url.slice(1)); if (req.url === "/fourth") setImmediate(finishFirst); return; - case "constructor": + case "constructor-response": + case "constructor-request": events.push(`request ${req.url}`); return void res.end(req.url.slice(1)); case "large-destroyed": @@ -122,15 +141,18 @@ for (const eventName of ["checkContinue", "checkExpectation"]) { } server.listen(0, "127.0.0.1", () => { - const client = net.connect(server.address().port, "127.0.0.1"); + // This client never answers a FIN, so the server has to close the connection on its own. + const client = net.connect({ port: server.address().port, host: "127.0.0.1", allowHalfOpen: true }); let wroteAgain = false; let head = ""; - client.on("error", () => {}); - client.on("close", () => { + const onServerClosedConnection = () => { clientClosed = true; if (mode === "unfinished") reportUnfinished(); else report(); - }); + }; + client.on("error", () => {}); + client.on("end", onServerClosedConnection); + client.on("close", onServerClosedConnection); client.on("data", chunk => { if (isLarge) { // Count the body of /first: no other response can follow it. @@ -157,7 +179,8 @@ server.listen(0, "127.0.0.1", () => { case "ended-later": if (received.includes("second")) report(); break; - case "constructor": + case "constructor-response": + case "constructor-request": // A response that took the turn of /second would show up here. if (received.includes("third")) report(); // fallthrough @@ -165,7 +188,7 @@ server.listen(0, "127.0.0.1", () => { // One more request, after the first response is complete. if (received.includes("first-done") && !wroteAgain) { wroteAgain = true; - client.write(get(mode === "constructor" ? "/third" : "/second")); + client.write(get(mode === "not-pipelined" ? "/second" : "/third")); } break; } diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index b49c946b1dbc..38ae8cb6a01e 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -2907,7 +2907,7 @@ describe("a dispatch that throws while an earlier response on the connection is return { result: stdout ? JSON.parse(stdout) : undefined, exitCode }; } - // node v26.3.0 gives the same result for the tests in these three loops. + // node v26.3.0 gives the same result for the tests in these four loops. for (const emitted of ["request", "checkContinue", "checkExpectation"]) { it.concurrent(`'${emitted}' listener: the response ahead still completes`, async () => { expect(await run(emitted)).toEqual({ @@ -2947,6 +2947,23 @@ describe("a dispatch that throws while an earlier response on the connection is }); } + // The throw comes before node:http has a response to queue. The turn of /second is still held: + // the connection closes when it comes, and no later response goes out in its place. Node ends + // with the same result, because it answers the next request with a 400 and closes. + for (const thrower of ["ServerResponse", "IncomingMessage"]) { + it.concurrent(`a throw from the ${thrower} constructor closes the connection at its turn`, async () => { + const mode = thrower === "ServerResponse" ? "constructor-response" : "constructor-request"; + expect(await run(mode)).toEqual({ + result: { + events: ["request /first", `${thrower} /second`, `uncaught: ${thrower} threw`], + bodies: ["first-done"], + closed: true, + }, + exitCode: 0, + }); + }); + } + // The tests below wait for a close. Node never makes it: it answers nothing and keeps the // connection, which then waits forever. Bun answers a throw with a close. For a queued response // that happens when its turn comes, after the responses ahead of it. The requests behind it are @@ -2963,19 +2980,6 @@ describe("a dispatch that throws while an earlier response on the connection is }); }); - // The throw comes before node:http queued a response, so nothing holds the turn of /second. - // A later response must not go out in its place. - it.concurrent("a throw from the ServerResponse constructor does not let a later response take the turn", async () => { - expect(await run("constructor")).toEqual({ - result: { - events: ["request /first", "constructor /second", "uncaught: constructor threw", "request /third"], - bodies: ["first-done"], - closed: true, - }, - exitCode: 0, - }); - }); - // Unchanged: with no response ahead, the throwing request is the current one and is answered at once. it.concurrent("a request that is not pipelined is still answered with a close", async () => { expect(await run("not-pipelined")).toEqual({ From 9b13e32f849c5ab53349322e72a366d94a63d3dd Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 22:29:41 +0000 Subject: [PATCH 3/5] node:http: keep the pipeline kick after the response is queued drainMicrotasks() in the dispatcher also runs the tick queue, so a kick that is scheduled only at the top of a pipelined dispatch runs while the turn is still held by the native handle, and nothing starts the response afterwards. The kick is back where the response is queued. The one at the top stays for a dispatch that throws before that point. The native tail now skips its whole block for a queued response (mark_dispatch_threw_if_queued), so the existing block is unchanged. Comments are one line each. --- src/js/node/_http_server.ts | 56 ++++++++----------- .../bindings/node/JSNodeHTTPServerSocket.cpp | 4 +- .../bindings/node/JSNodeHTTPServerSocket.h | 5 +- .../node/JSNodeHTTPServerSocketPrototype.cpp | 2 - src/runtime/server/NodeHTTPResponse.rs | 19 ++++--- src/runtime/server/mod.rs | 37 ++++++------ test/js/node/http/node-http.test.ts | 32 +++++++++++ 7 files changed, 84 insertions(+), 71 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 266f639d957a..d8bebb5a5edb 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -691,18 +691,10 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort } if (isPipelinedDispatch) { - // The native queue already holds this request's turn. Hold it here too, - // with the native handle, until the response exists and replaces it - // below: a throw on the way (the constructor of a user's IncomingMessage - // or ServerResponse subclass) must not let a later response take it. + // The native handle holds this request's turn until its response exists, also when a throw prevents that. (socket[kPipelinedResponses] ??= []).push(handle); - // 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); - } + // For a throw before the kick below. drainMicrotasks() further down runs this one when nothing throws. + kickPipelineIfIdle(server, socket); } // Like Node.js's resetSocketTimeout (parserOnIncoming): a new request @@ -904,10 +896,13 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // A previous response on this connection has not finished yet: like // 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). It takes - // the turn that the native handle held since the top of this dispatch. - socket[kPipelinedResponses]?.pop(); + // pipeline assigns it the socket (advanceResponsePipeline). + socket[kPipelinedResponses]?.pop(); // the turn that the native handle held 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. + 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. @@ -2508,6 +2503,12 @@ function releasePipelineOutgoingData(socket, bytes) { // connection's current response, is assigned the socket, and its buffered // output is flushed. const kPipelineKickScheduled = Symbol("kPipelineKickScheduled"); +function kickPipelineIfIdle(server, socket) { + if (socket._httpMessage == null && !socket[kPipelineKickScheduled]) { + socket[kPipelineKickScheduled] = true; + process.nextTick(advancePipelineIfIdleNT, server, socket); + } +} function advancePipelineIfIdleNT(server, socket) { socket[kPipelineKickScheduled] = false; if (socket._httpMessage == null && socket[kPipelinedResponses]?.length) { @@ -2543,8 +2544,7 @@ function abortQueuedPipelinedResponses(socket) { socket[kPipelinedResponses] = undefined; for (let i = 0; i < pipelinedLength; i++) { const queuedRes = pipelined[i]; - // A turn that the native handle still holds (its dispatch queued no - // response): nothing to abort here, the native close path notifies it. + // A turn that the native handle still holds: the native close path notifies that one. if (queuedRes[kPipelinedQueuedState] === undefined) continue; const queuedReq = queuedRes.req; if (queuedReq && !queuedReq.destroyed) { @@ -2564,11 +2564,7 @@ function abortQueuedPipelinedResponses(socket) { } } -// The response at the head of the queue can never be sent, and an HTTP/1.1 -// connection cannot skip its turn, so the response that just finished was the -// last one. A native socket closes once the bytes still buffered for that -// response have left: destroy() would discard them, and end() would wait for -// the client's FIN. The close path then aborts what is queued. +// The head of the queue can never be sent: close once the bytes of the responses ahead of it have left. function closeAfterLastSendableResponse(socket) { if (NodeHTTPServerSocket && socket instanceof NodeHTTPServerSocket) { socket[kHandle]?.closeWhenDrained(); @@ -2593,8 +2589,7 @@ function advanceResponsePipeline(server, socket) { const res = queue[0]; const queued = res[kPipelinedQueuedState]; if (queued === undefined) { - // Still the native handle that holds the turn of a pipelined dispatch: the - // dispatch queued no response. If it threw, none can come. + // The native handle still holds this turn: its dispatch queued no response, and after a throw none can come. if ((res.flags & NodeHTTPResponseFlags.dispatch_threw_while_queued) !== 0) { closeAfterLastSendableResponse(socket); } @@ -2605,18 +2600,15 @@ function advanceResponsePipeline(server, socket) { if ( res.destroyed || !handle || + // Its dispatch threw and nothing ended it since: natively only the current response is answered for a throw. (!queued.ended && (handle.flags & NodeHTTPResponseFlags.dispatch_threw_while_queued) !== 0) ) { - // The queued response was destroyed before it could be sent, or the - // dispatch of its request threw and nothing ended it since (the native - // dispatch tail answers a throw at once only for the connection's current - // response). The connection cannot produce a response for this slot, so it - // is unusable. Deliberate divergence from Node v26, which assigns the + // The queued response was destroyed before it could be sent; the + // connection cannot produce a response for this slot, so it is unusable. + // Deliberate divergence from Node v26, which assigns the destroyed // message and wedges the connection until requestTimeout: an HTTP/1.1 - // connection cannot skip a response slot, so reset it instead. The entry - // stays queued: nothing behind it can start, and the close path aborts it, - // and its request, with the rest. - closeAfterLastSendableResponse(socket); + // connection cannot skip a response slot, so reset it instead. + closeAfterLastSendableResponse(socket); // it stays queued for the close path return; } diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp index 73084656dc9d..d8ea4966902b 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp @@ -284,9 +284,7 @@ template static void closeWhenDrainedImpl(us_socket_t* socket) { auto* httpResponseData = reinterpret_cast*>(us_socket_ext(socket)); - /* What a response that closes the connection leaves behind: the close gates - * (here, or HttpContext::onWritable once the send buffer has flushed) - * shut down and close a connection marked like this with no response pending. */ + /* uWS's close gates (below, or onWritable after the flush) close a connection marked like this. */ httpResponseData->state |= uWS::HttpResponseData::HTTP_CONNECTION_CLOSE; /* A response that ended inside the read being parsed is still in the cork buffer. */ reinterpret_cast*>(socket)->uncork(); diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h index a3b18c2c6229..d45b65fe7413 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(); - /* node:http pipelining: the next queued response can never be sent, so the - * response that just ended was the last one. Close the connection once its - * bytes have left: now when none are buffered, otherwise from uWS's close - * gate. close() would discard them, and end() would wait for the peer's FIN. */ + /* Close once the bytes of the responses that ended have left. close() discards them, end() waits for the peer. */ void closeWhenDrained(); /* Switch the connection into CONNECT-style tunnel mode after an accepted diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp index 952cd5f096c4..eca0bf813856 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp @@ -205,8 +205,6 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketStopParsing, (JSC::JSGlob return JSValue::encode(JSC::jsUndefined()); } -// node:http pipelining: the response that just ended was the connection's last -// one (see JSNodeHTTPServerSocket::closeWhenDrained). JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketCloseWhenDrained, (JSC::JSGlobalObject * globalObject, JSC::CallFrame* callFrame)) { auto* thisObject = dynamicDowncast(callFrame->thisValue()); diff --git a/src/runtime/server/NodeHTTPResponse.rs b/src/runtime/server/NodeHTTPResponse.rs index 513fb7cbbd56..4c45d11fa82f 100644 --- a/src/runtime/server/NodeHTTPResponse.rs +++ b/src/runtime/server/NodeHTTPResponse.rs @@ -95,9 +95,7 @@ bitflags! { /// node:http handed this connection to a raw 'upgrade'/'connect' /// tunnel (JSNodeHTTPServerSocket::upgradeToTunnelMode). const TUNNELED = 1 << 8; - /// The dispatch of this request threw while the response was queued - /// behind another one. Nothing ended it natively: node:http closes the - /// connection at its turn unless JS ended it (advanceResponsePipeline). + /// Its dispatch threw while it was queued (pipelining): advanceResponsePipeline decides at its turn. const DISPATCH_THREW_WHILE_QUEUED = 1 << 9; } } @@ -446,13 +444,16 @@ impl NodeHTTPResponse { Bun__getNodeHTTPResponseThisValue(any_response_is_ssl(&raw), raw.socket().cast()) } - /// Pipelining: the connection has another current response, and this one - /// waits for its turn. Until then the state of `raw_response` (one per - /// connection) describes that other response. - pub(crate) fn is_queued_behind_current_response(&self) -> bool { - self.get_this_value() + /// Flags this response when another one is the connection's current response, and says so. + pub(crate) fn mark_dispatch_threw_if_queued(&self) -> bool { + let queued = self + .get_this_value() .as_class_ref::() - .is_some_and(|current| !ptr::eq(current, self)) + .is_some_and(|current| !ptr::eq(current, self)); + if queued { + self.update_flags(|f| f.insert(Flags::DISPATCH_THREW_WHILE_QUEUED)); + } + queued } fn get_server_socket_value(&self) -> JSValue { diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index cd722287ac26..297545249490 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1460,14 +1460,16 @@ impl NewServer { ) }; - if !node_http_response.is_null() { + // A pipelined response stays queued: `raw_response` describes the one ahead of it. + let threw_while_queued = !node_http_response.is_null() + // SAFETY: see `nhr` above. + && unsafe { &*node_http_response }.mark_dispatch_threw_if_queued(); + + if !node_http_response.is_null() && !threw_while_queued { // SAFETY: see `nhr` above. let nhr = unsafe { &*node_http_response }; let nhr_flags = nhr.flags.get(); - // A pipelined dispatch: the pending response that - // `raw_response` reports is the one ahead of this one. - let is_queued = nhr.is_queued_behind_current_response(); - if !nhr_flags.contains(NhrFlags::UPGRADED) && !is_queued { + 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() @@ -1481,22 +1483,15 @@ impl NewServer { } } } - if is_queued { - // Nothing was ended, so this response stays queued - // like any other one. - nhr.flags - .set(nhr.flags.get() | NhrFlags::DISPATCH_THREW_WHILE_QUEUED); - } else { - // The handler threw before `res.end()`; we just ended (or - // will never end) the raw response above. Mark ENDED so - // `on_request_complete()` → `mark_request_as_done()` runs - // and releases the `IS_REQUEST_PENDING` ref (one of the - // initial 3). Without this the box leaks: the later - // `on_abort` socket-close path early-returns once - // `REQUEST_HAS_COMPLETED` is set and never balances it. - nhr.flags.set(nhr.flags.get() | NhrFlags::ENDED); - nhr.on_request_complete(); - } + // The handler threw before `res.end()`; we just ended (or + // will never end) the raw response above. Mark ENDED so + // `on_request_complete()` → `mark_request_as_done()` runs + // and releases the `IS_REQUEST_PENDING` ref (one of the + // initial 3). Without this the box leaks: the later + // `on_abort` socket-close path early-returns once + // `REQUEST_HAS_COMPLETED` is set and never balances it. + nhr.flags.set(nhr.flags.get() | NhrFlags::ENDED); + nhr.on_request_complete(); } } HttpResult::Success | HttpResult::Pending => {} diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index 38ae8cb6a01e..b5bca7a00214 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -2892,6 +2892,38 @@ it("pipelined responses buffered past the high water mark pause reads on the con } }); +it("a pipelined response is started when no response is in flight to hand it the socket", async () => { + // The dispatcher kicks the pipeline when it queues a response and nothing is in flight (the + // previous response finished and detached while it still counts as pending). Clearing + // socket._httpMessage by hand reaches that state: only the kick can start /second then. + const server = createServer((req, res) => { + if (req.url === "/first") { + (req.socket as any)._httpMessage = null; + return; + } + res.end("second-response"); + }); + try { + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const socket = connect((server.address() as AddressInfo).port, "127.0.0.1"); + const { promise: started, resolve: onStarted, reject: onFailure } = Promise.withResolvers(); + let received = ""; + socket.on("data", chunk => { + received += chunk.toString("latin1"); + if (received.includes("second-response")) onStarted(); + }); + socket.on("error", onFailure); + socket.on("close", () => onFailure(new Error("closed before the queued response was started"))); + socket.write("GET /first HTTP/1.1\r\nHost: x\r\n\r\nGET /second HTTP/1.1\r\nHost: x\r\n\r\n"); + await started; + expect(received).toContain("second-response"); + socket.destroy(); + } finally { + server.close(); + } +}); + // The native dispatch tail answers a throw by ending the connection's current response. For a // pipelined request that was the response ahead of it: the client got "first" of a 10-byte body, // then the close. The throw is an uncaught exception, so each scenario runs in a child process. From 460f430f2dd31e1f5d284da5ad71098120bbb490 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 22:59:11 +0000 Subject: [PATCH 4/5] test(node:http): release the client socket of the pipeline kick test on failure, and check the response The last assertion restated what already resolved the awaited promise. It now checks that the stream is the one response to /second. The socket is destroyed in the finally block, so a rejection does not leave it open. --- test/js/node/http/node-http.test.ts | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index b5bca7a00214..b3e5507d5c0f 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -2896,6 +2896,7 @@ it("a pipelined response is started when no response is in flight to hand it the // The dispatcher kicks the pipeline when it queues a response and nothing is in flight (the // previous response finished and detached while it still counts as pending). Clearing // socket._httpMessage by hand reaches that state: only the kick can start /second then. + // Bun only: Node.js v26.3.0 fails an internal assertion (resOnFinish) on this use of its internals. const server = createServer((req, res) => { if (req.url === "/first") { (req.socket as any)._httpMessage = null; @@ -2903,10 +2904,11 @@ it("a pipelined response is started when no response is in flight to hand it the } res.end("second-response"); }); + let socket: ReturnType | undefined; try { server.listen(0, "127.0.0.1"); await once(server, "listening"); - const socket = connect((server.address() as AddressInfo).port, "127.0.0.1"); + socket = connect((server.address() as AddressInfo).port, "127.0.0.1"); const { promise: started, resolve: onStarted, reject: onFailure } = Promise.withResolvers(); let received = ""; socket.on("data", chunk => { @@ -2917,9 +2919,14 @@ it("a pipelined response is started when no response is in flight to hand it the socket.on("close", () => onFailure(new Error("closed before the queued response was started"))); socket.write("GET /first HTTP/1.1\r\nHost: x\r\n\r\nGET /second HTTP/1.1\r\nHost: x\r\n\r\n"); await started; - expect(received).toContain("second-response"); - socket.destroy(); + // /first never writes, so the stream is the one response to /second and nothing else. + const [head, ...bodies] = received.split("\r\n\r\n"); + expect({ status: head.split("\r\n")[0], bodies }).toEqual({ + status: "HTTP/1.1 200 OK", + bodies: ["second-response"], + }); } finally { + socket?.destroy(); server.close(); } }); From a88cedd1aaeb258401062a11b5f2fe7e07503161 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 23:58:02 +0000 Subject: [PATCH 5/5] ci: retrigger