From 10b496c6a40dfd468ebcb2771a02bde0c3100d8e Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sun, 13 Sep 2026 20:42:07 +0000 Subject: [PATCH 1/4] node:http2: destroy streams before the native sweep in session teardown A session teardown ran the native stream sweep while the JS streams were still live. The sweep drops DATA frames that are queued behind flow control and settles their write callbacks with no error, so the Writable emitted 'drain'. A producer parked on 'drain' woke up, and every later write() reported success because the native handle was gone. - ClientHttp2Session.destroy() destroys its open streams first, as the server session and node do. It goes through emitStreamErrorNT, so each stream gets the same error and rstCode as before. - ServerHttp2Session's socket-close handler closes and destroys the streams in JS, like the client session and node's socketOnClose. The native abort sweep it replaced has no caller left and is removed. - The server's error-GOAWAY handler no longer sweeps before destroy(), which does the same work after it has destroyed the streams. - A stream that session.destroy() destroys ends its writable without _final. The server session put DATA(END_STREAM) on the wire behind the GOAWAY and emitted 'finish' for an open response. --- src/js/node/http2.ts | 42 ++- src/runtime/api/bun/h2_frame_parser.rs | 39 --- src/runtime/api/h2.classes.ts | 4 - ...http2-session-destroy-backpressure.test.ts | 281 ++++++++++++++++++ 4 files changed, 315 insertions(+), 51 deletions(-) create mode 100644 test/js/node/http2/node-http2-session-destroy-backpressure.test.ts diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 9c7747d20737..6d9dccd10cda 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2044,6 +2044,9 @@ enum StreamState { // callback). Until then no 'error' listener can exist, so stream errors must not be emitted: // node never constructs the JS stream object before a complete header block arrives. Delivered = 1 << 8, // 100000000 = 256 + // session.destroy() is destroying the stream: node's _destroy ends the writable without + // _final, so no END_STREAM reaches the wire behind the GOAWAY and no 'finish' is emitted. + SessionDestroyed = 1 << 9, // 1000000000 = 512 } // native.writeStream() return-value flag (mirrors WRITE_FLUSHED_WITHOUT_CALLBACK in // h2_frame_parser.rs): the chunk was handed to the socket without queueing and the engine did @@ -2280,8 +2283,17 @@ function destroyStreamForSessionDestroy(error: Error | undefined, rstCode: numbe // listener would otherwise turn session.destroy(code) into an uncaught // exception (e.g. grpc-js forceShutdown destroying sessions with // NGHTTP2_CANCEL while unread UNIMPLEMENTED streams are still around). + stream[bunHTTP2StreamStatus] |= StreamState.SessionDestroyed; stream.destroy(error !== undefined && stream.listenerCount("error") > 0 ? error : undefined); } +// Client counterpart: the same synchronous teardown, through emitStreamErrorNT so the stream gets +// the error and rstCode the deferred streamError dispatch delivers. A stream still live when the +// native sweep settles its flow-control-queued write would emit 'drain' from the teardown. +function cancelStreamForSessionDestroy(session: ClientHttp2Session, rstCode: number, stream: Http2Stream) { + if (stream.destroyed || stream.closed) return; + stream[bunHTTP2StreamStatus] |= StreamState.SessionDestroyed; + emitStreamErrorNT(session, stream, rstCode, true, false); +} class Http2Stream extends Duplex { #id: number; [bunHTTP2Session]: ClientHttp2Session | ServerHttp2Session | null = null; @@ -2588,11 +2600,16 @@ class Http2Stream extends Duplex { this[kAborted] = true; this.emit("aborted"); } - // at this state destroyed will be true but we need to close the writable side - this._writableState.destroyed = false; - this.end(); - // we now restore the destroyed flag - this._writableState.destroyed = true; + if ((this[bunHTTP2StreamStatus] & StreamState.SessionDestroyed) !== 0) { + // destroyed stays set, so end() marks the writable ended and _final does not run. + this.end(); + } else { + // at this state destroyed will be true but we need to close the writable side + this._writableState.destroyed = false; + this.end(); + // we now restore the destroyed flag + this._writableState.destroyed = true; + } } const session = this[bunHTTP2Session]; @@ -4289,7 +4306,6 @@ class ServerHttp2Session extends Http2Session { // Windows agents the frame deterministically arrived first). self.destroy(); } else { - self.#parser?.emitErrorToAllStreams(errorCode); // Like Node, destroy with an error but send our own goaway with // NGHTTP2_NO_ERROR since this side had no error. self.destroy(sessionErrorFromCode(errorCode), constants.NGHTTP2_NO_ERROR); @@ -4317,7 +4333,9 @@ class ServerHttp2Session extends Http2Session { #onClose() { const parser = this.#parser; if (parser) { - parser.emitAbortToAllStreams(); + // Node's socketOnClose: close(NGHTTP2_CANCEL) every stream, then destroy it. A native abort + // sweep first would settle queued writes on streams that are still live ('drain'). + parser.forEachStream(streamCancel); parser.forEachStream(streamSocketClosed); parser.detach(); this.#parser = null; @@ -5870,9 +5888,17 @@ class ClientHttp2Session extends Http2Session { } // Like Node's Http2Stream._destroy: a received GOAWAY's code takes // precedence over the destroy code when streams are torn down. + const streamRstCode = this[kGoawayCode] || (code !== undefined ? code : constants.NGHTTP2_CANCEL); + // Destroy the open streams before the native sweep (see the server session). The sweep + // rejects a non-numeric code, and that throw must leave the streams for the retry. + if (typeof streamRstCode === "number") { + parser.forEachStream( + FunctionPrototypeBind.$call(cancelStreamForSessionDestroy, undefined, this, streamRstCode), + ); + } this[bunHTTP2SessionTeardownFrame] = $getInternalField($asyncContext, 0); try { - parser.emitErrorToAllStreams(this[kGoawayCode] || (code !== undefined ? code : constants.NGHTTP2_CANCEL)); + parser.emitErrorToAllStreams(streamRstCode); } finally { this[bunHTTP2SessionTeardownFrame] = kNoSessionTeardown; } diff --git a/src/runtime/api/bun/h2_frame_parser.rs b/src/runtime/api/bun/h2_frame_parser.rs index 0047486e23b6..f3999243a1d9 100644 --- a/src/runtime/api/bun/h2_frame_parser.rs +++ b/src/runtime/api/bun/h2_frame_parser.rs @@ -6519,45 +6519,6 @@ impl H2FrameParser { Ok(JSValue::UNDEFINED) } - #[bun_jsc::host_fn(method)] - pub(crate) fn emit_abort_to_all_streams( - this: &Self, - _global_object: &JSGlobalObject, - _callframe: &CallFrame, - ) -> JsResult { - // R-2: StreamResumableIterator stores a `ParentRef`; `streams` is `JsCell`-backed, - // so the loop body can keep using `this` (`&Self`) directly. - let mut it = StreamResumableIterator::init(this); - while let Some(stream_ptr) = it.next() { - // SAFETY: stream_ptr is a *mut Stream stored in self.streams (heap::alloc); valid for - // the lifetime of the entry. Separate heap allocation from `this`, so no aliasing. - let stream = unsafe { &mut *stream_ptr }; - // this is the oposite logic of emitErrorToallStreams, in this case we wanna to cancel this streams - if this.is_server.get() { - if stream.id % 2 == 0 { - continue; - } - } else if stream.id % 2 != 0 { - continue; - } - if stream.state != StreamState::CLOSED { - let old_state = stream.state; - stream.state = StreamState::CLOSED; - stream.rst_code = ErrorCode::CANCEL.0; - let identifier = stream.get_identifier(); - identifier.ensure_still_alive(); - stream.free_resources::(this); - this.dispatch_with_2_extra( - JSH2FrameParser::Gc::onAborted, - identifier, - JSValue::UNDEFINED, - JSValue::js_number(old_state as u8 as f64), - ); - } - } - Ok(JSValue::UNDEFINED) - } - #[bun_jsc::host_fn(method)] pub(crate) fn emit_error_to_all_streams( this: &Self, diff --git a/src/runtime/api/h2.classes.ts b/src/runtime/api/h2.classes.ts index 8b4749d6fb05..a3307e49a3d5 100644 --- a/src/runtime/api/h2.classes.ts +++ b/src/runtime/api/h2.classes.ts @@ -121,10 +121,6 @@ export default [ fn: "emitErrorToAllStreams", length: 1, }, - emitAbortToAllStreams: { - fn: "emitAbortToAllStreams", - length: 0, - }, getNextStream: { fn: "getNextStream", length: 0, diff --git a/test/js/node/http2/node-http2-session-destroy-backpressure.test.ts b/test/js/node/http2/node-http2-session-destroy-backpressure.test.ts new file mode 100644 index 000000000000..187eca1a5ab3 --- /dev/null +++ b/test/js/node/http2/node-http2-session-destroy-backpressure.test.ts @@ -0,0 +1,281 @@ +/** + * Session teardown while a stream's body is blocked on flow control: the teardown must not emit + * 'drain' on the stream, and the stream must not accept another write afterwards. Also, a stream + * that session.destroy() tears down must not end cleanly on the wire behind the GOAWAY. + * + * Works with both: + * bun bd test test/js/node/http2/node-http2-session-destroy-backpressure.test.ts + * node --test test/js/node/http2/node-http2-session-destroy-backpressure.test.ts + */ +import assert from "node:assert"; +import { once } from "node:events"; +import http2 from "node:http2"; +import net from "node:net"; +import { describe, test } from "node:test"; + +const INITIAL_WINDOW = 65535; +const PREFACE = Buffer.from("PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n"); +const F = { DATA: 0, HEADERS: 1, RST_STREAM: 3, SETTINGS: 4, PING: 6, GOAWAY: 7 }; +const NAMES = { [F.DATA]: "DATA", [F.HEADERS]: "HEADERS", [F.RST_STREAM]: "RST_STREAM", [F.GOAWAY]: "GOAWAY" }; + +function frame(type: number, flags: number, streamId: number, payload = Buffer.alloc(0)) { + const b = Buffer.alloc(9 + payload.length); + b.writeUIntBE(payload.length, 0, 3); + b[3] = type; + b[4] = flags; + b.writeUInt32BE(streamId >>> 0, 5); + payload.copy(b, 9); + return b; +} + +// HPACK "literal header field without indexing, new name" for each pair. +function hpack(headers: [string, string][]) { + return Buffer.concat( + headers.flatMap(([k, v]) => [Buffer.from([0x00, k.length]), Buffer.from(k), Buffer.from([v.length]), Buffer.from(v)]), + ); +} + +/** + * One raw h2c endpoint on `socket`. It answers SETTINGS and PING and never sends WINDOW_UPDATE, so + * a body larger than the initial window stays queued in the bun side under test. + * frames the DATA/HEADERS/RST_STREAM/GOAWAY frames received, as strings + * gotData settles on the first DATA frame with a payload + * windowExhausted settles once a full window of DATA arrived: the sender is now blocked + * closed settles when the socket closes + */ +function rawPeer(socket: net.Socket, { isServer }: { isServer: boolean }) { + const frames: string[] = []; + const gotData = Promise.withResolvers(); + const windowExhausted = Promise.withResolvers(); + const closed = Promise.withResolvers(); + let buf = Buffer.alloc(0); + let prefaceSeen = !isServer; + let received = 0; + socket.on("error", () => {}); + socket.on("close", () => closed.resolve()); + socket.on("data", chunk => { + buf = Buffer.concat([buf, chunk]); + if (!prefaceSeen) { + if (buf.length < PREFACE.length) return; + prefaceSeen = true; + buf = buf.subarray(PREFACE.length); + socket.write(frame(F.SETTINGS, 0, 0)); + } + while (buf.length >= 9) { + const len = buf.readUIntBE(0, 3); + if (buf.length < 9 + len) break; + const type = buf[3]; + const flags = buf[4]; + const payload = buf.subarray(9, 9 + len); + buf = buf.subarray(9 + len); + if (type === F.SETTINGS && !(flags & 1)) socket.write(frame(F.SETTINGS, 1, 0)); + else if (type === F.PING && !(flags & 1)) socket.write(frame(F.PING, 1, 0, payload)); + if (type in NAMES) { + const endStream = (type === F.DATA || type === F.HEADERS) && flags & 1 ? " END_STREAM" : ""; + frames.push(NAMES[type] + endStream); + } + if (type === F.DATA && len > 0) { + gotData.resolve(); + if ((received += len) >= INITIAL_WINDOW) windowExhausted.resolve(); + } + } + }); + return { + socket, + frames, + gotData: gotData.promise, + windowExhausted: windowExhausted.promise, + closed: closed.promise, + }; +} + +/** A raw h2c server for one connection, and the client session connected to it. */ +async function clientAgainstRawServer() { + const peer = Promise.withResolvers>(); + const server = net.createServer(socket => peer.resolve(rawPeer(socket, { isServer: true }))); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const session = http2.connect(`http://127.0.0.1:${(server.address() as net.AddressInfo).port}`); + session.on("error", () => {}); + await once(session, "remoteSettings"); + return { + session, + peer: await peer.promise, + close() { + session.destroy(); + server.close(); + }, + }; +} + +/** An http2 server, and a raw h2c client that has sent one request (stream 1) to it. */ +async function rawClientAgainstServer(onStream: (stream: http2.ServerHttp2Stream) => void) { + const session = Promise.withResolvers(); + const server = http2.createServer(); + server.on("session", s => { + s.on("error", () => {}); + session.resolve(s); + }); + server.on("stream", onStream); + 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"); + const peer = rawPeer(socket, { isServer: false }); + await once(socket, "connect"); + socket.write(PREFACE); + socket.write(frame(F.SETTINGS, 0, 0)); + const request: [string, string][] = [ + [":method", "POST"], + [":scheme", "http"], + [":path", "/"], + [":authority", "localhost"], + ]; + socket.write(frame(F.HEADERS, 0x4 /* END_HEADERS */, 1, hpack(request))); + return { + session: await session.promise, + peer, + close() { + socket.destroy(); + server.close(); + }, + }; +} + +/** The usual producer: on 'drain', write until write() reports backpressure again. */ +function produceOnDrain(stream: http2.Http2Stream) { + const seen = { drains: 0, accepted: 0 }; + stream.on("drain", () => { + seen.drains++; + for (let budget = 32; budget > 0 && stream.write(Buffer.alloc(16384)); budget--) seen.accepted++; + }); + return seen; +} + +/** + * Not events.once(): that rejects on the stream's 'error', which is not under test here. Settles a + * turn after 'close' so that a late 'drain' is still counted. + */ +function closedAndSettled(stream: http2.Http2Stream) { + const { promise, resolve } = Promise.withResolvers(); + stream.on("close", () => setImmediate(resolve)); + return promise; +} + +function recordEvents(stream: http2.Http2Stream) { + const events: string[] = []; + for (const name of ["aborted", "finish", "end", "close"]) stream.on(name, () => events.push(name)); + return events; +} + +// One window goes out at once. The other 32 KiB stays queued with the write callback held. +const BLOCKED_BODY = Buffer.alloc(INITIAL_WINDOW + 32768, 0x41); + +describe("a flow-control-blocked write gets no 'drain' and the stream accepts no more writes", () => { + test("when the client session is destroyed", async () => { + const { session, peer, close } = await clientAgainstRawServer(); + try { + const stream = session.request({ ":path": "/upload", ":method": "POST" }); + stream.on("error", () => {}); + const backpressured = !stream.write(BLOCKED_BODY); + const seen = produceOnDrain(stream); + await peer.windowExhausted; + + const closed = closedAndSettled(stream); + session.destroy(); + await closed; + + assert.deepStrictEqual({ backpressured, ...seen }, { backpressured: true, drains: 0, accepted: 0 }); + } finally { + close(); + } + }); + + test("when the server session's socket closes", async () => { + const result = Promise.withResolvers(); + const { peer, close } = await rawClientAgainstServer(stream => { + stream.on("error", () => {}); + stream.respond({ ":status": 200 }); + const backpressured = !stream.write(BLOCKED_BODY); + const seen = produceOnDrain(stream); + closedAndSettled(stream).then(() => result.resolve({ backpressured, ...seen })); + }); + try { + await peer.windowExhausted; + peer.socket.destroy(); + assert.deepStrictEqual(await result.promise, { backpressured: true, drains: 0, accepted: 0 }); + } finally { + close(); + } + }); + + test("when the server session receives a GOAWAY with an error code", async () => { + const result = Promise.withResolvers(); + const { peer, close } = await rawClientAgainstServer(stream => { + stream.on("error", () => {}); + stream.respond({ ":status": 200 }); + const backpressured = !stream.write(BLOCKED_BODY); + const seen = produceOnDrain(stream); + closedAndSettled(stream).then(() => result.resolve({ backpressured, ...seen })); + }); + try { + await peer.windowExhausted; + const goaway = Buffer.alloc(8); + goaway.writeUInt32BE(0, 0); // last stream id + goaway.writeUInt32BE(http2.constants.NGHTTP2_INTERNAL_ERROR, 4); + peer.socket.write(frame(F.GOAWAY, 0, 0, goaway)); + assert.deepStrictEqual(await result.promise, { backpressured: true, drains: 0, accepted: 0 }); + } finally { + close(); + } + }); +}); + +// The stream has no 'error' listener and nothing queued: the case where destroy() can still +// reach _final. An END_STREAM here tells the peer that a truncated body is complete. +describe("session.destroy() does not end an open stream on the wire", () => { + test("client session", async () => { + const { session, peer, close } = await clientAgainstRawServer(); + try { + const stream = session.request({ ":path": "/upload", ":method": "POST" }); + const events = recordEvents(stream); + stream.write(Buffer.alloc(1000, 0x41)); + await peer.gotData; + + const before = peer.frames.length; + const closed = closedAndSettled(stream); + session.destroy(); + await Promise.all([closed, peer.closed]); + + assert.deepStrictEqual( + { events, wire: peer.frames.slice(before) }, + { events: ["aborted", "close"], wire: ["GOAWAY"] }, + ); + } finally { + close(); + } + }); + + test("server session", async () => { + const result = Promise.withResolvers(); + const { session, peer, close } = await rawClientAgainstServer(stream => { + const events = recordEvents(stream); + stream.respond({ ":status": 200 }); + stream.write(Buffer.alloc(1000, 0x41)); + closedAndSettled(stream).then(() => result.resolve(events)); + }); + try { + await peer.gotData; + + const before = peer.frames.length; + session.destroy(); + const [events] = await Promise.all([result.promise, peer.closed]); + + assert.deepStrictEqual( + { events, wire: peer.frames.slice(before) }, + { events: ["aborted", "close"], wire: ["GOAWAY"] }, + ); + } finally { + close(); + } + }); +}); From 9b73710fffcedbd5b8ba6246404b30298d92a3e5 Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Mon, 14 Sep 2026 13:45:58 +0000 Subject: [PATCH 2/4] [autofix.ci] apply automated fixes --- .../http2/node-http2-session-destroy-backpressure.test.ts | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/test/js/node/http2/node-http2-session-destroy-backpressure.test.ts b/test/js/node/http2/node-http2-session-destroy-backpressure.test.ts index 187eca1a5ab3..31d26802778a 100644 --- a/test/js/node/http2/node-http2-session-destroy-backpressure.test.ts +++ b/test/js/node/http2/node-http2-session-destroy-backpressure.test.ts @@ -31,7 +31,12 @@ function frame(type: number, flags: number, streamId: number, payload = Buffer.a // HPACK "literal header field without indexing, new name" for each pair. function hpack(headers: [string, string][]) { return Buffer.concat( - headers.flatMap(([k, v]) => [Buffer.from([0x00, k.length]), Buffer.from(k), Buffer.from([v.length]), Buffer.from(v)]), + headers.flatMap(([k, v]) => [ + Buffer.from([0x00, k.length]), + Buffer.from(k), + Buffer.from([v.length]), + Buffer.from(v), + ]), ); } From 816d72f0242cc3b5e2292f7e0c9b9ec6144cef27 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Mon, 14 Sep 2026 13:49:19 +0000 Subject: [PATCH 3/4] node:http2: shorten the new comments to one line each --- src/js/node/http2.ts | 13 ++++--------- 1 file changed, 4 insertions(+), 9 deletions(-) diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 6d9dccd10cda..b3758425cf77 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2044,8 +2044,7 @@ enum StreamState { // callback). Until then no 'error' listener can exist, so stream errors must not be emitted: // node never constructs the JS stream object before a complete header block arrives. Delivered = 1 << 8, // 100000000 = 256 - // session.destroy() is destroying the stream: node's _destroy ends the writable without - // _final, so no END_STREAM reaches the wire behind the GOAWAY and no 'finish' is emitted. + // Destroyed by session.destroy(): like node's _destroy, the writable ends without _final. SessionDestroyed = 1 << 9, // 1000000000 = 512 } // native.writeStream() return-value flag (mirrors WRITE_FLUSHED_WITHOUT_CALLBACK in @@ -2286,9 +2285,7 @@ function destroyStreamForSessionDestroy(error: Error | undefined, rstCode: numbe stream[bunHTTP2StreamStatus] |= StreamState.SessionDestroyed; stream.destroy(error !== undefined && stream.listenerCount("error") > 0 ? error : undefined); } -// Client counterpart: the same synchronous teardown, through emitStreamErrorNT so the stream gets -// the error and rstCode the deferred streamError dispatch delivers. A stream still live when the -// native sweep settles its flow-control-queued write would emit 'drain' from the teardown. +// Client counterpart, through emitStreamErrorNT for the deferred path's error and rstCode. function cancelStreamForSessionDestroy(session: ClientHttp2Session, rstCode: number, stream: Http2Stream) { if (stream.destroyed || stream.closed) return; stream[bunHTTP2StreamStatus] |= StreamState.SessionDestroyed; @@ -4333,8 +4330,7 @@ class ServerHttp2Session extends Http2Session { #onClose() { const parser = this.#parser; if (parser) { - // Node's socketOnClose: close(NGHTTP2_CANCEL) every stream, then destroy it. A native abort - // sweep first would settle queued writes on streams that are still live ('drain'). + // Node's socketOnClose: close(NGHTTP2_CANCEL) every stream, then destroy it. parser.forEachStream(streamCancel); parser.forEachStream(streamSocketClosed); parser.detach(); @@ -5889,8 +5885,7 @@ class ClientHttp2Session extends Http2Session { // Like Node's Http2Stream._destroy: a received GOAWAY's code takes // precedence over the destroy code when streams are torn down. const streamRstCode = this[kGoawayCode] || (code !== undefined ? code : constants.NGHTTP2_CANCEL); - // Destroy the open streams before the native sweep (see the server session). The sweep - // rejects a non-numeric code, and that throw must leave the streams for the retry. + // The native sweep throws on a non-numeric code: the retry must still find the streams. if (typeof streamRstCode === "number") { parser.forEachStream( FunctionPrototypeBind.$call(cancelStreamForSessionDestroy, undefined, this, streamRstCode), From a2c536386e1c03dee7f4f26e46b0fa57eb52d113 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Mon, 14 Sep 2026 14:10:37 +0000 Subject: [PATCH 4/4] ci: retrigger