diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index f41ce68605fc..0cd0cf8e0be4 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -444,11 +444,22 @@ function emitEventNT(self: any, event: string, ...args: any[]) { // The frame is passed in: destroy() clears it off the session before emitting // 'error', so a throwing listener on either event cannot leave a retained // session pinning the store. -function emitSessionCloseNT(self: Http2Session, frame) { +function emitSessionCloseNT(self: Http2Session, error: Error | null | undefined, frame) { + if (error) { + runInFrame(frame, self.emit, self, "error", error); + } if (self.listenerCount("close") > 0) { runInFrame(frame, self.emit, self, "close"); } } +// Node's emitClose: the session reports 'error' and 'close' once its socket has closed. +function emitSessionCloseAfterSocket(self: Http2Session, socket, error: Error | null | undefined, frame) { + if (socket && !socket.destroyed) { + socket.once("close", () => emitSessionCloseNT(self, error, frame)); + } else { + process.nextTick(emitSessionCloseNT, self, error, frame); + } +} function emitErrorNT(self: any, error: any, destroy: boolean) { if (destroy) { if (self.listenerCount("error") > 0) { @@ -4750,7 +4761,8 @@ class ServerHttp2Session extends Http2Session { return; } this.#destroying = true; - emitHttp2SessionPerf(this, this.#parser, this[bunHTTP2Socket]); + const socket = this[bunHTTP2Socket]; + emitHttp2SessionPerf(this, this.#parser, socket); try { const server = this[kServer]; if (server) { @@ -4774,7 +4786,6 @@ class ServerHttp2Session extends Http2Session { this[kSessionDestroyError] = error; } - const socket = this[bunHTTP2Socket]; if (!this.#connected) return; this.#closed = true; this.#connected = false; @@ -4785,24 +4796,7 @@ class ServerHttp2Session extends Http2Session { // a destroy(err) after close() must still put the error GOAWAY on the wire. this.goaway(code || constants.NGHTTP2_NO_ERROR, 0, Buffer.alloc(0)); } - if (error) { - // node's finishSessionClose destroys the socket when the session dies - // with an error (a misbehaving peer must observe the connection going - // away) - but it still ends first and destroys a tick later, so the - // final GOAWAY flushes behind a FIN instead of an abortive close (see - // endThenDestroySessionSocket). - endThenDestroySessionSocket(socket, error); - } else { - // Node's finishSessionClose: "If we're gracefully closing the socket, - // call resume() so we can detect the peer closing in case - // binding.Http2Session is already gone." Without a reader, unread - // inbound bytes (a late GOAWAY from the peer) turn the close into an - // RST, which the peer surfaces as read ECONNRESET (routine on Windows - // loopback - the same reason Node delays the error-path destroy). - // https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188 - socket.resume(); - socket.end(); - } + closeSessionSocket(socket, this.#closeCalled && !error, error); } const parser = this.#parser; if (parser) { @@ -4828,16 +4822,9 @@ class ServerHttp2Session extends Http2Session { } this[bunHTTP2Socket] = null; - // Read-and-clear the frame first: emitting 'error' with no listener throws, - // which would skip the clear and leave a retained session pinning the store. const asyncFrame = this[bunHTTP2AsyncContextFrame]; this[bunHTTP2AsyncContextFrame] = undefined; - if (error) { - runInFrame(asyncFrame, this.emit, this, "error", error); - } - // node emits the session 'close' event asynchronously (a listener attached right after - // close()/destroy() returns must still observe it). - process.nextTick(emitSessionCloseNT, this, asyncFrame); + emitSessionCloseAfterSocket(this, socket, error, asyncFrame); } } function emitTimeout(session: ClientHttp2Session) { @@ -4880,19 +4867,21 @@ function setSessionTimeout(this: Http2Session, msecs, callback) { return this; } -// Node's finishSessionClose error path: socket.end() flushes and sends the FIN -// first, and the hard destroy runs a tick later - "If session.destroy() was -// called, destroy the underlying socket. Delay it a bit to try to avoid -// ECONNRESET on Windows" - so the peer reads our final frames off a FIN'd -// socket instead of observing an abortive close. -// https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188 function destroySessionSocketDelayedNT(socket, error) { if (!socket.destroyed) { socket.destroy(error); } } -function endThenDestroySessionSocket(socket, error) { - socket.end(() => setImmediate(destroySessionSocketDelayedNT, socket, error)); +// Node's finishSessionClose: end(), then destroy() unless close() asked for a graceful shutdown. +function closeSessionSocket(socket, graceful: boolean, error: Error | null | undefined) { + if (socket.destroyed) return; + if (graceful) { + // Unread inbound bytes would turn the close into an RST. + socket.resume(); + socket.end(); + } else { + socket.end(() => setImmediate(destroySessionSocketDelayedNT, socket, error)); + } } // node callTimeout (lib/internal/http2/core.js): when the timer expires while writes are still in // flight and bytes have reached the wire since the previous expiry, the session is not idle — @@ -5896,18 +5885,7 @@ class ClientHttp2Session extends Http2Session { // a destroy(err) after close() must still put the error GOAWAY on the wire. this.goaway(code || constants.NGHTTP2_NO_ERROR, 0, Buffer.alloc(0)); } - if (error) { - // See the client session: end first, destroy a tick later (node's - // finishSessionClose Windows-ECONNRESET avoidance). - endThenDestroySessionSocket(socket, error); - } else { - // See the client session's destroy: Node's finishSessionClose resumes - // the socket on a graceful close so unread inbound bytes cannot turn - // the FIN teardown into an RST. - // https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188 - socket.resume(); - socket.end(); - } + closeSessionSocket(socket, this.#closeCalled && !error, error); } const parser = this.#parser; if (parser) { @@ -5937,16 +5915,9 @@ class ClientHttp2Session extends Http2Session { this.#parser = null; this[bunHTTP2Socket] = null; - // Read-and-clear the frame first: emitting 'error' with no listener throws, - // which would skip the clear and leave a retained session pinning the store. const asyncFrame = this[bunHTTP2AsyncContextFrame]; this[bunHTTP2AsyncContextFrame] = undefined; - if (error) { - runInFrame(asyncFrame, this.emit, this, "error", error); - } - // node emits the session 'close' event asynchronously (a listener attached right after - // close()/destroy() returns must still observe it). - process.nextTick(emitSessionCloseNT, this, asyncFrame); + emitSessionCloseAfterSocket(this, socket, error, asyncFrame); } request(headers?: HeadersObject | any[] | null, options?: ClientRequestOptions) { diff --git a/test/js/node/async_hooks/AsyncLocalStorage.test.ts b/test/js/node/async_hooks/AsyncLocalStorage.test.ts index 004923688fbc..74eaaa86d6c9 100644 --- a/test/js/node/async_hooks/AsyncLocalStorage.test.ts +++ b/test/js/node/async_hooks/AsyncLocalStorage.test.ts @@ -1,7 +1,7 @@ import { AsyncLocalStorage, AsyncResource } from "async_hooks"; import { heapStats } from "bun:jsc"; import { describe, expect, test } from "bun:test"; -import { bunEnv, bunExe } from "harness"; +import { bunEnv, bunExe, isASAN, isDebug } from "harness"; import http2 from "http2"; describe("AsyncLocalStorage", () => { @@ -280,7 +280,10 @@ test("re-entering a storage inside run() does not grow the context", () => { }; als.run(0, () => { const before = objects(); - for (let i = 0; i < 100_000; i++) { + // A context that grows per re-entry adds at least one object per iteration, so 10k + // iterations still blow past the threshold. The full count is too slow under ASAN. + const iterations = isASAN || isDebug ? 10_000 : 100_000; + for (let i = 0; i < iterations; i++) { als.run(1, () => { using _ = als.withScope(2); }); @@ -1180,9 +1183,11 @@ describe("async context passes through", () => { expect(stderr).not.toContain("AssertionError"); }); - // destroy(err) with no 'error' listener throws out of the emit, which must - // not skip the frame clear (the 'close' tick after it never runs). - test("http2 clears the session frame when destroy(err) throws past the emit", async () => { + // Like node, destroy(err) emits 'error' later (once the socket has closed), so + // with no listener it surfaces as an uncaught exception rather than a throw out + // of destroy(); the frame must already be clear when destroy() returns, not + // only once the deferred emit has run. + test("http2 clears the session frame when destroy(err) has no 'error' listener", async () => { await using proc = Bun.spawn({ cmd: [ bunExe(), @@ -1196,13 +1201,18 @@ describe("async context passes through", () => { als.run({ marker: true }, () => { client = http2.connect("http://127.0.0.1:" + server.address().port); }); + // The throw IS the condition: the unlistened 'error' can only come from + // the deferred emit, which runs strictly after the read-and-clear. + process.on("uncaughtException", err => { + console.log("UNCAUGHT " + err.message); + // Drain rather than process.exit(), like the siblings. + server.close(); + }); client.on("connect", () => { - try { client.destroy(new Error("boom")); } catch {} + client.destroy(new Error("boom")); const sym = Object.getOwnPropertySymbols(client) .find(x => x.description === "::bunhttp2asynccontextframe::"); console.log(sym === undefined ? "SYMBOL-MISSING" : client[sym] === undefined ? "CLEARED" : "PINNED"); - // Drain rather than process.exit(), like the siblings. - server.close(); }); });`, ], @@ -1211,7 +1221,7 @@ describe("async context passes through", () => { stderr: "pipe", }); const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); - expect({ stdout: stdout.trim(), exitCode }).toEqual({ stdout: "CLEARED", exitCode: 0 }); + expect({ stdout: stdout.trim(), exitCode }).toEqual({ stdout: "CLEARED\nUNCAUGHT boom", exitCode: 0 }); expect(stderr).not.toContain("AssertionError"); }); diff --git a/test/js/node/http2/h2-conformance.test.ts b/test/js/node/http2/h2-conformance.test.ts index b121dbb06374..d34e0c3ca034 100644 --- a/test/js/node/http2/h2-conformance.test.ts +++ b/test/js/node/http2/h2-conformance.test.ts @@ -1943,10 +1943,18 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou const calm = (f: Frame) => f.type === FrameType.GOAWAY && goawayErrorCode(f) === ErrorCode.ENHANCE_YOUR_CALM; function respondingServer(options: Record = {}, rejectUploads = false) { - const state = { handlers: 0, sessionErrorCode: undefined as string | undefined }; + // A promise, not a value sampled when the GOAWAY arrives: like node, the session reports its + // error once its socket has closed, which is after the GOAWAY has reached the peer. + const sessionError = Promise.withResolvers(); + sessionError.promise.catch(() => {}); // the tests that expect no session error never read it + const state = { handlers: 0, sessionErrorCode: sessionError.promise }; const server = http2.createServer(options); - server.on("sessionError", (e: any) => (state.sessionErrorCode = e.code)); - server.on("session", s => s.on("error", () => {})); + server.on("sessionError", (e: any) => sessionError.resolve(e.code)); + server.on("session", s => { + s.on("error", () => {}); + // 'error' comes before 'close', so this only settles the promise when there was no error. + s.on("close", () => sessionError.reject(new Error("the session closed without a 'sessionError'"))); + }); server.on("stream", (stream: any, headers: any) => { state.handlers++; stream.on("error", () => {}); @@ -1992,7 +2000,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou expect(goaway.payload.subarray(8).toString()).toBe("too many stream resets"); // Like node, every request up to the bucket's edge still reaches the handler. expect(handlers).toBeGreaterThan(900); - expect(sessionErrorCode).toBe("ERR_HTTP2_ERROR"); + expect(await sessionErrorCode).toBe("ERR_HTTP2_ERROR"); }); test("streamResetBurst sets where the flood is detected", async () => { @@ -2049,7 +2057,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou expect(goawayErrorCode(goaway)).toBe(ErrorCode.ENHANCE_YOUR_CALM); expect(c.frames.filter(f => f.type === FrameType.RST_STREAM).length).toBeGreaterThanOrEqual(1000); expect(handlers).toBeGreaterThanOrEqual(1000); - expect(sessionErrorCode).toBe("ERR_HTTP2_ERROR"); + expect(await sessionErrorCode).toBe("ERR_HTTP2_ERROR"); }); } @@ -2150,7 +2158,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou c.send(pairs(1200, 1, rstStream)); const goaway = await c.waitFor(calm, 10_000); expect(goaway.payload.subarray(8).toString()).toBe("too many stream resets"); - expect(state.sessionErrorCode).toBe("ERR_HTTP2_ERROR"); + expect(await state.sessionErrorCode).toBe("ERR_HTTP2_ERROR"); }); }); @@ -2164,7 +2172,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou c.send(Buffer.concat(ids.map(madeYouReset["WINDOW_UPDATE with a 0 increment"]))); const goaway = await c.waitFor(calm, 10_000); expect(goaway.payload.subarray(8).toString()).toBe("too many stream resets"); - expect(state.sessionErrorCode).toBe("ERR_HTTP2_ERROR"); + expect(await state.sessionErrorCode).toBe("ERR_HTTP2_ERROR"); }); }); }); diff --git a/test/js/node/http2/node-http2.test.js b/test/js/node/http2/node-http2.test.js index a7d2c4d33238..596ce6323b9b 100644 --- a/test/js/node/http2/node-http2.test.js +++ b/test/js/node/http2/node-http2.test.js @@ -5828,6 +5828,262 @@ it("close() completes when the peer never ACKs an outstanding SETTINGS", async ( server.close(); }); +// Node's finishSessionClose: destroy() ends the socket and then hard-destroys it, so the +// connection is released even when the peer never closes its own side (the peer a timeout or +// an error handler is typically destroying). close() only ends it and waits for the peer. In +// both cases the session's 'error'/'close' are emitted once the socket itself has closed. +describe.concurrent("session teardown when the peer never closes its side of the connection", () => { + // allowHalfOpen keeps the peer's side open after our FIN arrives: from the server's point of + // view this is a peer that stays connected for as long as it likes, and `peer.end()` is the + // peer finally hanging up. + async function lingeringPeer(server) { + await new Promise(resolve => server.listen(0, "127.0.0.1", resolve)); + const { port } = server.address(); + const peer = net.connect({ port, host: "127.0.0.1", allowHalfOpen: true }); + peer.on("error", () => {}); + const sawFin = new Promise(resolve => peer.once("end", resolve)); + peer.resume(); // 'end' only fires on a socket that is reading + return { peer, sawFin }; + } + async function acceptH2c(server) { + const session = new Promise(resolve => server.once("session", resolve)); + // The server's own 'connection' listener (which creates the session) is registered first, + // so this resolves with the socket behind the session above. + const socket = new Promise(resolve => server.once("connection", resolve)); + const { peer, sawFin } = await lingeringPeer(server); + return { session: await session, socket: await socket, peer, sawFin }; + } + function connectionCount(server) { + const { promise, resolve, reject } = Promise.withResolvers(); + server.getConnections((err, count) => (err ? reject(err) : resolve(count))); + return promise; + } + // Records the session's terminal events, each tagged with whether its socket was already + // destroyed when the event fired; `closed` resolves with the list once 'close' has fired. + function recordTeardown(session, socket) { + const events = []; + const { promise: closed, resolve } = Promise.withResolvers(); + session.on("error", err => events.push(["error", err.message, socket.destroyed])); + session.on("close", () => { + events.push(["close", socket.destroyed]); + resolve(events); + }); + return { events, closed }; + } + // A server that accepts and then says and does nothing, and whose side of each connection + // stays open after the client's FIN (allowHalfOpen): the client-side shape of the same peer. + async function silentServer() { + const accepted = []; + const server = net.createServer({ allowHalfOpen: true }, socket => { + accepted.push(socket); + socket.resume(); // 'end' only fires on a socket that is reading + }); + await new Promise(resolve => server.listen(0, "127.0.0.1", resolve)); + const { port } = server.address(); + return { + port, + firstConnectionData: new Promise(resolve => server.once("connection", socket => socket.once("data", resolve))), + firstConnectionEnded: new Promise(resolve => server.once("connection", socket => socket.once("end", resolve))), + [Symbol.dispose]() { + for (const socket of accepted) socket.destroy(); + server.close(); + }, + }; + } + function connectOverOwnSocket(port) { + let socket; + const client = http2.connect(`http://127.0.0.1:${port}`, { + createConnection: () => (socket = net.connect({ port, host: "127.0.0.1" })), + }); + return { client, socket }; + } + + it("server session.destroy() destroys the socket and releases the server connection", async () => { + const server = http2.createServer(); + const { session, socket, peer, sawFin } = await acceptH2c(server); + try { + const { closed } = recordTeardown(session, socket); + session.destroy(); + expect(session.destroyed).toBeTrue(); + await sawFin; // the FIN still goes out first: the peer is told we are done before the hard close + expect(await closed).toEqual([["close", true]]); + expect(await connectionCount(server)).toBe(0); + await new Promise(resolve => server.close(resolve)); + } finally { + peer.destroy(); + server.close(); + } + }); + + it("server session.destroy(err) emits 'error' and 'close' only once the socket is gone", async () => { + const server = http2.createServer(); + const { session, socket, peer } = await acceptH2c(server); + try { + const { closed } = recordTeardown(session, socket); + session.destroy(new Error("boom")); + expect(await closed).toEqual([ + ["error", "boom", true], + ["close", true], + ]); + expect(await connectionCount(server)).toBe(0); + } finally { + peer.destroy(); + server.close(); + } + }); + + it("server.setTimeout() reaping an idle session releases its connection", async () => { + const server = http2.createServer(); + server.setTimeout(1); // armed on every accepted session; an unhandled 'timeout' destroys it + const { session, socket, peer } = await acceptH2c(server); + try { + const { closed } = recordTeardown(session, socket); + expect(await closed).toEqual([["close", true]]); + expect(await connectionCount(server)).toBe(0); + } finally { + peer.destroy(); + server.close(); + } + }); + + it("server session.close() leaves the socket to the peer and reports 'close' once it hangs up", async () => { + const server = http2.createServer(); + const { session, socket, peer, sawFin } = await acceptH2c(server); + try { + const { events, closed } = recordTeardown(session, socket); + session.close(); + await sawFin; // our GOAWAY + FIN are out and the peer has not answered either + expect({ events, socketDestroyed: socket.destroyed, connections: await connectionCount(server) }).toEqual({ + events: [], + socketDestroyed: false, + connections: 1, + }); + peer.end(); // the peer hangs up: now the socket closes, and the session reports it + expect(await closed).toEqual([["close", true]]); + expect(await connectionCount(server)).toBe(0); + } finally { + peer.destroy(); + server.close(); + } + }); + + it("secure server session.destroy() destroys the TLS socket and releases the server connection", async () => { + const server = http2.createSecureServer(TLS_CERT); + const session = new Promise(resolve => server.once("session", resolve)); + const socket = new Promise(resolve => server.once("secureConnection", resolve)); + const { peer, sawFin } = await lingeringPeer(server); + // The peer's TLS layer talks through a carrier Duplex that stops delivering inbound bytes + // once the handshake is done, so it never sees the server's close_notify and never answers + // it (a TLS peer that answered would close the connection itself and mask the bug). + let handshakeDone = false; + const carrier = new Duplex({ + read() {}, + write(chunk, _encoding, callback) { + peer.write(chunk, callback); + }, + }); + peer.on("data", chunk => { + if (!handshakeDone) carrier.push(chunk); + }); + await new Promise(resolve => peer.once("connect", resolve)); + const peerTls = tls.connect({ socket: carrier, ALPNProtocols: ["h2"], ...TLS_OPTIONS, servername: "localhost" }); + peerTls.on("error", () => {}); + await new Promise(resolve => peerTls.once("secureConnect", resolve)); + handshakeDone = true; + const serverSession = await session; + const serverSocket = await socket; + try { + const { closed } = recordTeardown(serverSession, serverSocket); + serverSession.destroy(); + await sawFin; + expect(await closed).toEqual([["close", true]]); + expect(await connectionCount(server)).toBe(0); + await new Promise(resolve => server.close(resolve)); + } finally { + peerTls.destroy(); + peer.destroy(); + server.close(); + } + }); + + it("client session.destroy() destroys the socket even though the server never closes", async () => { + using server = await silentServer(); + const { client, socket } = connectOverOwnSocket(server.port); + try { + await new Promise(resolve => client.once("connect", resolve)); + const { closed } = recordTeardown(client, socket); + client.destroy(); + expect(client.destroyed).toBeTrue(); + await server.firstConnectionEnded; // the FIN still goes out first + expect(await closed).toEqual([["close", true]]); + } finally { + client.destroy(); + } + }); + + it("client session.destroy(err) emits 'error' and 'close' only once the socket is gone", async () => { + using server = await silentServer(); + const { client, socket } = connectOverOwnSocket(server.port); + try { + await new Promise(resolve => client.once("connect", resolve)); + const { closed } = recordTeardown(client, socket); + client.destroy(new Error("boom")); + expect(await closed).toEqual([ + ["error", "boom", true], + ["close", true], + ]); + } finally { + client.destroy(); + } + }); + + // The session never connects: the server accepts the TCP connection and never answers the + // ClientHello (a load balancer that accepts and stalls). destroy() is the caller's connect + // deadline, and the socket it leaves behind must not outlive the session. + it("client session.destroy() during a TLS handshake the server never answers destroys the socket", async () => { + using server = await silentServer(); + let socket; + const client = http2.connect(`https://127.0.0.1:${server.port}`, { + createConnection: () => + (socket = tls.connect({ + port: server.port, + host: "127.0.0.1", + ALPNProtocols: ["h2"], + rejectUnauthorized: false, + })), + }); + let connected = false; + client.on("connect", () => (connected = true)); + try { + await server.firstConnectionData; // the ClientHello arrived; no ServerHello ever follows + const { closed } = recordTeardown(client, socket); + client.destroy(); + expect(client.destroyed).toBeTrue(); + await server.firstConnectionEnded; // the FIN still goes out first + expect(await closed).toEqual([["close", true]]); + expect(connected).toBeFalse(); + } finally { + client.destroy(); + } + }); + + it("a session whose socket the peer already closed still reports 'close'", async () => { + const server = http2.createServer(); + const { session, socket, peer } = await acceptH2c(server); + try { + const { closed } = recordTeardown(session, socket); + // A FIN, not an RST: peer.destroy() with unread inbound bytes resets the connection on + // macOS and Windows, and the session would then report that ECONNRESET as its 'error'. + peer.end(); + expect(await closed).toEqual([["close", true]]); + expect(await connectionCount(server)).toBe(0); + } finally { + peer.destroy(); + server.close(); + } + }); +}); + // A pull-mode consumer (pause() then on('readable')/read()) must reopen the receive // window via _read(): the 'resume' event never fires on that path, so without _read() // clearing the paused gate the peer stalls at the initial ~64KB stream window.