diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 9c7747d20737..b3758425cf77 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2044,6 +2044,8 @@ 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 + // 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 // h2_frame_parser.rs): the chunk was handed to the socket without queueing and the engine did @@ -2280,8 +2282,15 @@ 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, 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; + emitStreamErrorNT(session, stream, rstCode, true, false); +} class Http2Stream extends Duplex { #id: number; [bunHTTP2Session]: ClientHttp2Session | ServerHttp2Session | null = null; @@ -2588,11 +2597,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 +4303,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 +4330,8 @@ class ServerHttp2Session extends Http2Session { #onClose() { const parser = this.#parser; if (parser) { - parser.emitAbortToAllStreams(); + // Node's socketOnClose: close(NGHTTP2_CANCEL) every stream, then destroy it. + parser.forEachStream(streamCancel); parser.forEachStream(streamSocketClosed); parser.detach(); this.#parser = null; @@ -5870,9 +5884,16 @@ 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); + // 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), + ); + } 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..31d26802778a --- /dev/null +++ b/test/js/node/http2/node-http2-session-destroy-backpressure.test.ts @@ -0,0 +1,286 @@ +/** + * 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(); + } + }); +});