From 052cf976096e1c5a5c4472c08c192e4fa267a948 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 29 Sep 2026 03:38:33 +0000 Subject: [PATCH 1/4] node:http2: count a refused stream where the session refuses it, not in rstStream rstStream counted every REFUSED_STREAM reset that JS submits on a server against maxSessionRejectedStreams. close() and the destroy that follows it each submit one, so a server whose handler calls stream.close(NGHTTP2_REFUSED_STREAM) sent GOAWAY(ENHANCE_YOUR_CALM) at the 50th refusal and closed its own session. A RST_STREAM(REFUSED_STREAM) from the peer was counted too. The count was written for one case: the stream that streamStart refuses over SETTINGS_MAX_CONCURRENT_STREAMS. streamStart now returns the code for that stream, and the native caller resets the stream and counts it. rstStream counts nothing. --- src/js/node/http2.ts | 4 +- src/runtime/api/bun/h2_frame_parser.rs | 52 +++--- test/js/node/http2/h2-conformance.test.ts | 184 ++++++++++++++++++++++ test/js/node/http2/node-http2.test.js | 53 +++++++ 4 files changed, 260 insertions(+), 33 deletions(-) diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index f41ce68605fc..fc7e41610555 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -4051,8 +4051,8 @@ class ServerHttp2Session extends Http2Session { // SETTINGS_MAX_CONCURRENT_STREAMS. nghttp2 answers with RST_STREAM REFUSED_STREAM and never // surfaces the stream to the JS layer. if (stream_id % 2 === 1 && self.#peerInitiatedStreams >= self.#advertisedMaxConcurrentStreams) { - self.#parser?.rstStream(stream_id, constants.NGHTTP2_REFUSED_STREAM); - return; + // Native counts this against maxSessionRejectedStreams and resets the stream while budget remains. + return constants.NGHTTP2_REFUSED_STREAM; } self.#connections++; if (stream_id % 2 === 1) self.#peerInitiatedStreams++; diff --git a/src/runtime/api/bun/h2_frame_parser.rs b/src/runtime/api/bun/h2_frame_parser.rs index 46a7ea31dd49..efb73649a7c7 100644 --- a/src/runtime/api/bun/h2_frame_parser.rs +++ b/src/runtime/api/bun/h2_frame_parser.rs @@ -3452,10 +3452,30 @@ impl H2FrameParser { }); self.enter_stream_dispatch(stream) .set_context(returned, &global); + } else if returned.is_number() && self.count_rejected_stream(stream_identifier) { + // streamStart refused the stream and returned the RST_STREAM code that answers it. + // SAFETY: stream is *mut Stream from self.streams; valid while the map entry exists + self.end_stream(unsafe { &mut *stream }, ErrorCode(returned.to_u32())); } Some(stream) } + /// Returns false when this used up maxSessionRejectedStreams and the session sent its GOAWAY. + fn count_rejected_stream(&self, stream_id: u32) -> bool { + self.rejected_streams.set(self.rejected_streams.get() + 1); + if self.max_rejected_streams.get() <= self.rejected_streams.get() { + self.send_go_away( + stream_id, + ErrorCode::ENHANCE_YOUR_CALM, + b"ENHANCE_YOUR_CALM", + self.last_stream_id.get(), + true, + ); + return false; + } + true + } + fn to_writer(&self) -> DirectWriterStruct { DirectWriterStruct { writer: bun_ptr::BackRef::new(self), @@ -4165,18 +4185,7 @@ impl crate::api::h2::connection::Sink for H2FrameParser { } fn on_stream_rejected(&self, stream_id: u32) { - // maxSessionRejectedStreams: counts only locally-initiated rejections (oversized or - // malformed header blocks) - peer-sent RST_STREAM frames must not consume the budget. - self.rejected_streams.set(self.rejected_streams.get() + 1); - if self.max_rejected_streams.get() <= self.rejected_streams.get() { - self.send_go_away( - stream_id, - ErrorCode::ENHANCE_YOUR_CALM, - b"ENHANCE_YOUR_CALM", - self.last_stream_id.get(), - true, - ); - } + self.count_rejected_stream(stream_id); } fn on_stream_reset(&self, stream_id: u32, code: u32) { @@ -5044,25 +5053,6 @@ impl H2FrameParser { } let error_code = error_arg.to_u32(); - // maxSessionRejectedStreams: a REFUSED_STREAM reset from the JS layer (the - // max-concurrent-streams refusal in streamStart) is the same rejection class the engine - // counts; budget it identically so a flood of refused streams still tears the session - // down. Server-side only: a client's GOAWAY sweep resets its own unprocessed streams - // with REFUSED_STREAM and must not consume the budget. - if error_code == ErrorCode::REFUSED_STREAM.0 && this.is_server.get() { - this.rejected_streams.set(this.rejected_streams.get() + 1); - if this.max_rejected_streams.get() <= this.rejected_streams.get() { - this.send_go_away( - stream_id, - ErrorCode::ENHANCE_YOUR_CALM, - b"ENHANCE_YOUR_CALM", - this.last_stream_id.get(), - true, - ); - return Ok(JSValue::UNDEFINED); - } - } - let Some(stream) = this.streams.get().get(&stream_id).copied() else { // Streams the legacy bookkeeping never registered (e.g. peer-initiated pushed streams // surfaced by the rewrite engine) get the RST_STREAM written directly. The frame is diff --git a/test/js/node/http2/h2-conformance.test.ts b/test/js/node/http2/h2-conformance.test.ts index b121dbb06374..238c8fe9841f 100644 --- a/test/js/node/http2/h2-conformance.test.ts +++ b/test/js/node/http2/h2-conformance.test.ts @@ -1929,6 +1929,190 @@ describe("stream release after a queued END_STREAM", () => { }); }); +// The reset of a stream that reached the handler does not use maxSessionRejectedStreams, whatever +// code it carries and whoever asks for it. Node counts in one place only, where it rejects a stream +// at creation: https://github.com/nodejs/node/blob/v26.3.0/src/node_http2.cc#L1035-L1050 +describe.concurrent("maxSessionRejectedStreams and the reset of a delivered stream", () => { + const REFUSED_STREAM = http2.constants.NGHTTP2_REFUSED_STREAM; + const refused = Buffer.alloc(4); + refused.writeUInt32BE(ErrorCode.REFUSED_STREAM, 0); + const get = (c: RawH2, id: number) => c.sendFrame(FrameType.HEADERS, 0x5, id, requestHeaderBlock("GET")); + const upload = (c: RawH2, id: number) => c.sendFrame(FrameType.HEADERS, 0x4, id, requestHeaderBlock("POST")); + + /** + * Opens `count` requests on one connection, one after the other, and sends a PING after each. + * A stream is reset from setImmediate callbacks that close() and _destroy() queue, so the PING + * goes out one turn after the server's stream emitted 'close'. + * Resolves with the number of PINGs the server answered and the code of every GOAWAY it wrote. + */ + async function pingAfterEachRequest( + server: http2.Http2Server, + count: number, + open: (c: RawH2, id: number, opened: Promise) => void | Promise = get, + ) { + const streams = new Map; closed: PromiseWithResolvers }>(); + const events = (id: number) => { + if (!streams.has(id)) streams.set(id, { opened: Promise.withResolvers(), closed: Promise.withResolvers() }); + return streams.get(id)!; + }; + server.on("stream", (stream: any) => { + const { opened, closed } = events(stream.id); + stream.once("close", closed.resolve); + opened.resolve(); + }); + server.on("session", (session: any) => session.on("error", () => {})); + server.listen(0); + await once(server, "listening"); + const c = await RawH2.connect((server.address() as net.AddressInfo).port); + const ended = once(c.socket, "close"); + try { + c.sendPreface(); + c.sendEmptySettings(); + let acked = 0; + for (let i = 0; i < count; i++) { + const { opened, closed } = events(1 + 2 * i); + await Promise.race([open(c, 1 + 2 * i, opened.promise), ended]); + await Promise.race([closed.promise, ended]); + if (c.closed) break; + await new Promise(resolve => setImmediate(resolve)); + const tag = Buffer.alloc(8); + tag.writeUInt32BE(i + 1, 0); + c.sendFrame(FrameType.PING, 0, 0, tag); + const answer = await c.waitFor( + f => + f.type === FrameType.GOAWAY || + (f.type === FrameType.PING && (f.flags & 0x1) !== 0 && f.payload.equals(tag)), + ); + if (answer.type !== FrameType.PING) break; + acked++; + } + return { acked, goaways: c.frames.filter(f => f.type === FrameType.GOAWAY).map(goawayErrorCode) }; + } finally { + c.destroy(); + server.close(); + } + } + + const refusals: Record void> = { + "close(REFUSED_STREAM)": stream => stream.close(REFUSED_STREAM), + "respond() then close(REFUSED_STREAM)": stream => { + stream.respond({ ":status": 200 }); + stream.close(REFUSED_STREAM); + }, + "close(REFUSED_STREAM) from a later turn": stream => setImmediate(() => stream.close(REFUSED_STREAM)), + "close(REFUSED_STREAM) twice then destroy()": stream => { + stream.close(REFUSED_STREAM); + stream.close(REFUSED_STREAM); + stream.destroy(); + }, + }; + + describe.each(Object.keys(refusals))("%s in every handler", name => { + test.each([ + ["a request that has ended", get], + ["a request whose body is still open", upload], + ])("of %s keeps the session", async (_, open) => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }); + server.on("stream", (stream: any) => { + stream.on("error", () => {}); + refusals[name](stream); + }); + expect(await pingAfterEachRequest(server, 6, open)).toEqual({ acked: 6, goaways: [] }); + }); + }); + + test("close(REFUSED_STREAM) of every pushed stream keeps the session", async () => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }); + server.on("stream", (stream: any) => { + stream.on("error", () => {}); + stream.pushStream({ ":path": "/pushed" }, (err: Error | null, pushed: any) => { + if (err) return; + pushed.on("error", () => {}); + pushed.close(REFUSED_STREAM); + }); + stream.respond({ ":status": 200 }); + stream.end("ok"); + }); + expect(await pingAfterEachRequest(server, 6)).toEqual({ acked: 6, goaways: [] }); + }); + + test("req.stream.close(REFUSED_STREAM) in every request listener keeps the session", async () => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }, (req: any, res: any) => { + req.on("error", () => {}); + res.on("error", () => {}); + req.stream.on("error", () => {}); + req.stream.close(REFUSED_STREAM); + }); + expect(await pingAfterEachRequest(server, 6)).toEqual({ acked: 6, goaways: [] }); + }); + + test("RST_STREAM(REFUSED_STREAM) from the peer on every stream keeps the session", async () => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }); + server.on("stream", (stream: any) => stream.on("error", () => {})); + const result = await pingAfterEachRequest(server, 6, async (c, id, opened) => { + upload(c, id); + await opened; + c.sendFrame(FrameType.RST_STREAM, 0, id, refused); + }); + expect(result).toEqual({ acked: 6, goaways: [] }); + }); + + test.each([0, 1, 2])("maxSessionRejectedStreams: %d keeps the session after one refused request", async max => { + const server = http2.createServer({ maxSessionRejectedStreams: max }); + server.on("stream", (stream: any) => { + stream.on("error", () => {}); + stream.close(REFUSED_STREAM); + }); + expect(await pingAfterEachRequest(server, 1)).toEqual({ acked: 1, goaways: [] }); + }); + + // Bun-only, kept from 1.4.x: node v26.3.0 answers 3, 5 and 7 with RST_STREAM(REFUSED_STREAM), + // charges maxSessionInvalidFrames and keeps the session. Rewrite this test when the refusal + // moves to that budget. + test("a stream refused over maxConcurrentStreams is counted", async () => { + const seen: number[] = []; + let sessionErrorCode: string | undefined; + const server = http2.createServer({ settings: { maxConcurrentStreams: 1 }, maxSessionRejectedStreams: 3 }); + server.on("sessionError", (e: any) => (sessionErrorCode = e.code)); + server.on("session", (session: any) => session.on("error", () => {})); + server.on("stream", (stream: any) => { + seen.push(stream.id); + stream.on("error", () => {}); + stream.respond({ ":status": 200 }); + stream.write("hello"); + }); + server.listen(0); + await once(server, "listening"); + const c = await RawH2.connect((server.address() as net.AddressInfo).port); + try { + c.sendPreface(); + c.sendEmptySettings(); + get(c, 1); + await c.waitFor(f => f.type === FrameType.HEADERS && f.streamId === 1); + for (const id of [3, 5, 7]) get(c, id); + const goaway = await c.waitForGoaway(); + await c.waitClosed(); + expect({ + seen, + resets: c.frames.filter(f => f.type === FrameType.RST_STREAM).map(f => [f.streamId, f.payload.readUInt32BE(0)]), + goaway: goawayErrorCode(goaway), + sessionErrorCode, + }).toEqual({ + seen: [1], + resets: [ + [3, ErrorCode.REFUSED_STREAM], + [5, ErrorCode.REFUSED_STREAM], + ], + goaway: ErrorCode.ENHANCE_YOUR_CALM, + sessionErrorCode: "ERR_HTTP2_SESSION_ERROR", + }); + } finally { + c.destroy(); + server.close(); + } + }); +}); + // Stream resets are rate-limited per connection like nghttp2's stream_reset_ratelim (burst 1000, // refill 33/s): past the bucket the session dies with GOAWAY, surfaced as ERR_HTTP2_ERROR. describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYouReset)", () => { diff --git a/test/js/node/http2/node-http2.test.js b/test/js/node/http2/node-http2.test.js index a7d2c4d33238..f7c9a91dabb1 100644 --- a/test/js/node/http2/node-http2.test.js +++ b/test/js/node/http2/node-http2.test.js @@ -3,6 +3,7 @@ import { bunEnv, bunExe, isASAN, isCI, isDebug, nodeExe } from "harness"; import { createTest } from "node-harness"; import { AsyncLocalStorage } from "node:async_hooks"; import dc from "node:diagnostics_channel"; +import { once } from "node:events"; import fs from "node:fs"; import http2 from "node:http2"; import https from "node:https"; @@ -2195,6 +2196,58 @@ it("http2 client receives 'goaway' when the server rejects a stream", async () = } }); +it("http2 server keeps its session when the handler refuses requests with REFUSED_STREAM", async () => { + // A budget of 2 so that one refusal uses it up when resets of delivered streams count: close() + // and the destroy that follows it each submit a reset. + const server = http2.createServer({ maxSessionRejectedStreams: 2 }); + server.on("stream", (stream, headers) => { + stream.on("error", () => {}); + if (headers[":path"] === "/events") { + stream.respond({ ":status": 200 }); + stream.on("data", chunk => stream.write(chunk)); + return; + } + stream.close(http2.constants.NGHTTP2_REFUSED_STREAM); + }); + await new Promise(resolve => server.listen(0, "127.0.0.1", resolve)); + + const client = http2.connect(`http://127.0.0.1:${server.address().port}`); + try { + const seen = { goaway: [], sessionError: undefined, eventsError: undefined }; + client.on("goaway", code => seen.goaway.push(code)); + client.on("error", err => (seen.sessionError = err.message)); + + const events = client.request({ ":path": "/events", ":method": "POST" }); + events.on("error", err => (seen.eventsError = err.code)); + events.setEncoding("utf8"); + await once(events, "response"); + + for (let i = 0; i < 5; i++) { + const { promise: closed, resolve: onClose } = Promise.withResolvers(); + const req = client.request({ ":path": "/work" }); + // A refused request ends with an 'error' on node. + req.on("error", () => {}); + req.on("close", onClose); + req.end(); + await closed; + // The long-lived stream still carries data in both directions. + events.write(`${i}`); + expect((await once(events, "data"))[0]).toBe(`${i}`); + } + expect({ ...seen, closed: client.closed, destroyed: client.destroyed }).toEqual({ + goaway: [], + sessionError: undefined, + eventsError: undefined, + closed: false, + destroyed: false, + }); + events.close(); + } finally { + client.close(); + server.close(); + } +}); + // node's Http2Session#goawayCode / #goawayLastStreamID report the GOAWAY frame this side // received (0 / 0 until one arrives). Expected values below were checked against node v26.3.0. describe.concurrent("http2 session goawayCode / goawayLastStreamID", () => { From b45a84aaf7d31f510c1ee53ad08f7baeee10f284 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 29 Sep 2026 18:30:34 +0000 Subject: [PATCH 2/4] ci: retrigger From 5f739d7625fa9d3248269f9c7d1e5ddb84d0704d Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 29 Sep 2026 19:11:42 +0000 Subject: [PATCH 3/4] node:http2: put the rejected-stream tests in a file of their own The tests have their own raw client with no timers. A wait ends with a frame or with the end of the connection. --- test/js/node/http2/h2-conformance.test.ts | 184 ---------- .../http2/node-http2-rejected-streams.test.ts | 338 ++++++++++++++++++ test/js/node/http2/node-http2.test.js | 53 --- 3 files changed, 338 insertions(+), 237 deletions(-) create mode 100644 test/js/node/http2/node-http2-rejected-streams.test.ts diff --git a/test/js/node/http2/h2-conformance.test.ts b/test/js/node/http2/h2-conformance.test.ts index 238c8fe9841f..b121dbb06374 100644 --- a/test/js/node/http2/h2-conformance.test.ts +++ b/test/js/node/http2/h2-conformance.test.ts @@ -1929,190 +1929,6 @@ describe("stream release after a queued END_STREAM", () => { }); }); -// The reset of a stream that reached the handler does not use maxSessionRejectedStreams, whatever -// code it carries and whoever asks for it. Node counts in one place only, where it rejects a stream -// at creation: https://github.com/nodejs/node/blob/v26.3.0/src/node_http2.cc#L1035-L1050 -describe.concurrent("maxSessionRejectedStreams and the reset of a delivered stream", () => { - const REFUSED_STREAM = http2.constants.NGHTTP2_REFUSED_STREAM; - const refused = Buffer.alloc(4); - refused.writeUInt32BE(ErrorCode.REFUSED_STREAM, 0); - const get = (c: RawH2, id: number) => c.sendFrame(FrameType.HEADERS, 0x5, id, requestHeaderBlock("GET")); - const upload = (c: RawH2, id: number) => c.sendFrame(FrameType.HEADERS, 0x4, id, requestHeaderBlock("POST")); - - /** - * Opens `count` requests on one connection, one after the other, and sends a PING after each. - * A stream is reset from setImmediate callbacks that close() and _destroy() queue, so the PING - * goes out one turn after the server's stream emitted 'close'. - * Resolves with the number of PINGs the server answered and the code of every GOAWAY it wrote. - */ - async function pingAfterEachRequest( - server: http2.Http2Server, - count: number, - open: (c: RawH2, id: number, opened: Promise) => void | Promise = get, - ) { - const streams = new Map; closed: PromiseWithResolvers }>(); - const events = (id: number) => { - if (!streams.has(id)) streams.set(id, { opened: Promise.withResolvers(), closed: Promise.withResolvers() }); - return streams.get(id)!; - }; - server.on("stream", (stream: any) => { - const { opened, closed } = events(stream.id); - stream.once("close", closed.resolve); - opened.resolve(); - }); - server.on("session", (session: any) => session.on("error", () => {})); - server.listen(0); - await once(server, "listening"); - const c = await RawH2.connect((server.address() as net.AddressInfo).port); - const ended = once(c.socket, "close"); - try { - c.sendPreface(); - c.sendEmptySettings(); - let acked = 0; - for (let i = 0; i < count; i++) { - const { opened, closed } = events(1 + 2 * i); - await Promise.race([open(c, 1 + 2 * i, opened.promise), ended]); - await Promise.race([closed.promise, ended]); - if (c.closed) break; - await new Promise(resolve => setImmediate(resolve)); - const tag = Buffer.alloc(8); - tag.writeUInt32BE(i + 1, 0); - c.sendFrame(FrameType.PING, 0, 0, tag); - const answer = await c.waitFor( - f => - f.type === FrameType.GOAWAY || - (f.type === FrameType.PING && (f.flags & 0x1) !== 0 && f.payload.equals(tag)), - ); - if (answer.type !== FrameType.PING) break; - acked++; - } - return { acked, goaways: c.frames.filter(f => f.type === FrameType.GOAWAY).map(goawayErrorCode) }; - } finally { - c.destroy(); - server.close(); - } - } - - const refusals: Record void> = { - "close(REFUSED_STREAM)": stream => stream.close(REFUSED_STREAM), - "respond() then close(REFUSED_STREAM)": stream => { - stream.respond({ ":status": 200 }); - stream.close(REFUSED_STREAM); - }, - "close(REFUSED_STREAM) from a later turn": stream => setImmediate(() => stream.close(REFUSED_STREAM)), - "close(REFUSED_STREAM) twice then destroy()": stream => { - stream.close(REFUSED_STREAM); - stream.close(REFUSED_STREAM); - stream.destroy(); - }, - }; - - describe.each(Object.keys(refusals))("%s in every handler", name => { - test.each([ - ["a request that has ended", get], - ["a request whose body is still open", upload], - ])("of %s keeps the session", async (_, open) => { - const server = http2.createServer({ maxSessionRejectedStreams: 2 }); - server.on("stream", (stream: any) => { - stream.on("error", () => {}); - refusals[name](stream); - }); - expect(await pingAfterEachRequest(server, 6, open)).toEqual({ acked: 6, goaways: [] }); - }); - }); - - test("close(REFUSED_STREAM) of every pushed stream keeps the session", async () => { - const server = http2.createServer({ maxSessionRejectedStreams: 2 }); - server.on("stream", (stream: any) => { - stream.on("error", () => {}); - stream.pushStream({ ":path": "/pushed" }, (err: Error | null, pushed: any) => { - if (err) return; - pushed.on("error", () => {}); - pushed.close(REFUSED_STREAM); - }); - stream.respond({ ":status": 200 }); - stream.end("ok"); - }); - expect(await pingAfterEachRequest(server, 6)).toEqual({ acked: 6, goaways: [] }); - }); - - test("req.stream.close(REFUSED_STREAM) in every request listener keeps the session", async () => { - const server = http2.createServer({ maxSessionRejectedStreams: 2 }, (req: any, res: any) => { - req.on("error", () => {}); - res.on("error", () => {}); - req.stream.on("error", () => {}); - req.stream.close(REFUSED_STREAM); - }); - expect(await pingAfterEachRequest(server, 6)).toEqual({ acked: 6, goaways: [] }); - }); - - test("RST_STREAM(REFUSED_STREAM) from the peer on every stream keeps the session", async () => { - const server = http2.createServer({ maxSessionRejectedStreams: 2 }); - server.on("stream", (stream: any) => stream.on("error", () => {})); - const result = await pingAfterEachRequest(server, 6, async (c, id, opened) => { - upload(c, id); - await opened; - c.sendFrame(FrameType.RST_STREAM, 0, id, refused); - }); - expect(result).toEqual({ acked: 6, goaways: [] }); - }); - - test.each([0, 1, 2])("maxSessionRejectedStreams: %d keeps the session after one refused request", async max => { - const server = http2.createServer({ maxSessionRejectedStreams: max }); - server.on("stream", (stream: any) => { - stream.on("error", () => {}); - stream.close(REFUSED_STREAM); - }); - expect(await pingAfterEachRequest(server, 1)).toEqual({ acked: 1, goaways: [] }); - }); - - // Bun-only, kept from 1.4.x: node v26.3.0 answers 3, 5 and 7 with RST_STREAM(REFUSED_STREAM), - // charges maxSessionInvalidFrames and keeps the session. Rewrite this test when the refusal - // moves to that budget. - test("a stream refused over maxConcurrentStreams is counted", async () => { - const seen: number[] = []; - let sessionErrorCode: string | undefined; - const server = http2.createServer({ settings: { maxConcurrentStreams: 1 }, maxSessionRejectedStreams: 3 }); - server.on("sessionError", (e: any) => (sessionErrorCode = e.code)); - server.on("session", (session: any) => session.on("error", () => {})); - server.on("stream", (stream: any) => { - seen.push(stream.id); - stream.on("error", () => {}); - stream.respond({ ":status": 200 }); - stream.write("hello"); - }); - server.listen(0); - await once(server, "listening"); - const c = await RawH2.connect((server.address() as net.AddressInfo).port); - try { - c.sendPreface(); - c.sendEmptySettings(); - get(c, 1); - await c.waitFor(f => f.type === FrameType.HEADERS && f.streamId === 1); - for (const id of [3, 5, 7]) get(c, id); - const goaway = await c.waitForGoaway(); - await c.waitClosed(); - expect({ - seen, - resets: c.frames.filter(f => f.type === FrameType.RST_STREAM).map(f => [f.streamId, f.payload.readUInt32BE(0)]), - goaway: goawayErrorCode(goaway), - sessionErrorCode, - }).toEqual({ - seen: [1], - resets: [ - [3, ErrorCode.REFUSED_STREAM], - [5, ErrorCode.REFUSED_STREAM], - ], - goaway: ErrorCode.ENHANCE_YOUR_CALM, - sessionErrorCode: "ERR_HTTP2_SESSION_ERROR", - }); - } finally { - c.destroy(); - server.close(); - } - }); -}); - // Stream resets are rate-limited per connection like nghttp2's stream_reset_ratelim (burst 1000, // refill 33/s): past the bucket the session dies with GOAWAY, surfaced as ERR_HTTP2_ERROR. describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYouReset)", () => { diff --git a/test/js/node/http2/node-http2-rejected-streams.test.ts b/test/js/node/http2/node-http2-rejected-streams.test.ts new file mode 100644 index 000000000000..2e013578ae2a --- /dev/null +++ b/test/js/node/http2/node-http2-rejected-streams.test.ts @@ -0,0 +1,338 @@ +import { describe, expect, test } from "bun:test"; +import { once } from "node:events"; +import http2 from "node:http2"; +import net from "node:net"; + +// The reset of a stream that reached the handler does not use maxSessionRejectedStreams, whatever +// code it carries and whoever asks for it. Node counts in one place only, where it rejects a stream +// at creation: https://github.com/nodejs/node/blob/v26.3.0/src/node_http2.cc#L1035-L1050 + +const PREFACE = Buffer.from("PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n", "latin1"); +const HEADERS = 0x1; +const RST_STREAM = 0x3; +const SETTINGS = 0x4; +const PING = 0x6; +const GOAWAY = 0x7; +const END_STREAM = 0x1; +const END_HEADERS = 0x4; +const ACK = 0x1; +const REFUSED_STREAM = http2.constants.NGHTTP2_REFUSED_STREAM; +const ENHANCE_YOUR_CALM = http2.constants.NGHTTP2_ENHANCE_YOUR_CALM; + +type Frame = { type: number; flags: number; streamId: number; payload: Buffer }; + +function frame(type: number, flags: number, streamId: number, payload: Buffer = Buffer.alloc(0)): Buffer { + const header = Buffer.alloc(9); + header.writeUIntBE(payload.length, 0, 3); + header.writeUInt8(type, 3); + header.writeUInt8(flags, 4); + header.writeUInt32BE(streamId, 5); + return Buffer.concat([header, payload]); +} + +function u32(value: number): Buffer { + const buffer = Buffer.alloc(4); + buffer.writeUInt32BE(value, 0); + return buffer; +} + +// :method GET or POST, :scheme http, :path /, and a literal :authority. No dynamic table entry. +function requestBlock(method: "GET" | "POST"): Buffer { + const authority = Buffer.from("localhost", "latin1"); + return Buffer.concat([Buffer.from([method === "POST" ? 0x83 : 0x82, 0x86, 0x84, 0x01, authority.length]), authority]); +} + +/** A raw HTTP/2 client. It has no timers: a wait ends with a frame or with the end of the connection. */ +class RawClient { + frames: Frame[] = []; + closed = false; + #buffered = Buffer.alloc(0); + #waiters: Array<{ matches: (f: Frame) => boolean; resolve: (f: Frame | null) => void }> = []; + + constructor(readonly socket: net.Socket) { + socket.on("error", () => {}); + socket.on("close", () => { + this.closed = true; + for (const waiter of this.#waiters.splice(0)) waiter.resolve(null); + }); + socket.on("data", chunk => { + this.#buffered = Buffer.concat([this.#buffered, chunk]); + while (this.#buffered.length >= 9) { + const length = this.#buffered.readUIntBE(0, 3); + if (this.#buffered.length < 9 + length) break; + const received: Frame = { + type: this.#buffered.readUInt8(3), + flags: this.#buffered.readUInt8(4), + streamId: this.#buffered.readUInt32BE(5) & 0x7fffffff, + payload: Buffer.from(this.#buffered.subarray(9, 9 + length)), + }; + this.#buffered = this.#buffered.subarray(9 + length); + this.frames.push(received); + const index = this.#waiters.findIndex(waiter => waiter.matches(received)); + if (index !== -1) this.#waiters.splice(index, 1)[0].resolve(received); + } + }); + } + + static async connect(server: http2.Http2Server): Promise { + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const socket = net.connect((server.address() as net.AddressInfo).port, "127.0.0.1"); + await once(socket, "connect"); + const client = new RawClient(socket); + socket.write(Buffer.concat([PREFACE, frame(SETTINGS, 0, 0)])); + return client; + } + + /** Resolves with the first frame that matches, or with null when the connection ended first. */ + waitFor(matches: (f: Frame) => boolean): Promise { + const existing = this.frames.find(matches); + if (existing) return Promise.resolve(existing); + if (this.closed) return Promise.resolve(null); + return new Promise(resolve => this.#waiters.push({ matches, resolve })); + } + + get(streamId: number) { + this.socket.write(frame(HEADERS, END_STREAM | END_HEADERS, streamId, requestBlock("GET"))); + } + + /** A request whose body stays open. */ + upload(streamId: number) { + this.socket.write(frame(HEADERS, END_HEADERS, streamId, requestBlock("POST"))); + } + + /** True when the server answered the PING. False when it sent GOAWAY or closed the connection. */ + async ping(tag: number): Promise { + const payload = Buffer.alloc(8); + payload.writeUInt32BE(tag, 0); + this.socket.write(frame(PING, 0, 0, payload)); + const answer = await this.waitFor( + f => f.type === GOAWAY || (f.type === PING && (f.flags & ACK) !== 0 && f.payload.equals(payload)), + ); + return answer?.type === PING; + } + + goawayCodes(): number[] { + return this.frames.filter(f => f.type === GOAWAY).map(f => f.payload.readUInt32BE(4)); + } + + resets(): Array<[number, number]> { + return this.frames.filter(f => f.type === RST_STREAM).map(f => [f.streamId, f.payload.readUInt32BE(0)]); + } +} + +/** + * Opens `count` requests on one connection, one after the other, and sends a PING after each. + * A stream is reset from setImmediate callbacks that close() and _destroy() queue, so the PING + * goes out one turn after the server's stream emitted 'close'. + * Resolves with the number of PINGs the server answered and the code of every GOAWAY it wrote. + */ +async function pingAfterEachRequest( + server: http2.Http2Server, + count: number, + open: (client: RawClient, streamId: number, opened: Promise) => void | Promise = (client, streamId) => + client.get(streamId), +) { + const streams = new Map; closed: PromiseWithResolvers }>(); + const events = (streamId: number) => { + if (!streams.has(streamId)) { + streams.set(streamId, { opened: Promise.withResolvers(), closed: Promise.withResolvers() }); + } + return streams.get(streamId)!; + }; + server.on("stream", (stream: http2.ServerHttp2Stream) => { + const { opened, closed } = events(stream.id!); + stream.once("close", () => closed.resolve()); + opened.resolve(); + }); + server.on("session", session => session.on("error", () => {})); + const client = await RawClient.connect(server); + const ended = once(client.socket, "close"); + try { + let acked = 0; + for (let i = 0; i < count; i++) { + const streamId = 1 + 2 * i; + const { opened, closed } = events(streamId); + await Promise.race([open(client, streamId, opened.promise), ended]); + await Promise.race([closed.promise, ended]); + if (client.closed) break; + await new Promise(resolve => setImmediate(resolve)); + if (!(await client.ping(i + 1))) break; + acked++; + } + return { acked, goaways: client.goawayCodes() }; + } finally { + client.socket.destroy(); + server.close(); + } +} + +describe("maxSessionRejectedStreams and the reset of a delivered stream", () => { + const refusals: Record void> = { + "close(REFUSED_STREAM)": stream => stream.close(REFUSED_STREAM), + "respond() then close(REFUSED_STREAM)": stream => { + stream.respond({ ":status": 200 }); + stream.close(REFUSED_STREAM); + }, + "close(REFUSED_STREAM) from a later turn": stream => { + setImmediate(() => stream.close(REFUSED_STREAM)); + }, + "close(REFUSED_STREAM) twice then destroy()": stream => { + stream.close(REFUSED_STREAM); + stream.close(REFUSED_STREAM); + stream.destroy(); + }, + }; + + describe.each(Object.keys(refusals))("%s in every handler", name => { + test.each([ + ["a request that has ended", (client: RawClient, streamId: number) => client.get(streamId)], + ["a request whose body is still open", (client: RawClient, streamId: number) => client.upload(streamId)], + ])("of %s keeps the session", async (_, open) => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }); + server.on("stream", stream => { + stream.on("error", () => {}); + refusals[name](stream); + }); + expect(await pingAfterEachRequest(server, 6, open)).toEqual({ acked: 6, goaways: [] }); + }); + }); + + test("close(REFUSED_STREAM) of every pushed stream keeps the session", async () => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }); + server.on("stream", stream => { + stream.on("error", () => {}); + stream.pushStream({ ":path": "/pushed" }, (err, pushed) => { + if (err) return; + pushed.on("error", () => {}); + pushed.close(REFUSED_STREAM); + }); + stream.respond({ ":status": 200 }); + stream.end("ok"); + }); + expect(await pingAfterEachRequest(server, 6)).toEqual({ acked: 6, goaways: [] }); + }); + + test("req.stream.close(REFUSED_STREAM) in every request listener keeps the session", async () => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }, (req, res) => { + req.on("error", () => {}); + res.on("error", () => {}); + req.stream.on("error", () => {}); + req.stream.close(REFUSED_STREAM); + }); + expect(await pingAfterEachRequest(server, 6)).toEqual({ acked: 6, goaways: [] }); + }); + + test("RST_STREAM(REFUSED_STREAM) from the peer on every stream keeps the session", async () => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }); + server.on("stream", stream => { + stream.on("error", () => {}); + }); + const result = await pingAfterEachRequest(server, 6, async (client, streamId, opened) => { + client.upload(streamId); + await opened; + client.socket.write(frame(RST_STREAM, 0, streamId, u32(REFUSED_STREAM))); + }); + expect(result).toEqual({ acked: 6, goaways: [] }); + }); + + test.each([0, 1, 2])("maxSessionRejectedStreams: %d keeps the session after one refused request", async max => { + const server = http2.createServer({ maxSessionRejectedStreams: max }); + server.on("stream", stream => { + stream.on("error", () => {}); + stream.close(REFUSED_STREAM); + }); + expect(await pingAfterEachRequest(server, 1)).toEqual({ acked: 1, goaways: [] }); + }); + + test("a long-lived stream survives refused requests on its connection", async () => { + const server = http2.createServer({ maxSessionRejectedStreams: 2 }); + server.on("stream", (stream, headers) => { + stream.on("error", () => {}); + if (headers[":path"] === "/events") { + stream.respond({ ":status": 200 }); + stream.on("data", chunk => stream.write(chunk)); + return; + } + stream.close(REFUSED_STREAM); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + + const client = http2.connect(`http://127.0.0.1:${(server.address() as net.AddressInfo).port}`); + try { + const seen: { goaway: number[]; sessionError?: string; eventsError?: string } = { goaway: [] }; + client.on("goaway", code => seen.goaway.push(code)); + client.on("error", err => (seen.sessionError = err.message)); + + const events = client.request({ ":path": "/events", ":method": "POST" }); + events.on("error", (err: NodeJS.ErrnoException) => (seen.eventsError = err.code)); + events.setEncoding("utf8"); + await once(events, "response"); + + for (let i = 0; i < 5; i++) { + const closed = Promise.withResolvers(); + const req = client.request({ ":path": "/work" }); + // A refused request ends with an 'error' on node. + req.on("error", () => {}); + req.on("close", () => closed.resolve()); + req.end(); + await closed.promise; + // The long-lived stream still carries data in both directions. + events.write(`${i}`); + expect((await once(events, "data"))[0]).toBe(`${i}`); + } + expect({ ...seen, closed: client.closed, destroyed: client.destroyed }).toEqual({ + goaway: [], + closed: false, + destroyed: false, + }); + events.close(); + } finally { + client.close(); + server.close(); + } + }); + + // Bun-only, kept from 1.4.x: node v26.3.0 answers 3, 5 and 7 with RST_STREAM(REFUSED_STREAM), + // charges maxSessionInvalidFrames and keeps the session. Rewrite this test when the refusal + // moves to that budget. + test("a stream refused over maxConcurrentStreams is counted", async () => { + const seen: number[] = []; + let sessionErrorCode: string | undefined; + const server = http2.createServer({ settings: { maxConcurrentStreams: 1 }, maxSessionRejectedStreams: 3 }); + server.on("sessionError", (err: NodeJS.ErrnoException) => (sessionErrorCode = err.code)); + server.on("session", session => session.on("error", () => {})); + server.on("stream", stream => { + seen.push(stream.id!); + stream.on("error", () => {}); + stream.respond({ ":status": 200 }); + stream.write("hello"); + }); + const client = await RawClient.connect(server); + try { + client.get(1); + await client.waitFor(f => f.type === HEADERS && f.streamId === 1); + for (const streamId of [3, 5, 7]) client.get(streamId); + const goaway = await client.waitFor(f => f.type === GOAWAY); + if (!client.closed) await once(client.socket, "close"); + expect({ + seen, + resets: client.resets(), + goaway: goaway?.payload.readUInt32BE(4), + sessionErrorCode, + }).toEqual({ + seen: [1], + resets: [ + [3, REFUSED_STREAM], + [5, REFUSED_STREAM], + ], + goaway: ENHANCE_YOUR_CALM, + sessionErrorCode: "ERR_HTTP2_SESSION_ERROR", + }); + } finally { + client.socket.destroy(); + server.close(); + } + }); +}); diff --git a/test/js/node/http2/node-http2.test.js b/test/js/node/http2/node-http2.test.js index f7c9a91dabb1..a7d2c4d33238 100644 --- a/test/js/node/http2/node-http2.test.js +++ b/test/js/node/http2/node-http2.test.js @@ -3,7 +3,6 @@ import { bunEnv, bunExe, isASAN, isCI, isDebug, nodeExe } from "harness"; import { createTest } from "node-harness"; import { AsyncLocalStorage } from "node:async_hooks"; import dc from "node:diagnostics_channel"; -import { once } from "node:events"; import fs from "node:fs"; import http2 from "node:http2"; import https from "node:https"; @@ -2196,58 +2195,6 @@ it("http2 client receives 'goaway' when the server rejects a stream", async () = } }); -it("http2 server keeps its session when the handler refuses requests with REFUSED_STREAM", async () => { - // A budget of 2 so that one refusal uses it up when resets of delivered streams count: close() - // and the destroy that follows it each submit a reset. - const server = http2.createServer({ maxSessionRejectedStreams: 2 }); - server.on("stream", (stream, headers) => { - stream.on("error", () => {}); - if (headers[":path"] === "/events") { - stream.respond({ ":status": 200 }); - stream.on("data", chunk => stream.write(chunk)); - return; - } - stream.close(http2.constants.NGHTTP2_REFUSED_STREAM); - }); - await new Promise(resolve => server.listen(0, "127.0.0.1", resolve)); - - const client = http2.connect(`http://127.0.0.1:${server.address().port}`); - try { - const seen = { goaway: [], sessionError: undefined, eventsError: undefined }; - client.on("goaway", code => seen.goaway.push(code)); - client.on("error", err => (seen.sessionError = err.message)); - - const events = client.request({ ":path": "/events", ":method": "POST" }); - events.on("error", err => (seen.eventsError = err.code)); - events.setEncoding("utf8"); - await once(events, "response"); - - for (let i = 0; i < 5; i++) { - const { promise: closed, resolve: onClose } = Promise.withResolvers(); - const req = client.request({ ":path": "/work" }); - // A refused request ends with an 'error' on node. - req.on("error", () => {}); - req.on("close", onClose); - req.end(); - await closed; - // The long-lived stream still carries data in both directions. - events.write(`${i}`); - expect((await once(events, "data"))[0]).toBe(`${i}`); - } - expect({ ...seen, closed: client.closed, destroyed: client.destroyed }).toEqual({ - goaway: [], - sessionError: undefined, - eventsError: undefined, - closed: false, - destroyed: false, - }); - events.close(); - } finally { - client.close(); - server.close(); - } -}); - // node's Http2Session#goawayCode / #goawayLastStreamID report the GOAWAY frame this side // received (0 / 0 until one arrives). Expected values below were checked against node v26.3.0. describe.concurrent("http2 session goawayCode / goawayLastStreamID", () => { From 9ac99b967d529f134a5abf36021b35789f1a9140 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 29 Sep 2026 20:39:39 +0000 Subject: [PATCH 4/4] node:http2: reset a refused stream through enter_stream_dispatch The helper arms the dispatch guard for the borrow, as the branch above it does. No behaviour change. --- src/runtime/api/bun/h2_frame_parser.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/runtime/api/bun/h2_frame_parser.rs b/src/runtime/api/bun/h2_frame_parser.rs index efb73649a7c7..3f8a95cb3371 100644 --- a/src/runtime/api/bun/h2_frame_parser.rs +++ b/src/runtime/api/bun/h2_frame_parser.rs @@ -3454,8 +3454,8 @@ impl H2FrameParser { .set_context(returned, &global); } else if returned.is_number() && self.count_rejected_stream(stream_identifier) { // streamStart refused the stream and returned the RST_STREAM code that answers it. - // SAFETY: stream is *mut Stream from self.streams; valid while the map entry exists - self.end_stream(unsafe { &mut *stream }, ErrorCode(returned.to_u32())); + let mut refused = self.enter_stream_dispatch(stream); + self.end_stream(&mut refused, ErrorCode(returned.to_u32())); } Some(stream) }