From f3f9a22162089f6f7067decaa0c23d12bbf7a067 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 23:32:08 +0000 Subject: [PATCH 1/4] node:http: keep the body of a paused Upgrade request and resume reads at the switch to tunnel mode A read of the upgrade socket resumed the request at once, also one that its 'upgrade' listener paused. The listener runs before the rest of its read is parsed, so a body that arrived with the head then flowed away with no listener. A paused request is now resumed one event loop turn later, and only if its body is still incomplete then, like Node.js's UpgradeStream. A request that is not paused is resumed as before. The switch to tunnel mode at the end of the body now resumes reads that req.pause() or a full request buffer paused. pause() and resume() on the response do nothing in tunnel mode, so the socket stayed paused and later tunnel bytes never arrived. --- packages/bun-uws/src/HttpContext.h | 11 ++ src/js/node/_http_server.ts | 19 +- .../http/node-http-req-socket-pause.test.ts | 164 +++++++++++++++++- 3 files changed, 192 insertions(+), 2 deletions(-) diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index 6978d65dbe63..27f9ad3f5c38 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -657,6 +657,17 @@ struct HttpContext { httpResponseData->inStream = nullptr; } } + if constexpr (IsNodeHttp) { + /* Not after upgrade() from the body handler: the ext holds a WebSocketData then. */ + if (switchToTunnelAfterThisChunk && httpContextData->upgradedWebSocket != user) { + /* pause() and resume() on the response do nothing in tunnel mode: + * lift a read pause that the request body left (req.pause(), a full buffer). */ + Bun__NodeHTTP__onReadsResumable(SSL, (struct us_socket_t *) user); + if (us_socket_is_closed((struct us_socket_t *) user)) { + return nullptr; + } + } + } return user; }); diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 5e2c41acc30b..e3165f0c2ff6 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -1319,6 +1319,20 @@ function clearUpgradeIncoming(socket) { socket[kUpgradeIncoming] = undefined; } +// Like Node.js, a read of the raw socket resumes a paused request only while its body is +// incomplete. Checked one event loop turn after the read: the 'upgrade' listener runs before +// the rest of its read is parsed, and that rest can complete the body. +function resumePausedUpgradeIncoming(socket) { + const req = socket[kUpgradeIncoming]; + if (req === undefined) return; + const response = socket[kHandle]?.response; + if (response && (response.hasBody & NodeHTTPBodyReadState.done) !== 0) { + socket[kUpgradeIncoming] = undefined; + } else { + req.resume(); + } +} + // Node.js hands the connection over to 'connect'/'upgrade' listeners with the // connection-listener set removed (onParserExecuteCommon removes its data/end/ // close/drain/error/timeout listeners) and only net.Socket's own 'end' listener @@ -1779,7 +1793,10 @@ function getNodeHTTPServerSocket() { upgradeIncoming.push(resumed); } } - upgradeIncoming.resume(); + // Not paused: resumed at once, also with a complete body. req.complete is still false + // inside the listener then, so a listener can wait for 'end' and never read the body. + if (upgradeIncoming.readableFlowing !== false) upgradeIncoming.resume(); + else setImmediate(resumePausedUpgradeIncoming, this); return; } if (response) { diff --git a/test/js/node/http/node-http-req-socket-pause.test.ts b/test/js/node/http/node-http-req-socket-pause.test.ts index 7d3c2ac4d373..be150595a3d6 100644 --- a/test/js/node/http/node-http-req-socket-pause.test.ts +++ b/test/js/node/http/node-http-req-socket-pause.test.ts @@ -3,9 +3,10 @@ */ import { describe, expect, it } from "bun:test"; import { once } from "node:events"; -import { Agent, createServer, request, type Server } from "node:http"; +import { Agent, createServer, request, type IncomingMessage, type Server } from "node:http"; import type { AddressInfo, Socket } from "node:net"; import { connect } from "node:net"; +import type { Duplex } from "node:stream"; it("req.socket emits 'pause' once an unread request body fills the IncomingMessage buffer", async () => { // Node's test-http-no-read-no-dump: a handler that never reads the body sees @@ -294,3 +295,164 @@ it("upgrade request whose whole body arrived while it was paused still hands the if (server.listening) server.close(); } }); + +describe("upgrade request whose whole body arrived with its head", () => { + const upgradeHeaders = "Upgrade: test\r\nConnection: Upgrade\r\n"; + const switchingProtocols = `HTTP/1.1 101 Switching Protocols\r\n${upgradeHeaders}\r\n`; + const body = Buffer.alloc(100, "y").toString(); + const fixedLengthPost = `POST /upgrade HTTP/1.1\r\nHost: a\r\n${upgradeHeaders}Content-Length: ${body.length}\r\n\r\n${body}`; + + /** Collects what the upgrade socket receives. */ + function tunnelReader() { + let bytes = ""; + let wake: (() => void) | undefined; + return { + onData(chunk: Buffer) { + bytes += chunk; + wake?.(); + }, + get bytes() { + return bytes; + }, + async receives(expected: string) { + while (bytes.length < expected.length) { + const { promise, resolve } = Promise.withResolvers(); + wake = resolve; + await promise; + } + }, + }; + } + + for (const readInListener of [true, false]) { + it(`a paused request keeps its body when the upgrade socket is read ${readInListener ? "in" : "after"} the listener`, async () => { + const tunnel = tunnelReader(); + const { promise: handedOff, resolve: onUpgrade } = Promise.withResolvers<[IncomingMessage, Duplex]>(); + const server = createServer(); + server.on("upgrade", (req, socket) => { + req.pause(); + socket.on("error", () => {}); + if (readInListener) socket.on("data", tunnel.onData); + socket.write(switchingProtocols); + onUpgrade([req, socket]); + }); + let client: Awaited> | undefined; + try { + client = await connectTo(server); + client.socket.write(fixedLengthPost); + const [req, socket] = await handedOff; + // The 101 is a round trip: the server has parsed the whole first read. + await client.receive("101 Switching Protocols"); + client.socket.write("ping-1;"); + if (!readInListener) socket.on("data", tunnel.onData); + await tunnel.receives("ping-1;"); + expect(req.readableFlowing).toBe(false); + + let received = ""; + req.on("data", chunk => (received += chunk)); + req.resume(); + await once(req, "end"); + expect(received).toBe(body); + + client.socket.write("ping-2;"); + await tunnel.receives("ping-1;ping-2;"); + expect(tunnel.bytes).toBe("ping-1;ping-2;"); + } finally { + client?.socket.destroy(); + server.closeAllConnections(); + if (server.listening) server.close(); + } + }); + } + + it("the upgrade socket gets the bytes after a body that paused the connection", async () => { + // As above, the first chunk fills the request's buffer, and the end of the body is + // received while the connection is paused. Node.js v26.3.0 delivers the body, but + // never reads the socket again. + const tunnel = tunnelReader(); + const { promise: handedOff, resolve: onUpgrade } = Promise.withResolvers<[IncomingMessage, Duplex]>(); + const server = createServer({ highWaterMark: 1024 }); + server.on("upgrade", (req, socket) => { + socket.on("error", () => {}); + socket.write(switchingProtocols); + onUpgrade([req, socket]); + }); + let client: Awaited> | undefined; + try { + client = await connectTo(server); + client.socket.write(chunkedPost("/upgrade", upgradeHeaders)); + const [req, socket] = await handedOff; + await client.receive("101 Switching Protocols"); + + let received = ""; + req.on("data", chunk => (received += chunk)); + await once(req, "end"); + expect(received).toBe(BODY_HEAD + BODY_TAIL); + + client.socket.write("ping;"); + socket.on("data", tunnel.onData); + await tunnel.receives("ping;"); + expect(tunnel.bytes).toBe("ping;"); + } finally { + client?.socket.destroy(); + server.closeAllConnections(); + if (server.listening) server.close(); + } + }); + + it("a read of the upgrade socket still resumes a paused request whose body is incomplete", async () => { + // Like Node.js's UpgradeStream._read: the listener never resumes the request, and the part + // of the body that came with the head fills its buffer. The socket still gets its bytes. + const tunnel = tunnelReader(); + const server = createServer({ highWaterMark: 1024 }); + server.on("upgrade", (req, socket) => { + req.pause(); + socket.on("error", () => {}); + socket.on("data", tunnel.onData); + socket.write(switchingProtocols); + }); + const request = chunkedPost("/upgrade", upgradeHeaders); + const cut = request.lastIndexOf(`${BODY_TAIL.length.toString(16)}\r\n${BODY_TAIL}`); + let client: Awaited> | undefined; + try { + client = await connectTo(server); + client.socket.write(request.slice(0, cut)); + await client.receive("101 Switching Protocols"); + client.socket.write(request.slice(cut) + "ping;"); + await tunnel.receives("ping;"); + expect(tunnel.bytes).toBe("ping;"); + } finally { + client?.socket.destroy(); + server.closeAllConnections(); + if (server.listening) server.close(); + } + }); + + it("a read of the upgrade socket still resumes a request that is not paused", async () => { + // Inside the listener req.complete is still false here (Node.js has true), so a listener + // that waits for the end of the message depends on this resume. + const { promise: completeWhenAccepted, resolve: onAccept } = Promise.withResolvers(); + const server = createServer(); + server.on("upgrade", (req, socket) => { + socket.on("error", () => {}); + socket.on("data", () => {}); + const accept = () => { + socket.write(switchingProtocols); + onAccept(req.complete); + }; + if (req.complete) accept(); + else req.once("end", accept); + }); + let client: Awaited> | undefined; + try { + client = await connectTo(server); + client.socket.write(fixedLengthPost); + expect(await completeWhenAccepted).toBe(true); + await client.receive("101 Switching Protocols"); + } finally { + client?.socket.destroy(); + server.closeAllConnections(); + if (server.listening) server.close(); + } + }); +}); From c9e34cc8f56d08dba22837b20408e31effce1fa9 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 22:14:52 +0000 Subject: [PATCH 2/4] node:http: tell an upgraded socket by its kind, and wire failures in the new tests The resume at the switch to tunnel mode skipped a socket that upgrade() adopted inside the body callback by comparing it to upgradedWebSocket. That field is per context, and another connection's upgrade() in the same drain overwrites it. The check now reads the kind of the socket. The new tests reject their waits when a socket errors or closes, and the late-read case checks req.readableFlowing right after the read. --- packages/bun-uws/src/HttpContext.h | 5 +- .../http/node-http-req-socket-pause.test.ts | 59 +++++++++++++------ 2 files changed, 44 insertions(+), 20 deletions(-) diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index 27f9ad3f5c38..4ced5f98bed8 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -658,8 +658,9 @@ struct HttpContext { } } if constexpr (IsNodeHttp) { - /* Not after upgrade() from the body handler: the ext holds a WebSocketData then. */ - if (switchToTunnelAfterThisChunk && httpContextData->upgradedWebSocket != user) { + /* Not after upgrade() from the body handler: the socket is a WebSocket then, and + * its ext holds a WebSocketData. upgradedWebSocket can name another connection. */ + if (switchToTunnelAfterThisChunk && us_socket_kind((struct us_socket_t *) user) == socketKind()) { /* pause() and resume() on the response do nothing in tunnel mode: * lift a read pause that the request body left (req.pause(), a full buffer). */ Bun__NodeHTTP__onReadsResumable(SSL, (struct us_socket_t *) user); diff --git a/test/js/node/http/node-http-req-socket-pause.test.ts b/test/js/node/http/node-http-req-socket-pause.test.ts index be150595a3d6..888412540d48 100644 --- a/test/js/node/http/node-http-req-socket-pause.test.ts +++ b/test/js/node/http/node-http-req-socket-pause.test.ts @@ -2,7 +2,7 @@ * All tests in this file should also run in Node.js. */ import { describe, expect, it } from "bun:test"; -import { once } from "node:events"; +import { once, type EventEmitter } from "node:events"; import { Agent, createServer, request, type IncomingMessage, type Server } from "node:http"; import type { AddressInfo, Socket } from "node:net"; import { connect } from "node:net"; @@ -324,14 +324,28 @@ describe("upgrade request whose whole body arrived with its head", () => { }; } + /** `orFail(p)` settles like `p`, or rejects when a watched socket errors or closes first. */ + function failureWatcher() { + const { promise: failed, reject } = Promise.withResolvers(); + failed.catch(() => {}); // the teardown closes the sockets after the last race + return { + watch(what: string, socket: EventEmitter) { + socket.on("error", reject); + socket.on("close", () => reject(new Error(`${what} closed`))); + }, + orFail: (promise: Promise) => Promise.race([promise, failed]), + }; + } + for (const readInListener of [true, false]) { it(`a paused request keeps its body when the upgrade socket is read ${readInListener ? "in" : "after"} the listener`, async () => { const tunnel = tunnelReader(); + const { watch, orFail } = failureWatcher(); const { promise: handedOff, resolve: onUpgrade } = Promise.withResolvers<[IncomingMessage, Duplex]>(); const server = createServer(); server.on("upgrade", (req, socket) => { req.pause(); - socket.on("error", () => {}); + watch("the upgrade socket", socket); if (readInListener) socket.on("data", tunnel.onData); socket.write(switchingProtocols); onUpgrade([req, socket]); @@ -339,23 +353,26 @@ describe("upgrade request whose whole body arrived with its head", () => { let client: Awaited> | undefined; try { client = await connectTo(server); + watch("the client socket", client.socket); client.socket.write(fixedLengthPost); - const [req, socket] = await handedOff; + const [req, socket] = await orFail(handedOff); // The 101 is a round trip: the server has parsed the whole first read. - await client.receive("101 Switching Protocols"); + await orFail(client.receive("101 Switching Protocols")); client.socket.write("ping-1;"); if (!readInListener) socket.on("data", tunnel.onData); - await tunnel.receives("ping-1;"); + // The read of the socket does not resume the request, at once or a turn later. + expect(req.readableFlowing).toBe(false); + await orFail(tunnel.receives("ping-1;")); expect(req.readableFlowing).toBe(false); let received = ""; req.on("data", chunk => (received += chunk)); req.resume(); - await once(req, "end"); + await orFail(once(req, "end")); expect(received).toBe(body); client.socket.write("ping-2;"); - await tunnel.receives("ping-1;ping-2;"); + await orFail(tunnel.receives("ping-1;ping-2;")); expect(tunnel.bytes).toBe("ping-1;ping-2;"); } finally { client?.socket.destroy(); @@ -370,28 +387,30 @@ describe("upgrade request whose whole body arrived with its head", () => { // received while the connection is paused. Node.js v26.3.0 delivers the body, but // never reads the socket again. const tunnel = tunnelReader(); + const { watch, orFail } = failureWatcher(); const { promise: handedOff, resolve: onUpgrade } = Promise.withResolvers<[IncomingMessage, Duplex]>(); const server = createServer({ highWaterMark: 1024 }); server.on("upgrade", (req, socket) => { - socket.on("error", () => {}); + watch("the upgrade socket", socket); socket.write(switchingProtocols); onUpgrade([req, socket]); }); let client: Awaited> | undefined; try { client = await connectTo(server); + watch("the client socket", client.socket); client.socket.write(chunkedPost("/upgrade", upgradeHeaders)); - const [req, socket] = await handedOff; - await client.receive("101 Switching Protocols"); + const [req, socket] = await orFail(handedOff); + await orFail(client.receive("101 Switching Protocols")); let received = ""; req.on("data", chunk => (received += chunk)); - await once(req, "end"); + await orFail(once(req, "end")); expect(received).toBe(BODY_HEAD + BODY_TAIL); client.socket.write("ping;"); socket.on("data", tunnel.onData); - await tunnel.receives("ping;"); + await orFail(tunnel.receives("ping;")); expect(tunnel.bytes).toBe("ping;"); } finally { client?.socket.destroy(); @@ -404,10 +423,11 @@ describe("upgrade request whose whole body arrived with its head", () => { // Like Node.js's UpgradeStream._read: the listener never resumes the request, and the part // of the body that came with the head fills its buffer. The socket still gets its bytes. const tunnel = tunnelReader(); + const { watch, orFail } = failureWatcher(); const server = createServer({ highWaterMark: 1024 }); server.on("upgrade", (req, socket) => { req.pause(); - socket.on("error", () => {}); + watch("the upgrade socket", socket); socket.on("data", tunnel.onData); socket.write(switchingProtocols); }); @@ -416,10 +436,11 @@ describe("upgrade request whose whole body arrived with its head", () => { let client: Awaited> | undefined; try { client = await connectTo(server); + watch("the client socket", client.socket); client.socket.write(request.slice(0, cut)); - await client.receive("101 Switching Protocols"); + await orFail(client.receive("101 Switching Protocols")); client.socket.write(request.slice(cut) + "ping;"); - await tunnel.receives("ping;"); + await orFail(tunnel.receives("ping;")); expect(tunnel.bytes).toBe("ping;"); } finally { client?.socket.destroy(); @@ -431,10 +452,11 @@ describe("upgrade request whose whole body arrived with its head", () => { it("a read of the upgrade socket still resumes a request that is not paused", async () => { // Inside the listener req.complete is still false here (Node.js has true), so a listener // that waits for the end of the message depends on this resume. + const { watch, orFail } = failureWatcher(); const { promise: completeWhenAccepted, resolve: onAccept } = Promise.withResolvers(); const server = createServer(); server.on("upgrade", (req, socket) => { - socket.on("error", () => {}); + watch("the upgrade socket", socket); socket.on("data", () => {}); const accept = () => { socket.write(switchingProtocols); @@ -446,9 +468,10 @@ describe("upgrade request whose whole body arrived with its head", () => { let client: Awaited> | undefined; try { client = await connectTo(server); + watch("the client socket", client.socket); client.socket.write(fixedLengthPost); - expect(await completeWhenAccepted).toBe(true); - await client.receive("101 Switching Protocols"); + expect(await orFail(completeWhenAccepted)).toBe(true); + await orFail(client.receive("101 Switching Protocols")); } finally { client?.socket.destroy(); server.closeAllConnections(); From c74ed4543db50f1bd6fa06c8351e66781baaca96 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 22:36:31 +0000 Subject: [PATCH 3/4] node:http: shorten the comments of the upgrade body resume --- packages/bun-uws/src/HttpContext.h | 6 ++---- src/js/node/_http_server.ts | 7 ++----- 2 files changed, 4 insertions(+), 9 deletions(-) diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index 4ced5f98bed8..09513eb23a3e 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -658,11 +658,9 @@ struct HttpContext { } } if constexpr (IsNodeHttp) { - /* Not after upgrade() from the body handler: the socket is a WebSocket then, and - * its ext holds a WebSocketData. upgradedWebSocket can name another connection. */ + /* The kind check: upgrade() from the body handler turns the ext into a WebSocketData. */ if (switchToTunnelAfterThisChunk && us_socket_kind((struct us_socket_t *) user) == socketKind()) { - /* pause() and resume() on the response do nothing in tunnel mode: - * lift a read pause that the request body left (req.pause(), a full buffer). */ + /* The response cannot resume reads in tunnel mode: lift the pause the body left. */ Bun__NodeHTTP__onReadsResumable(SSL, (struct us_socket_t *) user); if (us_socket_is_closed((struct us_socket_t *) user)) { return nullptr; diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index e3165f0c2ff6..81da49c44a23 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -1319,9 +1319,7 @@ function clearUpgradeIncoming(socket) { socket[kUpgradeIncoming] = undefined; } -// Like Node.js, a read of the raw socket resumes a paused request only while its body is -// incomplete. Checked one event loop turn after the read: the 'upgrade' listener runs before -// the rest of its read is parsed, and that rest can complete the body. +// Deferred: the 'upgrade' listener runs before the rest of its read, which can complete the body. function resumePausedUpgradeIncoming(socket) { const req = socket[kUpgradeIncoming]; if (req === undefined) return; @@ -1793,8 +1791,7 @@ function getNodeHTTPServerSocket() { upgradeIncoming.push(resumed); } } - // Not paused: resumed at once, also with a complete body. req.complete is still false - // inside the listener then, so a listener can wait for 'end' and never read the body. + // Not paused: req.complete is false in the listener, so it can wait for 'end' with no reader. if (upgradeIncoming.readableFlowing !== false) upgradeIncoming.resume(); else setImmediate(resumePausedUpgradeIncoming, this); return; From 58ef47d4d086853bbe0fbc45cbf149849d3e1bbb Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 19 Sep 2026 22:52:50 +0000 Subject: [PATCH 4/4] node:http: use it.each for the two paused Upgrade request cases --- test/js/node/http/node-http-req-socket-pause.test.ts | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/test/js/node/http/node-http-req-socket-pause.test.ts b/test/js/node/http/node-http-req-socket-pause.test.ts index 888412540d48..75425bb423c5 100644 --- a/test/js/node/http/node-http-req-socket-pause.test.ts +++ b/test/js/node/http/node-http-req-socket-pause.test.ts @@ -337,8 +337,12 @@ describe("upgrade request whose whole body arrived with its head", () => { }; } - for (const readInListener of [true, false]) { - it(`a paused request keeps its body when the upgrade socket is read ${readInListener ? "in" : "after"} the listener`, async () => { + it.each([ + { when: "in", readInListener: true }, + { when: "after", readInListener: false }, + ])( + "a paused request keeps its body when the upgrade socket is read $when the listener", + async ({ readInListener }) => { const tunnel = tunnelReader(); const { watch, orFail } = failureWatcher(); const { promise: handedOff, resolve: onUpgrade } = Promise.withResolvers<[IncomingMessage, Duplex]>(); @@ -379,8 +383,8 @@ describe("upgrade request whose whole body arrived with its head", () => { server.closeAllConnections(); if (server.listening) server.close(); } - }); - } + }, + ); it("the upgrade socket gets the bytes after a body that paused the connection", async () => { // As above, the first chunk fills the request's buffer, and the end of the body is