diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index c3fbda0d72c7..9dc9e49e04b0 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -844,6 +844,12 @@ struct HttpContext { && httpResponseData->onWritable == nullptr) { responseDone = true; } + /* socket.destroySoon() issued while bytes were queued: Node's + * destroy() on 'finish' closes once they are out, whether or + * not the response in flight ever ends. */ + if (httpResponseData->state & HttpResponseData::HTTP_NODE_CLOSE_AFTER_DRAIN) { + responseDone = true; + } } if (responseDone && asyncSocket->hasFullyDrained()) { asyncSocket->shutdown(); diff --git a/packages/bun-uws/src/HttpResponseData.h b/packages/bun-uws/src/HttpResponseData.h index 756007a69d28..daafe79bb79c 100644 --- a/packages/bun-uws/src/HttpResponseData.h +++ b/packages/bun-uws/src/HttpResponseData.h @@ -150,13 +150,18 @@ struct HttpResponseData : AsyncSocketData, HttpParser { * shutdown sweep; the shouldCloseConnection() gates act on it once the * in-flight work completes. */ HTTP_CLOSE_WHEN_IDLE = 1 << 17, + /* node:http socket.destroySoon() with outgoing bytes still queued: shut + * down and close as soon as they have flushed, whether or not the + * response in flight has ended (Node's destroy() on 'finish'). */ + HTTP_NODE_CLOSE_AFTER_DRAIN = 1 << 18, /* Bits that describe the connection rather than the response in flight. * There is one HttpResponseData per socket, reused by every request on a * keep-alive connection, so starting a new response clears the rest of the * word (resetResponseState) - these have to survive that. */ HTTP_CONNECTION_SCOPED = HTTP_NODE_PARSING_STOPPED | HTTP_NODE_READS_PAUSED - | HTTP_NODE_TUNNEL_AFTER_BODY | HTTP_NODE_RECEIVED_FIN | HTTP_CLOSE_WHEN_IDLE, + | HTTP_NODE_TUNNEL_AFTER_BODY | HTTP_NODE_RECEIVED_FIN | HTTP_CLOSE_WHEN_IDLE + | HTTP_NODE_CLOSE_AFTER_DRAIN, }; /* Begin a new response on this connection. Clearing the word in one go is @@ -228,7 +233,7 @@ struct HttpResponseData : AsyncSocketData, HttpParser { /* Whether the connection should be torn down once the in-flight response (if * any) has completed and all buffered outgoing data has been flushed. */ bool shouldCloseConnection() const { - return (state & HTTP_CONNECTION_CLOSE) + return (state & (HTTP_CONNECTION_CLOSE | HTTP_NODE_CLOSE_AFTER_DRAIN)) || ((state & HTTP_NODE_RECEIVED_FIN) && nodeHttpQueuedPipelinedCount == 0) || ((state & HTTP_CLOSE_WHEN_IDLE) && this->isIdle); } diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 45fb072718c6..035742523bfa 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -1298,9 +1298,10 @@ function onServerClientError(ssl: boolean, socket: unknown, errorCode: number, r // reaches socketOnError and 'clientError' never fires for it. function replyMissingHostHeader(socket) { if (!socket.writable) return; - socket.end( + socket.write( `HTTP/1.1 400 Bad Request\r\nConnection: close\r\nDate: ${new Date().toUTCString()}\r\nTransfer-Encoding: chunked\r\n\r\n0\r\n\r\n`, ); + socket.destroySoon(); } const kBytesWritten = Symbol("kBytesWritten"); @@ -1333,6 +1334,8 @@ function resolveHandoffPromise(promise) { } const kSocketTimeoutTimer = Symbol("socketTimeoutTimer"); const kStreamingEnabled = Symbol("kStreamingEnabled"); +// destroySoon() was called: _final closes the connection right behind the FIN. +const kDestroySoon = Symbol("kDestroySoon"); // Scratch options object for the builtin ServerResponse (see the dispatcher). const scratchResponseOptions = { [kHandle]: undefined, @@ -1485,6 +1488,7 @@ function getNodeHTTPServerSocket() { [kBytesWritten] = 0; [kHandle]; [kUpgradeIncoming] = undefined; + [kDestroySoon] = false; server: Server; _httpMessage; _secureEstablished = false; @@ -1765,10 +1769,22 @@ function getNodeHTTPServerSocket() { callback(); return; } - handle.end(); + handle.end(this[kDestroySoon]); callback(); } + // Not destroy() on 'finish' like net.Socket: 'finish' does not wait for bytes uWS still queues. + destroySoon() { + if (this[kDestroySoon]) return; + const handle = this[kHandle]; + if (this.writable && handle && !handle.closed) { + this[kDestroySoon] = true; + this.end(); + return; + } + super.destroySoon(); + } + get localAddress() { return this[kHandle]?.localAddress?.address; } @@ -2405,7 +2421,13 @@ function emitResponseFinish() { // is eventually closed. function onResponseFinishHandleSocket(server, socket, res) { if (res[kMustCloseConnection]) { - socket?.end(); + if (socket != null) { + if (typeof socket.destroySoon === "function") { + socket.destroySoon(); + } else { + socket.end(); + } + } return; } if (!socket || socket.destroyed || typeof socket.setTimeout !== "function") { diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp index 92aba2c148f1..86cb0f69e94b 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp @@ -243,7 +243,7 @@ bool JSNodeHTTPServerSocket::isClosed() const } template -static bool deferShutdownUntilResponseDrains(us_socket_t* socket) +static bool deferShutdownUntilResponseDrains(us_socket_t* socket, bool thenClose) { if (reinterpret_cast*>(socket)->getBufferedAmount() == 0) { return false; @@ -253,18 +253,22 @@ static bool deferShutdownUntilResponseDrains(us_socket_t* socket) * is sequenced after the response bytes (like Node's destroySoon). */ auto* httpResponseData = reinterpret_cast*>(us_socket_ext(socket)); httpResponseData->state |= uWS::HttpResponseData::HTTP_CONNECTION_CLOSE; + if (thenClose) { + /* destroySoon(): and closes it there, even if the response never ends. */ + httpResponseData->state |= uWS::HttpResponseData::HTTP_NODE_CLOSE_AFTER_DRAIN; + } return true; } -bool JSNodeHTTPServerSocket::shutdownAfterResponseDrains() +bool JSNodeHTTPServerSocket::shutdownAfterResponseDrains(bool thenClose) { if (!socket || upgraded || us_socket_is_closed(socket) || us_socket_is_shut_down(socket)) { return false; } if (is_ssl) { - return deferShutdownUntilResponseDrains(socket); + return deferShutdownUntilResponseDrains(socket, thenClose); } - return deferShutdownUntilResponseDrains(socket); + return deferShutdownUntilResponseDrains(socket, thenClose); } template diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h index eef438fce6b5..5eaac4ad8ce9 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocket.h +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocket.h @@ -106,7 +106,7 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject { /* node:http socket.end(): when the in-flight response still has bytes in * uWS's send buffer, a shutdown now would put the FIN ahead of them and * truncate the response. Returns true after handing the close to uWS. */ - bool shutdownAfterResponseDrains(); + bool shutdownAfterResponseDrains(bool thenClose); /* Switch the connection into CONNECT-style tunnel mode after an accepted * Upgrade: subsequent bytes bypass the HTTP parser and stream to the diff --git a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp index bdfc8aca98a0..d3e7201de8ab 100644 --- a/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp +++ b/src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp @@ -13,6 +13,7 @@ extern "C" uint64_t uws_res_get_remote_address_info(void* res, const char** dest extern "C" uint64_t uws_res_get_local_address_info(void* res, const char** dest, int* port, bool* is_ipv6); extern "C" void us_socket_resume(us_socket_t*); extern "C" void us_socket_pause(us_socket_t*); +extern "C" void us_socket_shutdown(us_socket_t*); namespace Bun { @@ -216,12 +217,24 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketEnd, (JSC::JSGlobalObject } thisObject->ended = true; + // end(true) is Node's destroySoon(): close right behind the FIN instead of waiting for the peer's. + bool destroySoon = callFrame->argument(0).toBoolean(globalObject); // The response's buffered body must reach the kernel before the FIN; uWS // performs the shutdown after its send buffer drains. - if (thisObject->shutdownAfterResponseDrains()) { + if (thisObject->shutdownAfterResponseDrains(destroySoon)) { return JSValue::encode(JSC::jsUndefined()); } auto bufferedSize = thisObject->streamBuffer.bufferedSize(); + if (destroySoon) { + // One flush for raw socket.write() bytes; the close drops the rest, as destroy() did. + if (bufferedSize == 0) { + us_socket_shutdown(thisObject->socket); + } else { + us_socket_buffered_js_write(thisObject->socket, thisObject->is_ssl, thisObject->ended, &thisObject->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined())); + } + thisObject->close(); + return JSValue::encode(JSC::jsUndefined()); + } if (bufferedSize == 0) { // onNodeHTTPRequest no longer pauses at dispatch; pause here so the // shutdown+resume below still cycles kqueue's EVFILT_READ (delete then diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index c555aa1b5576..a221a2d744ad 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -2708,34 +2708,142 @@ it("removing only Content-Length falls back to chunked encoding and keeps the co } }); -it("an explicit Connection: close response header closes the server-side socket after finish", async () => { - // Node.js's matchHeader sets _last for a user-set Connection: close, and - // resOnFinish then ends the socket; the transport must match the header. - const server = createServer((req, res) => { - res.setHeader("Connection", "close"); - res.end("ok"); - }); - try { +// Node.js's resOnFinish destroySoon()s the socket after a response that ends the +// connection: the server sends its FIN and releases the socket right behind it +// ('close' on the server-side socket, fd gone) without waiting for the peer to +// close its side. A half-close that waits for the client's FIN lets a client +// that never closes pin one socket per completed request. +describe("a response that closes the connection releases the server-side socket without waiting for the peer", () => { + const cases = [ + { + name: "Connection: close response header", + request: "GET / HTTP/1.1\r\nHost: localhost\r\n\r\n", + prepare: (res: ServerResponse) => res.setHeader("Connection", "close"), + expected: /^HTTP\/1\.1 200 OK\r\n[^]*ok$/, + }, + { + name: "Connection: close request header", + request: "GET / HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n", + expected: /^HTTP\/1\.1 200 OK\r\n[^]*ok$/, + }, + { + name: "HTTP/1.0 request", + request: "GET / HTTP/1.0\r\nHost: localhost\r\n\r\n", + expected: /^HTTP\/1\.1 200 OK\r\n[^]*ok$/, + }, + { + // Node answers this itself with res.writeHead(400, ['Connection', 'close']). + name: "HTTP/1.1 request without a Host header", + request: "GET / HTTP/1.1\r\n\r\n", + expected: /^HTTP\/1\.1 400 Bad Request\r\n[^]*\r\n0\r\n\r\n$/, + }, + ]; + for (const { name, request, prepare, expected } of cases) { + it.concurrent(name, async () => { + const { promise: serverSocketClosed, resolve: onServerSocketClose } = Promise.withResolvers(); + const server = createServer((req, res) => { + prepare?.(res); + res.end("ok"); + }); + server.on("connection", socket => socket.on("close", onServerSocketClose)); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const { port } = server.address() as AddressInfo; + // allowHalfOpen: the client keeps its side open after the server's FIN, so + // only the server can release the server-side socket. + const client = connect({ port, host: "127.0.0.1", allowHalfOpen: true }); + try { + const out = await new Promise((resolve, reject) => { + let data = ""; + client.setEncoding("latin1"); + client.on("data", chunk => (data += chunk)); + client.on("end", () => resolve(data)); + client.on("error", reject); + client.write(request); + }); + expect(out).toMatch(expected); + expect(out).toContain("\r\nConnection: close\r\n"); + // The client has the whole response and the server's FIN, and has not + // sent its own. The server-side socket must close on its own now. + await serverSocketClosed; + expect(client.writableEnded).toBe(false); + } finally { + client.destroy(); + server.close(); + } + }); + } + + // net.Socket.destroySoon() is what Node's resOnFinish uses: end(), then + // destroy() once the write side is done, without waiting for the peer. + it.concurrent("socket.destroySoon() from a handler", async () => { + const { promise: serverSocketClosed, resolve: onServerSocketClose } = Promise.withResolvers(); + const server = createServer((req, res) => { + req.socket.write("raw bytes\r\n"); + req.socket.destroySoon(); + }); + server.on("connection", socket => socket.on("close", onServerSocketClose)); server.listen(0, "127.0.0.1"); await once(server, "listening"); const { port } = server.address() as AddressInfo; + const client = connect({ port, host: "127.0.0.1", allowHalfOpen: true }); + try { + const out = await new Promise((resolve, reject) => { + let data = ""; + client.setEncoding("latin1"); + client.on("data", chunk => (data += chunk)); + client.on("end", () => resolve(data)); + client.on("error", reject); + client.write("GET / HTTP/1.1\r\nHost: localhost\r\n\r\n"); + }); + expect(out).toBe("raw bytes\r\n"); + await serverSocketClosed; + expect(client.writableEnded).toBe(false); + } finally { + client.destroy(); + server.close(); + } + }); - const out = await new Promise((resolve, reject) => { - const socket = connect(port, "127.0.0.1"); - let data = ""; - socket.on("data", chunk => (data += chunk)); - // The server must send FIN on its own; the client never half-closes. - socket.on("end", () => resolve(data)); - socket.on("error", reject); - socket.write("GET / HTTP/1.1\r\nHost: localhost\r\n\r\n"); + // destroySoon() with response bytes still queued in the transport (the client + // is not reading yet) and no res.end(): the queued bytes are delivered first, + // then the connection closes, like Node's destroy() on the socket's 'finish'. + it.concurrent("socket.destroySoon() behind a backed-up response that never ends", async () => { + const CHUNK = Buffer.alloc(16 * 1024, "a"); + const CHUNKS = 1024; // 16 MB: more than the loopback socket buffers absorb + const { promise: serverSocketClosed, resolve: onServerSocketClose } = Promise.withResolvers(); + const { promise: destroyed, resolve: onDestroySoonCalled } = Promise.withResolvers(); + const server = createServer((req, res) => { + res.on("error", () => {}); + for (let i = 0; i < CHUNKS; i++) res.write(CHUNK); + req.socket.destroySoon(); + onDestroySoonCalled(); }); - - expect(out).toContain("HTTP/1.1 200"); - expect(out).toContain("Connection: close"); - expect(out).toEndWith("ok"); - } finally { - server.close(); - } + server.on("connection", socket => socket.on("close", onServerSocketClose)); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const { port } = server.address() as AddressInfo; + const client = connect({ port, host: "127.0.0.1", allowHalfOpen: true }); + try { + let received = 0; + const ended = new Promise((resolve, reject) => { + client.on("data", chunk => (received += chunk.length)); + client.on("end", resolve); + client.on("error", reject); + }); + client.pause(); + client.write("GET / HTTP/1.1\r\nHost: localhost\r\n\r\n"); + await destroyed; + client.resume(); + await ended; + expect(received).toBeGreaterThan(CHUNK.length * CHUNKS); + await serverSocketClosed; + expect(client.writableEnded).toBe(false); + } finally { + client.destroy(); + server.close(); + } + }); }); it("a pipelined request behind Connection: close is never dispatched (clientError HPE_CLOSED_CONNECTION)", async () => {