diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 282248b5675b..623d85be3cbf 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -86,6 +86,7 @@ const OutgoingMessagePrototype = OutgoingMessage.prototype; const { kIncomingMessage } = require("node:_http_common"); const kConnectionsCheckingInterval = Symbol("http.server.connectionsCheckingInterval"); const kTrackedConnections = Symbol("http.server.trackedConnections"); +const kPendingDrainClose = Symbol("http.server.pendingDrainClose"); const kHttpAllowHalfOpen = Symbol("http.server.httpAllowHalfOpen"); // node.http trace events ('http.server.request' b/e). The agent module is @@ -120,7 +121,14 @@ const DateNow = Date.now; let cluster; function emitCloseServer(self: Server) { - callCloseCallback(self); + // The native all-closed promise tracks pending requests, not open connections. + if (self[serverSymbol]) return; + const connections = self[kTrackedConnections]; + if (connections && connections.size > 0) { + self[kPendingDrainClose] = true; + return; + } + self[kPendingDrainClose] = false; self.emit("close"); } function emitCloseNTServer(this: Server) { @@ -309,6 +317,7 @@ function Server(options, callback): void { defineHttpAllowHalfOpen(this); this[kInternalSocketData] = undefined; this[kTrackedConnections] = new Set(); + this[kPendingDrainClose] = false; this[tlsSymbol] = null; this.noDelay = true; if (typeof options === "function") { @@ -473,14 +482,18 @@ Server.prototype.unref = function () { Server.prototype.closeAllConnections = function () { const server = this[serverSymbol]; - if (!server) { + if (server) { + this[serverSymbol] = undefined; + clearInterval(this[kConnectionsCheckingInterval]); + this.listening = false; + server.stop(true); return; } - this[serverSymbol] = undefined; - clearInterval(this[kConnectionsCheckingInterval]); - this.listening = false; - - server.stop(true); + // close() already dropped the native handle; destroy what is still tracked. + const tracked = this[kTrackedConnections]; + if (tracked && tracked.size > 0) { + for (const socket of Array.from(tracked)) socket.destroy(); + } }; Server.prototype.getConnections = function (callback) { @@ -495,7 +508,16 @@ Server.prototype.getConnections = function (callback) { Server.prototype.closeIdleConnections = function () { const server = this[serverSymbol]; - server?.closeIdleConnections(); + if (server) { + server.closeIdleConnections(); + return; + } + const tracked = this[kTrackedConnections]; + if (tracked && tracked.size > 0) { + for (const socket of Array.from(tracked)) { + if (!socket._httpMessage) socket.destroy(); + } + } }; Server.prototype.close = function (optionalCallback?) { @@ -509,7 +531,7 @@ Server.prototype.close = function (optionalCallback?) { return this; } this[serverSymbol] = undefined; - if (typeof optionalCallback === "function") setCloseCallback(this, optionalCallback); + if (typeof optionalCallback === "function") this.once("close", optionalCallback); this.listening = false; server.closeIdleConnections(); server.stop(); @@ -1148,6 +1170,7 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // }, }); + this[kPendingDrainClose] = false; getBunServerAllClosedPromise(this[serverSymbol]).$then(emitCloseNTServer.bind(this)); isHTTPS = this[serverSymbol].protocol === "https"; applyServerCustomOptions(this); @@ -1667,7 +1690,14 @@ const NodeHTTPServerSocket = class Socket extends NetSocket { // released parser (free() invoked, kOnTimeout nulled). releaseServerParserShim(this); this[kHandle] = null; - this.server?.[kTrackedConnections]?.delete(this); + const server = this.server; + const tracked = server?.[kTrackedConnections]; + if (tracked) { + tracked.delete(this); + if (tracked.size === 0 && server[kPendingDrainClose]) { + process.nextTick(emitCloseServer, server); + } + } const timer = this[kSocketTimeoutTimer]; if (timer) { clearTimeout(timer); diff --git a/test/js/node/http/node-http-server-close-drain.test.ts b/test/js/node/http/node-http-server-close-drain.test.ts new file mode 100644 index 000000000000..c23c9638c1c5 --- /dev/null +++ b/test/js/node/http/node-http-server-close-drain.test.ts @@ -0,0 +1,211 @@ +import { expect, test } from "bun:test"; +import { once } from "node:events"; +import { createServer } from "node:http"; +import type { AddressInfo } from "node:net"; +import { connect } from "node:net"; + +// Node's net.Server#close callback (and the 'close' event) only fires once +// every accepted connection has ended. A connection that was mid-request +// when close() ran stays open after the response is delivered, so the +// callback must be withheld until that connection closes. +test("server.close(cb) does not fire while a keep-alive connection is still open", async () => { + const inHandler = Promise.withResolvers(); + let releaseResponse!: () => void; + const paths: string[] = []; + const server = createServer((req, res) => { + paths.push(req.url as string); + if (paths.length === 1) { + inHandler.resolve(); + releaseResponse = () => res.end("resp:" + req.url); + } else { + res.end("resp:" + req.url); + } + }); + server.keepAliveTimeout = 60000; + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const { port } = server.address() as AddressInfo; + + const socket = connect(port, "127.0.0.1"); + try { + await once(socket, "connect"); + let body = ""; + socket.on("data", chunk => (body += chunk)); + socket.on("error", () => {}); + + // First request: handler is entered but response held until after close(). + socket.write("GET /first HTTP/1.1\r\nHost: x\r\n\r\n"); + await inHandler.promise; + + let closeEventFired = false; + server.once("close", () => (closeEventFired = true)); + const closed = Promise.withResolvers(); + let closeCbFired = false; + server.close(() => { + closeCbFired = true; + closed.resolve(); + }); + + // Handler finishes: response is delivered, connection stays open. + // Yield a few event-loop turns after the bytes arrive so the server's + // "all requests done" task chain has run before the callback is checked. + releaseResponse(); + while (!body.includes("resp:/first")) await once(socket, "data"); + for (let i = 0; i < 4; i++) await new Promise(r => setImmediate(r)); + expect(closeCbFired).toBe(false); + expect(closeEventFired).toBe(false); + + // A second request on the same connection is still served (matching Node), + // and the close callback must still be withheld afterwards. + socket.write("GET /second HTTP/1.1\r\nHost: x\r\n\r\n"); + while (!body.includes("resp:/second")) await once(socket, "data"); + for (let i = 0; i < 4; i++) await new Promise(r => setImmediate(r)); + expect(closeCbFired).toBe(false); + expect(closeEventFired).toBe(false); + expect(paths).toEqual(["/first", "/second"]); + + // Closing the connection drains the server and releases the callback. + socket.destroy(); + await closed.promise; + expect(closeCbFired).toBe(true); + expect(closeEventFired).toBe(true); + } finally { + socket.destroy(); + server.closeAllConnections(); + } +}); + +// Sanity: an idle keep-alive connection at close() time is reaped by +// closeIdleConnections(), so the callback fires promptly like before. +test("server.close(cb) fires once an idle keep-alive connection is reaped", async () => { + const server = createServer((req, res) => res.end("ok")); + server.keepAliveTimeout = 60000; + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const { port } = server.address() as AddressInfo; + + const socket = connect(port, "127.0.0.1"); + try { + await once(socket, "connect"); + let body = ""; + socket.on("data", chunk => (body += chunk)); + socket.on("error", () => {}); + socket.write("GET / HTTP/1.1\r\nHost: x\r\n\r\n"); + while (!body.includes("ok")) await once(socket, "data"); + for (let i = 0; i < 4; i++) await new Promise(r => setImmediate(r)); + + const closed = Promise.withResolvers(); + server.close(() => closed.resolve()); + await closed.promise; + } finally { + socket.destroy(); + server.closeAllConnections(); + } +}); + +// The graceful-drain-with-deadline pattern: close(), then force via +// closeAllConnections() once the caller has waited long enough. The force +// step must work even though close() already dropped the native handle. +test("closeAllConnections() after close() force-drains the withheld callback", async () => { + const inHandler = Promise.withResolvers(); + let releaseResponse!: () => void; + const server = createServer((req, res) => { + inHandler.resolve(); + releaseResponse = () => res.end("ok"); + }); + server.keepAliveTimeout = 60000; + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const { port } = server.address() as AddressInfo; + + const socket = connect(port, "127.0.0.1"); + try { + await once(socket, "connect"); + let body = ""; + socket.on("data", chunk => (body += chunk)); + socket.on("error", () => {}); + socket.write("GET / HTTP/1.1\r\nHost: x\r\n\r\n"); + await inHandler.promise; + + const closed = Promise.withResolvers(); + let closeCbFired = false; + server.close(() => { + closeCbFired = true; + closed.resolve(); + }); + releaseResponse(); + while (!body.includes("ok")) await once(socket, "data"); + for (let i = 0; i < 4; i++) await new Promise(r => setImmediate(r)); + expect(closeCbFired).toBe(false); + + server.closeAllConnections(); + await closed.promise; + expect(closeCbFired).toBe(true); + } finally { + socket.destroy(); + server.closeAllConnections(); + } +}); + +// Re-listening after close() while a keep-alive connection from the previous +// cycle is still open must not fire 'close' on the new (listening) server +// when that old connection finally ends. +test("no 'close' is emitted on a re-listened server when an earlier connection ends", async () => { + const inHandler = Promise.withResolvers(); + let releaseResponse!: () => void; + let requests = 0; + const server = createServer((req, res) => { + if (++requests === 1) { + inHandler.resolve(); + releaseResponse = () => res.end("ok"); + } else { + res.end("ok"); + } + }); + server.keepAliveTimeout = 60000; + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const port1 = (server.address() as AddressInfo).port; + + const socket = connect(port1, "127.0.0.1"); + try { + await once(socket, "connect"); + let body = ""; + socket.on("data", chunk => (body += chunk)); + socket.on("error", () => {}); + socket.write("GET / HTTP/1.1\r\nHost: x\r\n\r\n"); + await inHandler.promise; + + let cb1Fired = false; + server.close(() => (cb1Fired = true)); + releaseResponse(); + while (!body.includes("ok")) await once(socket, "data"); + for (let i = 0; i < 4; i++) await new Promise(r => setImmediate(r)); + + // Re-listen while the old connection is still open. + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + + let closeEmitted = 0; + server.on("close", () => closeEmitted++); + + // Old connection ends. No 'close' must fire on the listening server. + socket.destroy(); + for (let i = 0; i < 4; i++) await new Promise(r => setImmediate(r)); + expect(closeEmitted).toBe(0); + expect(cb1Fired).toBe(false); + expect(server.listening).toBe(true); + + // Closing the new server then emits 'close' exactly once. The first + // cycle's callback was registered via once('close') and fires now too, + // like Node (and passing a second callback does not throw). + const closed = Promise.withResolvers(); + server.close(() => closed.resolve()); + await closed.promise; + expect(closeEmitted).toBe(1); + expect(cb1Fired).toBe(true); + } finally { + socket.destroy(); + server.closeAllConnections(); + } +});