diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 1e0bcc316d13..25069be08790 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -755,6 +755,9 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // Node.js's parserOnIncoming: req.upgrade is true for CONNECT // regardless of shouldUpgradeCallback. http_req.upgrade = true; + // llhttp completes a CONNECT request at the end of its headers, so + // Node's 'connect' listener already sees req.complete === true. + http_req.complete = true; // Node frees the parser before handing the raw socket to 'connect'. releaseServerParserShim(socket, http_req); server.emit("connect", http_req, socket, head); @@ -970,6 +973,11 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort if (hasBody) { socket[kUpgradeIncoming] = http_req; http_req.once("end", clearUpgradeIncoming.bind(undefined, socket)); + } else { + // llhttp completes an Upgrade request without a body at the end + // of its headers, so Node's 'upgrade' listener already sees + // req.complete === true. + http_req.complete = true; } const upgradeHead = !hasBody && connectHead ? connectHead : kEmptyBuffer; let upgradeHandled; @@ -2369,8 +2377,17 @@ function emitResponseFinish() { // req.socket is nulled by the stream destroyer (pipeline/compose cleanup); // the response's own socket (set by assignSocket, cleared only by // detachSocket) still references the connection then. - const socket = this.req?.socket ?? this.socket; + const req = this.req; + const socket = req?.socket ?? this.socket; onResponseFinishHandleSocket(socket?.server, socket, this); + // Like Node's clearIncoming: a request that already ended (for example one + // that optimizeEmptyRequests pre-dumped, which never reaches + // emitEOFIncomingMessageOuter) must not stay reachable from socket.parser + // while the kept-alive connection idles. + if (req != null && req.readableEnded) { + const parser = socket?.parser; + if (parser != null && parser.incoming === req) parser.incoming = null; + } // The dispatcher detached a synchronously-finished response itself; // advancing the pipeline again here would skip a queued response. if (this[kDispatcherDetached]) return; diff --git a/src/runtime/server/NodeHTTPResponse.rs b/src/runtime/server/NodeHTTPResponse.rs index 15977391c7e4..69609d8c4ba8 100644 --- a/src/runtime/server/NodeHTTPResponse.rs +++ b/src/runtime/server/NodeHTTPResponse.rs @@ -7,7 +7,6 @@ use bstr::BStr; use bun_collections::VecExt; use bun_core::scoped_log; -use bun_http::Method as HttpMethod; use bun_jsc::JsCell; use bun_ptr::AsCtxPtr; use bun_uws as uws; @@ -2596,23 +2595,22 @@ pub(crate) unsafe extern "C" fn NodeHTTPResponse__createForJS( let request_ref = bun_opaque::opaque_deref(request.cast_const()); let vm = bun_vm_mut(global_object); - let method = HttpMethod::which(request_ref.method()).unwrap_or(HttpMethod::OPTIONS); - // GET in node.js can have a body - if method.has_request_body() || method == HttpMethod::GET { - let req_len: usize = 'brk: { - if let Some(content_length) = request_ref.header(b"content-length") { - scoped_log!( - NodeHTTPResponse, - "content-length: {}", - BStr::new(content_length) - ); - break 'brk bun_http_types::parse_content_length(content_length); - } - break 'brk 0; - }; + // Like llhttp, the framing headers decide whether a request has a body + // for every method: node delivers the body of a HEAD or TRACE request + // that declares one. + let req_len: usize = 'brk: { + if let Some(content_length) = request_ref.header(b"content-length") { + scoped_log!( + NodeHTTPResponse, + "content-length: {}", + BStr::new(content_length) + ); + break 'brk bun_http_types::parse_content_length(content_length); + } + break 'brk 0; + }; - *has_body = req_len > 0 || request_ref.has_transfer_encoding(); - } + *has_body = req_len > 0 || request_ref.has_transfer_encoding(); let raw_response = if is_ssl != 0 { uws::AnyResponse::SSL(response_ptr.cast()) diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index 71ad7e6435d6..aeec37b0aaa4 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -4668,3 +4668,94 @@ it("connectionListener pauses reads when queued pipelined responses back up", as clientSide.destroy(); serverSide.destroy(); }); + +describe("request completion (req.complete, socket.parser.incoming)", () => { + function rawClient(server: Server, request: string) { + const { port, address } = server.address() as AddressInfo; + const client = connect(port, address, () => client.write(request)); + client.on("error", () => {}); + client.on("data", () => {}); + return client; + } + + test("req.complete is true inside a 'connect' listener", async () => { + await using server = createServer(); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + let sawComplete: boolean | undefined; + server.on("connect", (req, socket) => { + sawComplete = req.complete; + socket.destroy(); + }); + const client = rawClient(server, "CONNECT example.com:443 HTTP/1.1\r\nHost: example.com:443\r\n\r\n"); + const [req] = await once(server, "connect"); + expect(sawComplete).toBe(true); + await new Promise(resolve => process.nextTick(resolve)); + expect(req.complete).toBe(true); + await once(client, "close"); + }); + + test("req.complete is true inside an 'upgrade' listener for a request without a body", async () => { + await using server = createServer(); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + let sawComplete: boolean | undefined; + server.on("upgrade", (req, socket) => { + sawComplete = req.complete; + socket.destroy(); + }); + const client = rawClient(server, "GET / HTTP/1.1\r\nHost: x\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\n"); + const [req] = await once(server, "upgrade"); + expect(sawComplete).toBe(true); + await new Promise(resolve => process.nextTick(resolve)); + expect(req.complete).toBe(true); + await once(client, "close"); + }); + + for (const method of ["HEAD", "TRACE"]) { + test(`a ${method} request that declares a body delivers it and completes only once it arrived`, async () => { + await using server = createServer(); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const chunks: Buffer[] = []; + const body = Promise.withResolvers(); + server.on("request", (req, res) => { + req.on("data", c => chunks.push(c)); + req.on("end", () => { + body.resolve(Buffer.concat(chunks).toString()); + res.end(); + }); + }); + // The headers go first. The declared body follows only once the + // request has been dispatched. + const client = rawClient(server, `${method} / HTTP/1.1\r\nHost: x\r\nContent-Length: 5\r\n\r\n`); + const [req] = await once(server, "request"); + await new Promise(resolve => process.nextTick(resolve)); + expect(req.complete).toBe(false); + client.write("hello"); + expect(await body.promise).toBe("hello"); + expect(req.complete).toBe(true); + client.destroy(); + await once(client, "close"); + }); + } + + test("optimizeEmptyRequests: socket.parser.incoming does not keep the request once the response finished", async () => { + const closed = Promise.withResolvers<{ atRequest: boolean; atClose: boolean }>(); + await using server = createServer({ optimizeEmptyRequests: true }, (req, res) => { + const socket = req.socket; + const atRequest = socket.parser.incoming === req; + // The keep-alive connection idles after this response. Node cleared + // parser.incoming in resOnFinish because the pre-dumped request had + // already ended. + res.on("close", () => closed.resolve({ atRequest, atClose: socket.parser.incoming === req })); + res.end("ok"); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const client = rawClient(server, "GET / HTTP/1.1\r\nHost: x\r\n\r\n"); + expect(await closed.promise).toEqual({ atRequest: true, atClose: false }); + client.destroy(); + await once(client, "close"); + }); +});