diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 14302e78b572..4a36eea30322 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2050,6 +2050,10 @@ enum StreamState { // The native side fully closed and freed the stream (state 7 delivered): there is // nothing left to send on the wire for it. NativeClosed = 1 << 6, // 1000000 = 64 + // Set only while respond() has a HEADERS frame carrying END_STREAM in flight: the writable + // side is finished even though this.end() has not run yet. Read by _destroy, which the native + // request() reaches synchronously when that frame closes an already half-closed stream. + EndStreamOnHeaders = 1 << 7, // 10000000 = 128 } function markWritableDone(stream: Http2Stream) { const _final = stream[bunHTTP2StreamFinal]; @@ -2338,8 +2342,10 @@ class Http2Stream extends Duplex { // closed by definition — closing it is not an abort and nothing must be sent on the wire. if (!ending && !this[kPush]) { // If the writable side of the Http2Stream is still open, emit the - // 'aborted' event and set the aborted flag. - if (!this.aborted) { + // 'aborted' event and set the aborted flag. EndStreamOnHeaders means respond() is + // submitting END_STREAM on the HEADERS frame right now and will call end() once the + // native call returns, so the response completed - nothing was cut short. + if (!this.aborted && (this[bunHTTP2StreamStatus] & StreamState.EndStreamOnHeaders) === 0) { this[kAborted] = true; this.emit("aborted"); } @@ -3116,21 +3122,30 @@ class ServerHttp2Stream extends Http2Stream { } const wireHeaders = rawHeadersList !== null ? rawHeadersList : headers; - if (typeof options === "undefined") { - session[bunHTTP2Native]?.request(this.id, undefined, wireHeaders, sensitiveNames); - } else { - session[bunHTTP2Native]?.request(this.id, undefined, wireHeaders, sensitiveNames, options); - // Only track waitForTrailers when the HEADERS frame above did NOT end - // the stream. Status codes 204/205/304 and HEAD requests force - // endStream=true earlier in this method, which means the native - // request() already wrote END_STREAM on the HEADERS frame — driving - // the wantTrailers path from `_final` on such a stream would call - // `noTrailers`/`emit("wantTrailers")` on an already-half-closed - // stream and corrupt state. Use optional chaining: `options` may be - // `null` here (typeof null === "object" enters this else branch). - if (options?.waitForTrailers && !endStream) { - this[bunHTTP2WaitForTrailers] = true; + // An END_STREAM HEADERS frame can close an already half-closed stream outright, and the + // native request() then destroys it synchronously - before the this.end() below runs. Tell + // _destroy the writable side is finished, and clear it again once the frame is out: a + // request() that threw submitted nothing, so the writable side is still genuinely open. + if (endStream) this[bunHTTP2StreamStatus] |= StreamState.EndStreamOnHeaders; + try { + if (typeof options === "undefined") { + session[bunHTTP2Native]?.request(this.id, undefined, wireHeaders, sensitiveNames); + } else { + session[bunHTTP2Native]?.request(this.id, undefined, wireHeaders, sensitiveNames, options); + // Only track waitForTrailers when the HEADERS frame above did NOT end + // the stream. Status codes 204/205/304 and HEAD requests force + // endStream=true earlier in this method, which means the native + // request() already wrote END_STREAM on the HEADERS frame — driving + // the wantTrailers path from `_final` on such a stream would call + // `noTrailers`/`emit("wantTrailers")` on an already-half-closed + // stream and corrupt state. Use optional chaining: `options` may be + // `null` here (typeof null === "object" enters this else branch). + if (options?.waitForTrailers && !endStream) { + this[bunHTTP2WaitForTrailers] = true; + } } + } finally { + this[bunHTTP2StreamStatus] &= ~StreamState.EndStreamOnHeaders; } this.headersSent = true; if (onServerStreamFinishChannel.hasSubscribers) { diff --git a/test/js/node/http2/node-http2.test.js b/test/js/node/http2/node-http2.test.js index 8f658be0d0f6..86e80116f614 100644 --- a/test/js/node/http2/node-http2.test.js +++ b/test/js/node/http2/node-http2.test.js @@ -11,10 +11,12 @@ import { Duplex } from "stream"; import http2utils from "./helpers"; import { nodeEchoServer, TLS_CERT, TLS_OPTIONS } from "./http2-helpers"; const { describe, expect, it, beforeAll, afterAll, createCallCheckCtx } = createTest(import.meta.path); -// bun-debug ships with ASAN but isn't named bun-asan, so isASAN is false -// there; the 10k-request maxSessionMemory stress test takes ~90s under -// debug+ASAN vs ~2s release, so scale for either. +// bun-debug ships with ASAN but isn't named bun-asan, so isASAN is false there. const ASAN_MULTIPLIER = isDebug ? 10 : isASAN ? 3 : 1; +// The maxSessionMemory stress test costs ~2s for 10k requests in release but ~90s under +// debug+ASAN, close enough to its timeout that a loaded machine trips it. Shrink the workload +// rather than the deadline; 1k sequential requests still exercises the memory accounting. +const MAX_SESSION_MEMORY_REQUESTS = isDebug || isASAN ? 1_000 : 10_000; function invalidArgTypeHelper(input) { if (input === null) return " Received null"; @@ -1874,7 +1876,7 @@ it( const client = http2.connect(`http://localhost:${port}`); function next(i) { - if (i === 10000) { + if (i === MAX_SESSION_MEMORY_REQUESTS) { client.close(); server.close(); resolve(); @@ -3186,3 +3188,138 @@ it("http2 allowHTTP1 fallback omits the Connection header on a close-delimited r server.close(); } }); + +// Collects the server stream's lifecycle events for one request, plus `stream.aborted` as read +// from the 'close' handler. 'end' is filtered out: bun and node disagree on whether it lands +// before or after 'finish', which is unrelated to what these tests assert. +async function serverStreamLifecycle(onStream, requestHeaders = { ":path": "/" }) { + const events = []; + let aborted = null; + const { promise: streamClosed, resolve: streamDidClose } = Promise.withResolvers(); + + const server = http2.createServer(); + server.on("stream", stream => { + stream.on("aborted", () => events.push("aborted")); + stream.on("finish", () => events.push("finish")); + stream.on("close", () => { + aborted = stream.aborted; + events.push("close"); + streamDidClose(); + }); + onStream(stream); + }); + + await new Promise(resolve => server.listen(0, resolve)); + const client = http2.connect(`http://localhost:${server.address().port}`); + client.on("error", () => {}); + try { + const req = client.request(requestHeaders); + const responseHeaders = await new Promise((resolve, reject) => { + req.on("error", reject); + req.on("response", resolve); + req.end(); + }); + const body = []; + req.on("data", chunk => body.push(chunk)); + await new Promise((resolve, reject) => { + req.on("error", reject); + req.on("end", resolve); + }); + await streamClosed; + return { events, aborted, status: responseHeaders[":status"], body: Buffer.concat(body).toString() }; + } finally { + client.close(); + server.close(); + } +} + +it("http2 respondWithFile statCheck returning false does not emit a spurious 'aborted'", async () => { + // The documented statCheck veto: reply 304 from a cache validator and cancel the file send. + // Node treats that as a completed response, so no 'aborted' fires and stream.aborted stays false. + const result = await serverStreamLifecycle(stream => { + stream.respondWithFile( + import.meta.path, + { "content-type": "text/plain" }, + { + statCheck() { + stream.respond({ ":status": 304 }); + stream.end(); + return false; + }, + }, + ); + }); + + expect(result).toEqual({ events: ["finish", "close"], aborted: false, status: 304, body: "" }); +}); + +it("http2 server respond() with END_STREAM after the request body does not emit a spurious 'aborted'", async () => { + // Any respond() whose HEADERS frame carries END_STREAM (204/205/304, HEAD, endStream: true) + // fully closes a stream whose peer already half-closed. That is a completed response, so the + // writable side must be ended before the frame is submitted or the teardown looks like an abort. + const cases = [ + { name: "304", status: 304, respond: stream => stream.respond({ ":status": 304 }) }, + { name: "204", status: 204, respond: stream => stream.respond({ ":status": 204 }) }, + { name: "205", status: 205, respond: stream => stream.respond({ ":status": 205 }) }, + { name: "endStream", status: 200, respond: stream => stream.respond({ ":status": 200 }, { endStream: true }) }, + { name: "head", status: 200, respond: stream => stream.respond({ ":status": 200 }), method: "HEAD" }, + ]; + + const results = []; + for (const { name, respond, method } of cases) { + // Respond once the request body is fully received, which is when the stream is already + // half-closed by the peer and submitting END_STREAM closes it outright. + const { events, aborted, status, body } = await serverStreamLifecycle( + stream => stream.on("end", () => respond(stream)), + method ? { ":path": "/", ":method": method } : { ":path": "/" }, + ); + results.push({ name, events, aborted, status, body }); + } + + expect(results).toEqual( + cases.map(({ name, status }) => ({ name, events: ["finish", "close"], aborted: false, status, body: "" })), + ); +}); + +it("http2 a respond() that throws on invalid headers leaves the writable side open", async () => { + // 304 forces END_STREAM, and the invalid response pseudo-header makes respond() throw. Nothing + // was submitted, so the handler can still recover with a fresh respond() + end(). + const server = http2.createServer(); + let thrownCode = null; + server.on("stream", stream => { + try { + stream.respond({ ":status": 304, ":path": "/" }); + } catch (err) { + thrownCode = err.code; + stream.respond({ ":status": 500 }); + stream.end("err"); + } + }); + + await new Promise(resolve => server.listen(0, resolve)); + const client = http2.connect(`http://localhost:${server.address().port}`); + client.on("error", () => {}); + try { + const req = client.request({ ":path": "/" }); + const headers = await new Promise((resolve, reject) => { + req.on("error", reject); + req.on("response", resolve); + req.end(); + }); + const body = []; + req.on("data", chunk => body.push(chunk)); + await new Promise((resolve, reject) => { + req.on("error", reject); + req.on("end", resolve); + }); + + expect({ + thrownCode, + status: headers[":status"], + body: Buffer.concat(body).toString(), + }).toEqual({ thrownCode: "ERR_HTTP2_INVALID_PSEUDOHEADER", status: 500, body: "err" }); + } finally { + client.close(); + server.close(); + } +});