diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 282248b5675b..37f6cbb0adc9 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -86,6 +86,8 @@ const OutgoingMessagePrototype = OutgoingMessage.prototype; const { kIncomingMessage } = require("node:_http_common"); const kConnectionsCheckingInterval = Symbol("http.server.connectionsCheckingInterval"); const kTrackedConnections = Symbol("http.server.trackedConnections"); +const kClosing = Symbol("http.server.closing"); +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 @@ -119,7 +121,19 @@ const DateNow = Date.now; let cluster; +// Like Node.js's net.Server#_emitCloseIfDrained: the native all-closed promise +// resolves on pending_requests == 0 (not connections == 0), so gate the actual +// 'close' on kTrackedConnections draining; #onClose reschedules once it does. function emitCloseServer(self: Server) { + if (!self[kClosing]) return; + const connections = self[kTrackedConnections]; + if (connections && connections.size > 0) { + self[kPendingDrainClose] = true; + return; + } + self[kPendingDrainClose] = false; + self[kClosing] = false; + self[serverSymbol] = undefined; callCloseCallback(self); self.emit("close"); } @@ -276,7 +290,7 @@ function emitRequestCloseNT(self) { } function emitListeningNextTick(self, hostname, port) { - if ((self.listening = !!self[serverSymbol])) { + if ((self.listening = !!self[serverSymbol] && !self[kClosing])) { // TODO: remove the arguments // Note does not pass any arguments. self.emit("listening", null, hostname, port); @@ -309,6 +323,8 @@ function Server(options, callback): void { defineHttpAllowHalfOpen(this); this[kInternalSocketData] = undefined; this[kTrackedConnections] = new Set(); + this[kClosing] = false; + this[kPendingDrainClose] = false; this[tlsSymbol] = null; this.noDelay = true; if (typeof options === "function") { @@ -471,16 +487,13 @@ Server.prototype.unref = function () { return this; }; +// Node destroys every tracked connection and leaves the listen socket alone. Server.prototype.closeAllConnections = function () { - const server = this[serverSymbol]; - if (!server) { - return; + const connections = this[kTrackedConnections]; + if (!connections) return; + for (const socket of connections) { + socket.destroy(); } - this[serverSymbol] = undefined; - clearInterval(this[kConnectionsCheckingInterval]); - this.listening = false; - - server.stop(true); }; Server.prototype.getConnections = function (callback) { @@ -494,8 +507,21 @@ Server.prototype.getConnections = function (callback) { }; Server.prototype.closeIdleConnections = function () { - const server = this[serverSymbol]; - server?.closeIdleConnections(); + // Native sweep is authoritative on a live server (its isIdle flag spares + // mid-parse connections); the kTrackedConnections pass covers the + // post-close() window once the native app has deinit'd. + this[serverSymbol]?.closeIdleConnections(); + if (!this[kClosing]) return; + const connections = this[kTrackedConnections]; + if (!connections) return; + for (const socket of connections) { + if (socket.destroyed) continue; + const message = socket._httpMessage; + if (message && !message.finished) continue; + if (socket[kPipelinedResponses]?.length) continue; + if (socket[kHandle]?.response) continue; + socket.destroy(); + } }; Server.prototype.close = function (optionalCallback?) { @@ -503,12 +529,14 @@ Server.prototype.close = function (optionalCallback?) { // Node.js's httpServerPreClose clears the connections-checking interval // even when the server was never listening. clearInterval(this[kConnectionsCheckingInterval]); - if (!server) { + if (!server || this[kClosing]) { if (typeof optionalCallback === "function") process.nextTick(optionalCallback, $ERR_SERVER_NOT_RUNNING()); // Like Node.js's net.Server#close, close() returns the server. return this; } - this[serverSymbol] = undefined; + // kClosing is Node's "_handle == null" stand-in; serverSymbol stays set + // until emitCloseServer so the drain helpers keep working after close(). + this[kClosing] = true; if (typeof optionalCallback === "function") setCloseCallback(this, optionalCallback); this.listening = false; server.closeIdleConnections(); @@ -553,7 +581,7 @@ Server.prototype[Symbol.asyncDispose] = function () { }; Server.prototype.address = function () { - if (!this[serverSymbol]) return null; + if (!this[serverSymbol] || this[kClosing]) return null; return this[serverSymbol].address; }; @@ -1148,6 +1176,11 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // }, }); + // Cancel any prior close()'s pending drain now that the new handle is + // installed (after the fallible Bun.serve() so a throw leaves state intact). + this[kClosing] = false; + this[kPendingDrainClose] = false; + this[kCloseCallback] = undefined; getBunServerAllClosedPromise(this[serverSymbol]).$then(emitCloseNTServer.bind(this)); isHTTPS = this[serverSymbol].protocol === "https"; applyServerCustomOptions(this); @@ -1667,7 +1700,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/bun/test/parallel/test-http-server.listening-should-work.ts b/test/js/bun/test/parallel/test-http-server.listening-should-work.ts index 8f9e565558b8..872e9d1cd87e 100644 --- a/test/js/bun/test/parallel/test-http-server.listening-should-work.ts +++ b/test/js/bun/test/parallel/test-http-server.listening-should-work.ts @@ -6,5 +6,9 @@ const { expect } = createTest(import.meta.path); const server = http.createServer(); await once(server.listen(0), "listening"); expect(server.listening).toBe(true); +// closeAllConnections() destroys the connections, it does not stop listening. server.closeAllConnections(); +expect(server.listening).toBe(true); +server.close(); expect(server.listening).toBe(false); +await once(server, "close"); diff --git a/test/js/bun/test/parallel/test-http-timeout-destruction-should-be-visible-using-kConnectionsCheckingInterval.ts b/test/js/bun/test/parallel/test-http-timeout-destruction-should-be-visible-using-kConnectionsCheckingInterval.ts index 70ac11ae6895..4c73ba57a6ea 100644 --- a/test/js/bun/test/parallel/test-http-timeout-destruction-should-be-visible-using-kConnectionsCheckingInterval.ts +++ b/test/js/bun/test/parallel/test-http-timeout-destruction-should-be-visible-using-kConnectionsCheckingInterval.ts @@ -6,5 +6,10 @@ const { expect } = createTest(import.meta.path); const { kConnectionsCheckingInterval } = require("_http_server"); const server = http.createServer(); await once(server.listen(0), "listening"); +expect(server[kConnectionsCheckingInterval]._destroyed).toBe(false); +// Only close() tears the interval down; closeAllConnections() keeps listening. server.closeAllConnections(); +expect(server[kConnectionsCheckingInterval]._destroyed).toBe(false); +server.close(); expect(server[kConnectionsCheckingInterval]._destroyed).toBe(true); +await once(server, "close"); diff --git a/test/js/first_party/ws/ws.test.ts b/test/js/first_party/ws/ws.test.ts index 060bc7a1c917..5aa40fcd3363 100644 --- a/test/js/first_party/ws/ws.test.ts +++ b/test/js/first_party/ws/ws.test.ts @@ -771,6 +771,7 @@ it("Server should be able to send empty pings", async () => { return await promise; } finally { httpServer.closeAllConnections(); + httpServer.close(); } } { diff --git a/test/js/node/http/node-http-with-ws.test.ts b/test/js/node/http/node-http-with-ws.test.ts index a3ef8cac6a29..b64787a6b9a4 100644 --- a/test/js/node/http/node-http-with-ws.test.ts +++ b/test/js/node/http/node-http-with-ws.test.ts @@ -94,6 +94,7 @@ test.concurrent("should not crash when closing sockets after upgrade", async () http_socket?.destroy(); }); server.closeAllConnections(); + server.close(); resolve(); }, 10); } diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index a204e6d37259..9945cd62e327 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -1543,6 +1543,7 @@ describe("HTTP Server Security Tests - Advanced", () => { // Close the server if it's still running if (server.listening) { server.closeAllConnections(); + server.close(); } }); @@ -3467,6 +3468,293 @@ it("server.close(cb) completes after a raw upgrade once both sockets are destroy await closed; }); +// Node's server.close(cb) waits for every accepted connection to end +// (net.Server#_emitCloseIfDrained), not for the in-flight request count to +// reach zero. These four scenarios cover graceful-shutdown shapes that +// otherwise look identical to the "pending requests == 0" condition. +describe("server.close() drains connections, not requests", () => { + async function startServer(handler: http.RequestListener) { + const server = createServer(handler); + server.keepAliveTimeout = 60_000; + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const { port } = server.address() as AddressInfo; + const sock = connect(port, "127.0.0.1"); + sock.on("error", () => {}); + await once(sock, "connect"); + return { server, sock, port }; + } + + async function drainTicks() { + for (let i = 0; i < 6; i++) await new Promise(r => setImmediate(r)); + } + + it("A: close() FINs an idle keep-alive connection and fires", async () => { + const { server, sock } = await startServer((req, res) => res.end("ok:" + req.url)); + try { + let body = ""; + sock.on("data", c => (body += c)); + sock.write("GET /a HTTP/1.1\r\nHost: x\r\n\r\n"); + while (!body.includes("ok:/a")) await once(sock, "data"); + await drainTicks(); + + const ended = once(sock, "end"); + const closed = Promise.withResolvers(); + server.close(() => closed.resolve()); + await closed.promise; + await ended; + } finally { + sock.destroy(); + server.closeAllConnections(); + server.close(); + } + }); + + it("C: close(cb) waits while a keep-alive connection is still open", async () => { + const inHandler = Promise.withResolvers(); + let finishFirst!: () => void; + const paths: string[] = []; + const { server, sock } = await startServer((req, res) => { + paths.push(req.url as string); + if (paths.length === 1) { + inHandler.resolve(); + finishFirst = () => res.end("ok:" + req.url); + } else { + res.end("ok:" + req.url); + } + }); + try { + let body = ""; + sock.on("data", c => (body += c)); + sock.write("GET /first HTTP/1.1\r\nHost: x\r\n\r\n"); + await inHandler.promise; + + let closeCbFired = false; + let closeEventFired = false; + server.once("close", () => (closeEventFired = true)); + const closed = Promise.withResolvers(); + server.close(() => { + closeCbFired = true; + closed.resolve(); + }); + + finishFirst(); + while (!body.includes("ok:/first")) await once(sock, "data"); + await drainTicks(); + expect({ closeCbFired, closeEventFired }).toEqual({ closeCbFired: false, closeEventFired: false }); + + sock.write("GET /second HTTP/1.1\r\nHost: x\r\n\r\n"); + while (!body.includes("ok:/second")) await once(sock, "data"); + await drainTicks(); + expect({ closeCbFired, closeEventFired }).toEqual({ closeCbFired: false, closeEventFired: false }); + expect(paths).toEqual(["/first", "/second"]); + + sock.destroy(); + await closed.promise; + expect({ closeCbFired, closeEventFired }).toEqual({ closeCbFired: true, closeEventFired: true }); + } finally { + sock.destroy(); + server.closeAllConnections(); + server.close(); + } + }); + + it("B: closeIdleConnections() after close() drains a now-idle connection", async () => { + const inHandler = Promise.withResolvers(); + let finishFirst!: () => void; + const { server, sock } = await startServer((req, res) => { + inHandler.resolve(); + finishFirst = () => res.end("ok:" + req.url); + }); + try { + let body = ""; + sock.on("data", c => (body += c)); + const ended = once(sock, "end"); + sock.write("GET /b HTTP/1.1\r\nHost: x\r\n\r\n"); + await inHandler.promise; + + let closeCbFired = false; + const closed = Promise.withResolvers(); + server.close(() => { + closeCbFired = true; + closed.resolve(); + }); + + finishFirst(); + while (!body.includes("ok:/b")) await once(sock, "data"); + await drainTicks(); + expect(closeCbFired).toBe(false); + + // The response has been delivered and the connection is idle; a + // post-close closeIdleConnections() must reach it and let the + // callback fire. + server.closeIdleConnections(); + await closed.promise; + await ended; + } finally { + sock.destroy(); + server.closeAllConnections(); + server.close(); + } + }); + + it("B': closeAllConnections() after close() drains the in-flight connection", async () => { + const inHandler = Promise.withResolvers(); + const { server, sock } = await startServer((req, res) => { + inHandler.resolve(); + void res; + }); + try { + sock.write("GET /b2 HTTP/1.1\r\nHost: x\r\n\r\n"); + await inHandler.promise; + + const closed = Promise.withResolvers(); + server.close(() => closed.resolve()); + await drainTicks(); + + // The request handler never responded; closeAllConnections() must + // destroy the connection regardless and let the callback fire. + server.closeAllConnections(); + await closed.promise; + expect(sock.destroyed || sock.readableEnded).toBe(true); + } finally { + sock.destroy(); + server.closeAllConnections(); + server.close(); + } + }); + + it("D: close(cb) never fires while a keep-alive client keeps the connection busy", async () => { + // SIGTERM shape: close() arrives while a request is in flight, the client + // then keeps issuing keep-alive requests on that same connection. The + // close callback must not fire (and no request is served "after cb") + // until the client releases the connection. + const inHandler = Promise.withResolvers(); + let finishFirst!: () => void; + const paths: string[] = []; + const { server, sock } = await startServer((req, res) => { + paths.push(req.url as string); + if (paths.length === 1) { + inHandler.resolve(); + finishFirst = () => res.end("ok:" + req.url); + } else { + res.end("ok:" + req.url); + } + }); + try { + let body = ""; + sock.on("data", c => (body += c)); + let servedAfterCb = 0; + let closeCbFired = false; + + sock.write("GET /d0 HTTP/1.1\r\nHost: x\r\n\r\n"); + await inHandler.promise; + + const closed = Promise.withResolvers(); + server.close(() => { + closeCbFired = true; + closed.resolve(); + }); + finishFirst(); + while (!body.includes("ok:/d0")) await once(sock, "data"); + await drainTicks(); + if (closeCbFired) servedAfterCb = -1; + + for (let i = 1; i <= 5; i++) { + const marker = "ok:/d" + i; + sock.write(`GET /d${i} HTTP/1.1\r\nHost: x\r\n\r\n`); + while (!body.includes(marker)) await once(sock, "data"); + if (closeCbFired) servedAfterCb++; + await drainTicks(); + } + expect({ closeCbFired, servedAfterCb }).toEqual({ closeCbFired: false, servedAfterCb: 0 }); + + sock.destroy(); + await closed.promise; + expect(paths).toEqual(["/d0", "/d1", "/d2", "/d3", "/d4", "/d5"]); + } finally { + sock.destroy(); + server.closeAllConnections(); + server.close(); + } + }); + + it("closeAllConnections() leaves the listener running", async () => { + const { server, sock, port } = await startServer((req, res) => res.end("ok:" + req.url)); + try { + let body = ""; + sock.on("data", c => (body += c)); + sock.write("GET /l HTTP/1.1\r\nHost: x\r\n\r\n"); + while (!body.includes("ok:/l")) await once(sock, "data"); + + let closeEventFired = false; + server.once("close", () => (closeEventFired = true)); + const clientClosed = once(sock, "close"); + server.closeAllConnections(); + await clientClosed; + await drainTicks(); + expect(server.listening).toBe(true); + expect(closeEventFired).toBe(false); + + // A fresh connection is still accepted. + const sock2 = connect(port, "127.0.0.1"); + sock2.on("error", () => {}); + await once(sock2, "connect"); + let body2 = ""; + sock2.on("data", c => (body2 += c)); + sock2.write("GET /l2 HTTP/1.1\r\nHost: x\r\n\r\n"); + while (!body2.includes("ok:/l2")) await once(sock2, "data"); + sock2.destroy(); + + const closed = Promise.withResolvers(); + server.close(() => closed.resolve()); + await closed.promise; + } finally { + sock.destroy(); + server.closeAllConnections(); + server.close(); + } + }); + + it("closeIdleConnections() on a live server spares a connection mid-header-parse", async () => { + // Node's ConnectionsList.idle() excludes parsers with last_message_start != 0 + // (test-http-server-close-idle.js client1). The JS kTrackedConnections sweep + // would treat such a connection as idle, so it must not run on a live server. + const server = createServer((req, res) => res.end("ok:" + req.url)); + server.keepAliveTimeout = 60_000; + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const { port } = server.address() as AddressInfo; + const connected = once(server, "connection"); + const sock = connect(port, "127.0.0.1"); + sock.on("error", () => {}); + await once(sock, "connect"); + try { + let ended = false; + sock.on("end", () => (ended = true)); + sock.write("GET /p HTTP/1.1"); + await connected; + await drainTicks(); + + server.closeIdleConnections(); + await drainTicks(); + expect({ socketDestroyed: sock.destroyed, socketEnded: ended }).toEqual({ + socketDestroyed: false, + socketEnded: false, + }); + + let body = ""; + sock.on("data", c => (body += c)); + sock.write("\r\nHost: x\r\n\r\n"); + while (!body.includes("ok:/p")) await once(sock, "data"); + } finally { + sock.destroy(); + server.closeAllConnections(); + server.close(); + } + }); +}); + it("req.upgrade is true inside the 'connect' listener", async () => { let upgradeValue: unknown = "unset"; const { promise: sawConnect, resolve: onConnect } = Promise.withResolvers(); diff --git a/test/js/web/fetch/client-fetch.test.ts b/test/js/web/fetch/client-fetch.test.ts index 37cf159bbc09..2b90d8a2c14c 100644 --- a/test/js/web/fetch/client-fetch.test.ts +++ b/test/js/web/fetch/client-fetch.test.ts @@ -85,6 +85,7 @@ test("pre aborted with readable request body", async () => { ).rejects.toThrow(); } finally { server.closeAllConnections(); + server.close(); } }); @@ -559,6 +560,7 @@ test("fetching with Request object - issue #1527", async () => { expect(await fetch(request)).resolves.pass(); } finally { server.closeAllConnections(); + server.close(); } }); diff --git a/test/js/web/fetch/fetch.stream.test.ts b/test/js/web/fetch/fetch.stream.test.ts index ca7ee6af46dd..181d610077fb 100644 --- a/test/js/web/fetch/fetch.stream.test.ts +++ b/test/js/web/fetch/fetch.stream.test.ts @@ -243,6 +243,7 @@ describe.concurrent("fetch() with streaming", () => { expect(true).toBe(true); } finally { server?.closeAllConnections(); + server?.close(); } }); }