diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 1e97096a70bc..3db6b5f3584f 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2060,6 +2060,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 + // Native dispatched streamError or aborted: it closed the stream and wrote the RST_STREAM if one was due. + NativeReset = 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 @@ -2226,8 +2228,9 @@ function markStreamClosed(stream: Http2Stream) { markWritableDone(stream); } } -function rstNextTick(id: number, rstCode: number) { - const session = this as Http2Session; +function rstNextTick(this: Http2Stream, session: Http2Session, id: number, rstCode: number) { + // Native drops this call only via the stream's table entry: evicted on the next read, absent for a pushed stream. + if ((this[bunHTTP2StreamStatus] & (StreamState.NativeClosed | StreamState.NativeReset)) !== 0) return; session[bunHTTP2Native]?.rstStream(id, rstCode); } // node streamOnPause/streamOnResume (lib/internal/http2/core.js): the readable's flow state @@ -2247,7 +2250,7 @@ function streamOnResume(this: Http2Stream) { // A close() on a stream that has not been submitted yet (no id): the RST_STREAM has to follow the // HEADERS frame, which is sent when the queued request becomes ready (node's finishCloseStream). function sendRstOnReady(this: Http2Stream, session: Http2Session, code: number) { - setImmediate(rstNextTick.bind(session, this.id, code)); + setImmediate(rstNextTick.bind(this, session, this.id, code)); } function uncorkNT(stream: Http2Stream) { stream.uncork(); @@ -2591,9 +2594,9 @@ class Http2Stream extends (Duplex as Http2StreamBase) { // RST_STREAM has to be sent after the HEADERS frame, once the id is assigned. this.once("ready", sendRstOnReady.bind(this, session, code)); } else if (this.writableFinished || code) { - setImmediate(rstNextTick.bind(session, this.#id, code)); + setImmediate(rstNextTick.bind(this, session, this.#id, code)); } else { - this.once("finish", rstNextTick.bind(session, this.#id, code)); + this.once("finish", rstNextTick.bind(this, session, this.#id, code)); } // node destroys the stream once both halves have finished; without this a stream closed // while idle never emits 'close'. @@ -2680,11 +2683,10 @@ class Http2Stream extends (Duplex as Http2StreamBase) { session && typeof this.#id === "number" && !this[kNeverAnnounced] && - // A cleanly closed stream the native side already freed has nothing to send: - // the deferred rstStream would be a guaranteed no-op host call per request. - (rstCode !== 0 || (this[bunHTTP2StreamStatus] & StreamState.NativeClosed) === 0) + // Native already closed or reset the stream: nothing is left to send, whatever the rstCode. + (this[bunHTTP2StreamStatus] & (StreamState.NativeClosed | StreamState.NativeReset)) === 0 ) { - setImmediate(rstNextTick.bind(session, this.#id, rstCode)); + setImmediate(rstNextTick.bind(this, session, this.#id, rstCode)); } // Diagnostics channels: published after the stream is closed and destroyed, with the same error @@ -4071,6 +4073,7 @@ class ServerHttp2Session extends Http2Session { }, aborted(self: ServerHttp2Session, stream: ServerHttp2Stream, error: any, old_state: number) { if (!self || typeof stream !== "object") return; + stream[bunHTTP2StreamStatus] |= StreamState.NativeReset; stream.rstCode = constants.NGHTTP2_CANCEL; // if writable and not closed emit aborted if (old_state != 5 && old_state != 7) { @@ -4083,6 +4086,7 @@ class ServerHttp2Session extends Http2Session { }, streamError(self: ServerHttp2Session, stream: ServerHttp2Stream, error: number) { if (!self || typeof stream !== "object") return; + stream[bunHTTP2StreamStatus] |= StreamState.NativeReset; self.#connections--; if (stream.id % 2 === 1) self.#peerInitiatedStreams--; process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed); @@ -5077,6 +5081,7 @@ class ClientHttp2Session extends Http2Session { ), aborted: withStreamFrame((self: ClientHttp2Session, stream: ClientHttp2Stream, error: any, old_state: number) => { if (!self || typeof stream !== "object") return; + stream[bunHTTP2StreamStatus] |= StreamState.NativeReset; stream.rstCode = constants.NGHTTP2_CANCEL; // if writable and not closed emit aborted if (old_state != 5 && old_state != 7) { @@ -5088,7 +5093,7 @@ class ClientHttp2Session extends Http2Session { }), streamError: withStreamFrame((self: ClientHttp2Session, stream: ClientHttp2Stream, error: number) => { if (!self || typeof stream !== "object") return; - + stream[bunHTTP2StreamStatus] |= StreamState.NativeReset; self.#connections--; process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed); }), diff --git a/src/runtime/api/bun/h2/connection.rs b/src/runtime/api/bun/h2/connection.rs index c5955c232829..7458b82d521f 100644 --- a/src/runtime/api/bun/h2/connection.rs +++ b/src/runtime/api/bun/h2/connection.rs @@ -1377,6 +1377,10 @@ impl Connection { discard = true; } Some(st) => { + if st.state == State::ReservedRemote { + self.data_on_reserved_stream(sink, hdr.stream_id); + return StreamedDataStart::Fatal; + } if !stream::can_receive_data(st.state) { self.send_rst_stream(sink, hdr.stream_id, ErrorCode::StreamClosed); if let Some(st2) = self.streams.get_mut(&hdr.stream_id) { @@ -1494,6 +1498,7 @@ impl Connection { // below don't alias the streams map. enum DataDecision { Rst(ErrorCode), + ReservedStream, FlowControlViolation, Deliver(u32), } @@ -1513,7 +1518,9 @@ impl Connection { // §5.1: DATA for an unknown/closed stream is a STREAM_CLOSED error. None => DataDecision::Rst(ErrorCode::StreamClosed), Some(s) => { - if !stream::can_receive_data(s.state) { + if s.state == State::ReservedRemote { + DataDecision::ReservedStream + } else if !stream::can_receive_data(s.state) { DataDecision::Rst(ErrorCode::StreamClosed) } else { s.recv_window.on_data(consumed); @@ -1536,6 +1543,10 @@ impl Connection { sink.on_stream_reset(hdr.stream_id, code.as_u32()); return false; } + DataDecision::ReservedStream => { + self.data_on_reserved_stream(sink, hdr.stream_id); + return true; + } DataDecision::FlowControlViolation => { // nghttp2 (nghttp2_session_update_recv_stream_window_size): a stream flow-control // violation terminates the whole session with FLOW_CONTROL_ERROR; node surfaces @@ -1580,6 +1591,16 @@ impl Connection { false } + /// RFC 9113 §5.1 reserved (remote): DATA is a connection PROTOCOL_ERROR, as in nghttp2. + fn data_on_reserved_stream(&mut self, sink: &impl Sink, stream_id: u32) { + if let Some(s) = self.streams.get_mut(&stream_id) { + s.state = State::Closed; + } + // Session teardown in the embedder misses a promised stream: report this one first. + sink.on_stream_reset(stream_id, ErrorCode::InternalError.as_u32()); + self.send_go_away(sink, ErrorCode::ProtocolError, b"DATA: stream in reserved"); + } + /// RFC 9113 §8.1.1: once END_STREAM arrives, a request whose received DATA total contradicts /// its declared `content-length` is malformed. Resets the stream with PROTOCOL_ERROR instead /// of signalling end-of-stream and returns true if it did so. diff --git a/test/js/node/http2/h2-conformance.test.ts b/test/js/node/http2/h2-conformance.test.ts index 414bf6c3b68b..d8e987e32620 100644 --- a/test/js/node/http2/h2-conformance.test.ts +++ b/test/js/node/http2/h2-conformance.test.ts @@ -7,12 +7,12 @@ // Connection-level cases only here (no HPACK required): preface, SETTINGS handshake/ack, PING, // WINDOW_UPDATE, frame-size and stream-id rules. HPACK/HEADERS cases live in a sibling file. -import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import { afterAll, beforeAll, describe, expect, jest, test } from "bun:test"; import { bunEnv, bunExe, gcTick, normalizeBunSnapshot } from "harness"; import { once } from "node:events"; import http2 from "node:http2"; import net from "node:net"; -import { Writable } from "node:stream"; +import { Duplex, Writable } from "node:stream"; const PREFACE = Buffer.from("PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n", "latin1"); @@ -54,23 +54,41 @@ function encodeFrame(type: number, flags: number, streamId: number, payload: Buf return Buffer.concat([header, payload]); } +/** + * Two in-memory sockets. A write on one is a synchronous push on the other, so each write is + * exactly one read for the session on the far side, with no event-loop turn between two writes. + * Destroying one end destroys the other. + */ +function memoryPair(): [Duplex, Duplex] { + const end = (other: () => Duplex) => + new Duplex({ + read() {}, + write: (chunk, _enc, cb) => (other().push(chunk), cb()), + destroy: (err, cb) => (other().destroy(), cb(err)), + }); + const a: Duplex = end(() => b); + const b: Duplex = end(() => a); + return [a, b]; +} + /** A minimal raw HTTP/2 client: send arbitrary frames, collect parsed inbound frames. */ class RawH2 { - socket: net.Socket; + socket: Duplex; private buf: Buffer = Buffer.alloc(0); frames: Frame[] = []; closed = false; private waiters: Array<{ pred: (f: Frame) => boolean; resolve: (f: Frame) => void }> = []; - constructor(port: number) { - this.socket = net.connect(port, "127.0.0.1"); + /** `socket` is a TCP socket, or one end of a memoryPair() whose other end the server was given. */ + constructor(socket: Duplex) { + this.socket = socket; this.socket.on("data", d => this.onData(d)); this.socket.on("close", () => (this.closed = true)); this.socket.on("error", () => {}); } static async connect(port: number): Promise { - const c = new RawH2(port); + const c = new RawH2(net.connect(port, "127.0.0.1")); await once(c.socket, "connect"); return c; } @@ -675,32 +693,41 @@ describe.concurrent("header block decoding errors (RFC 9113 §4.3)", () => { /** A minimal raw HTTP/2 server: accept one connection, collect parsed inbound frames. */ class RawH2Server { - server: net.Server; - socket: net.Socket | null = null; + server: net.Server | null; + socket: Duplex | null = null; private buf: Buffer = Buffer.alloc(0); private sawPreface = false; frames: Frame[] = []; private waiters: Array<{ pred: (f: Frame) => boolean; resolve: (f: Frame) => void }> = []; - private constructor(server: net.Server) { + private constructor(server: net.Server | null) { this.server = server; } static async listen(): Promise { const server = net.createServer(); const s = new RawH2Server(server); - server.on("connection", socket => { - s.socket = socket; - socket.on("data", d => s.onData(d)); - socket.on("error", () => {}); - }); + server.on("connection", socket => s.accept(socket)); server.listen(0, "127.0.0.1"); await once(server, "listening"); return s; } + /** No listener: `socket` is one end of a memoryPair() whose other end is the client's socket. */ + static over(socket: Duplex): RawH2Server { + const s = new RawH2Server(null); + s.accept(socket); + return s; + } + + private accept(socket: Duplex) { + this.socket = socket; + socket.on("data", d => this.onData(d)); + socket.on("error", () => {}); + } + get port(): number { - return (this.server.address() as net.AddressInfo).port; + return (this.server!.address() as net.AddressInfo).port; } private onData(d: Buffer) { @@ -752,7 +779,7 @@ class RawH2Server { close() { this.socket?.destroy(); - this.server.close(); + this.server?.close(); } } @@ -763,18 +790,31 @@ function hpackLiteral(str: string): Buffer { } describe("push stream states (checklist §5.1, RFC 9113 §6.4/§8.4)", () => { - test("DATA on a promised stream before its response HEADERS is refused, not delivered", async () => { + // RFC 9113 §5.1, reserved (remote): any frame other than HEADERS, RST_STREAM or PRIORITY is a + // connection error of type PROTOCOL_ERROR. nghttp2 does the same ("DATA: stream in reserved"). + // With only part of its payload sent, the frame can reach the engine only as an incomplete + // frame, which the engine parses on a path of its own. + test.each([ + ["a whole DATA frame", 11], + ["the head of a DATA frame", 10], + ])("%s on a promised stream before its response HEADERS is a connection error", async (_, bytesSent) => { const raw = await RawH2Server.listen(); const client = http2.connect(`http://127.0.0.1:${raw.port}`); - client.on("error", () => {}); + const sessionError = Promise.withResolvers(); + client.on("error", err => sessionError.resolve(err)); const pushedData: Buffer[] = []; + const pushedError = jest.fn(); + const pushedClosed = Promise.withResolvers(); client.on("stream", pushed => { - pushed.on("error", () => {}); + pushed.on("error", pushedError); pushed.on("data", (d: Buffer) => pushedData.push(d)); + pushed.on("close", () => pushedClosed.resolve(pushed.rstCode)); }); try { const req = client.request({ ":path": "/" }); req.on("error", () => {}); + const reqClosed = Promise.withResolvers(); + req.on("close", () => reqClosed.resolve(req.rstCode)); await raw.waitFor(f => f.type === FrameType.HEADERS && f.streamId === 1); raw.sendFrame(FrameType.SETTINGS, 0, 0); // server SETTINGS raw.sendFrame(FrameType.SETTINGS, 0x1, 0); // ACK the client's @@ -784,14 +824,20 @@ describe("push stream states (checklist §5.1, RFC 9113 §6.4/§8.4)", () => { promised.writeUInt32BE(2, 0); const block = Buffer.concat([Buffer.from([0x82, 0x86, 0x84, 0x01]), hpackLiteral("localhost")]); raw.sendFrame(FrameType.PUSH_PROMISE, 0x4 /* END_HEADERS */, 1, Buffer.concat([promised, block])); - // DATA on the promised stream while it is still reserved (remote) - §5.1 forbids this - // before the pushed response HEADERS. - raw.sendFrame(FrameType.DATA, 0, 2, Buffer.from("x")); - const rst = await raw.waitFor(f => f.type === FrameType.RST_STREAM && f.streamId === 2); - expect(rst.payload.readUInt32BE(0)).toBe(ErrorCode.STREAM_CLOSED); - // The payload never reaches the pushed stream, and the connection survives. - raw.sendFrame(FrameType.PING, 0, 0, Buffer.alloc(8)); - await raw.waitFor(f => f.type === FrameType.PING && (f.flags & 0x1) !== 0); + // DATA on the promised stream while it is still reserved (remote), before the pushed + // response HEADERS. + const socketClosed = new Promise(resolve => raw.socket!.once("close", () => resolve())); + raw.socket!.write(encodeFrame(FrameType.DATA, 0, 2, Buffer.from("xy")).subarray(0, bytesSent)); + const goaway = await raw.waitFor(f => f.type === FrameType.GOAWAY); + expect(goawayErrorCode(goaway)).toBe(ErrorCode.PROTOCOL_ERROR); + expect((await sessionError.promise).code).toBe("ERR_HTTP2_ERROR"); + expect(await reqClosed.promise).toBe(http2.constants.NGHTTP2_INTERNAL_ERROR); + expect(await pushedClosed.promise).toBe(http2.constants.NGHTTP2_INTERNAL_ERROR); + expect(pushedError).toHaveBeenCalledTimes(1); + // Once the client's socket is gone, every frame it wrote is in raw.frames. The reserved + // stream gets no RST_STREAM of its own, and its payload never reaches the pushed stream. + await socketClosed; + expect(raw.frames.filter(f => f.type === FrameType.RST_STREAM)).toEqual([]); expect(Buffer.concat(pushedData).length).toBe(0); } finally { client.destroy(); @@ -1718,6 +1764,267 @@ describe("inbound stream lifecycle", () => { }); }); +// Once the native layer has closed a stream (it reset the stream, the peer reset it, or both +// END_STREAM flags went by), nothing more belongs on the wire for that stream. The JS stream is +// destroyed afterwards, and a reset it submits then is dropped by the native layer only while the +// stream's table entry exists. A client's pushed stream never has an entry and any other stream +// loses its entry on the next read, so that reset used to become a second RST_STREAM, or an answer +// to the peer's RST_STREAM. Node v26.3.0 sends neither. +describe("RST_STREAM on a stream the native layer already closed", () => { + const END_HEADERS = 0x4; + // `:status: 200`, then `connection: close`. A connection-specific field makes the block + // malformed (RFC 9113 §8.2.2): a stream error of type PROTOCOL_ERROR. + const MALFORMED_RESPONSE = Buffer.concat([ + Buffer.from([0x88, 0x00]), + hpackLiteral("connection"), + hpackLiteral("close"), + ]); + + function u32(value: number): Buffer { + const buf = Buffer.alloc(4); + buf.writeUInt32BE(value, 0); + return buf; + } + + let pings = 0; + /** Once the ACK is in, so is every frame the peer wrote before it read this PING. */ + async function pingRoundTrip(peer: RawH2 | RawH2Server) { + const payload = u32(++pings); + const opaque = Buffer.concat([payload, payload]); + peer.sendFrame(FrameType.PING, 0, 0, opaque); + await peer.waitFor(f => f.type === FrameType.PING && (f.flags & 0x1) !== 0 && f.payload.equals(opaque)); + } + + /** + * The codes of the RST_STREAM frames the peer wrote on `streamId`. Call it after the stream's + * 'close': a reset that close() or _destroy deferred with setImmediate was queued before that + * event, so it runs before the immediate awaited here. + */ + async function resetsWritten(peer: RawH2 | RawH2Server, streamId: number): Promise { + await new Promise(resolve => setImmediate(resolve)); + await pingRoundTrip(peer); + return peer.frames + .filter(f => f.type === FrameType.RST_STREAM && f.streamId === streamId) + .map(f => f.payload.readUInt32BE(0)); + } + + /** Opens a request on stream 1 and completes the raw server's half of the handshake. */ + async function openRequest( + raw: RawH2Server, + client: http2.ClientHttp2Session, + headers: http2.OutgoingHttpHeaders = { ":path": "/" }, + ) { + client.on("error", () => {}); + const req = client.request(headers); + req.on("error", () => {}); + await raw.waitFor(f => f.type === FrameType.HEADERS && f.streamId === 1); + raw.sendFrame(FrameType.SETTINGS, 0, 0); + raw.sendFrame(FrameType.SETTINGS, 0x1, 0); + return req; + } + + /** Resolves on 'close'. events.once() would reject on the stream's 'error' first. */ + function closeOf(stream: http2.Http2Stream): Promise { + const { promise, resolve } = Promise.withResolvers(); + stream.on("close", () => resolve()); + return promise; + } + + /** PUSH_PROMISE on stream 1 that reserves stream 2. */ + function sendPushPromise(raw: RawH2Server) { + raw.sendFrame(FrameType.PUSH_PROMISE, END_HEADERS, 1, Buffer.concat([u32(2), requestHeaderBlock("GET")])); + } + + /** Resolves to the pushed stream's rstCode once that stream has closed. */ + function pushedStreamClosed(client: http2.ClientHttp2Session): Promise { + const { promise, resolve } = Promise.withResolvers(); + client.on("stream", pushed => { + pushed.on("error", () => {}); + pushed.on("close", () => resolve(pushed.rstCode)); + }); + return promise; + } + + /** + * Sends a POST on stream 1 from the raw client and resolves once the server has the stream. + * `closed` resolves to the server stream's rstCode. + */ + async function openServerStream(server: http2.Http2Server, c: RawH2) { + const opened = Promise.withResolvers(); + const closed = Promise.withResolvers(); + server.on("stream", (stream: any) => { + stream.on("error", () => {}); + stream.on("close", () => closed.resolve(stream.rstCode)); + opened.resolve(); + }); + c.sendPreface(); + c.sendEmptySettings(); + c.sendFrame(FrameType.HEADERS, END_HEADERS, 1, requestHeaderBlock("POST")); + await opened.promise; + return { closed: closed.promise }; + } + + test("a client resets a pushed stream once for a malformed response block", async () => { + const raw = await RawH2Server.listen(); + const client = http2.connect(`http://127.0.0.1:${raw.port}`); + try { + const closed = pushedStreamClosed(client); + await openRequest(raw, client); + sendPushPromise(raw); + raw.sendFrame(FrameType.HEADERS, END_HEADERS, 2, MALFORMED_RESPONSE); + expect(await closed).toBe(ErrorCode.PROTOCOL_ERROR); + expect(await resetsWritten(raw, 2)).toEqual([ErrorCode.PROTOCOL_ERROR]); + } finally { + client.destroy(); + raw.close(); + } + }); + + test.each([ + ["INTERNAL_ERROR", ErrorCode.INTERNAL_ERROR], + ["CANCEL", ErrorCode.CANCEL], + ])("a client does not answer the peer's RST_STREAM(%s) on a pushed stream", async (_, code) => { + const raw = await RawH2Server.listen(); + const client = http2.connect(`http://127.0.0.1:${raw.port}`); + try { + const closed = pushedStreamClosed(client); + await openRequest(raw, client); + sendPushPromise(raw); + raw.sendFrame(FrameType.HEADERS, END_HEADERS, 2, Buffer.from([0x88])); + raw.sendFrame(FrameType.RST_STREAM, 0, 2, u32(code)); + expect(await closed).toBe(code); + expect(await resetsWritten(raw, 2)).toEqual([]); + } finally { + client.destroy(); + raw.close(); + } + }); + + test("a client resets a request stream once when another read follows the reset", async () => { + const [near, far] = memoryPair(); + const raw = RawH2Server.over(near); + const client = http2.connect("http://localhost", { createConnection: () => far }); + try { + const req = await openRequest(raw, client); + const closed = closeOf(req); + raw.sendFrame(FrameType.HEADERS, END_HEADERS, 1, MALFORMED_RESPONSE); + // A second read in the same turn: stream 1 loses its table entry before it is destroyed. + raw.sendFrame(FrameType.WINDOW_UPDATE, 0, 0, u32(1)); + await closed; + expect(await resetsWritten(raw, 1)).toEqual([ErrorCode.PROTOCOL_ERROR]); + } finally { + client.destroy(); + raw.close(); + } + }); + + // close(code) and the _destroy that follows it each submit a reset. Native answers the first + // one with a streamError dispatch, so the second one is for a stream it already closed. + test("close(code) resets a request stream once when a read follows the RST_STREAM", async () => { + const [near, far] = memoryPair(); + const raw = RawH2Server.over(near); + const client = http2.connect("http://localhost", { createConnection: () => far }); + try { + // The request body stays open, so close() has a stream to reset. + const req = await openRequest(raw, client, { ":method": "POST", ":path": "/" }); + const closed = closeOf(req); + req.close(http2.constants.NGHTTP2_CANCEL); + await raw.waitFor(f => f.type === FrameType.RST_STREAM && f.streamId === 1); + // This read takes stream 1's table entry away before the reset from _destroy runs. + raw.sendFrame(FrameType.WINDOW_UPDATE, 0, 0, u32(1)); + await closed; + expect(await resetsWritten(raw, 1)).toEqual([ErrorCode.CANCEL]); + } finally { + client.destroy(); + raw.close(); + } + }); + + test("a close() from an 'aborted' listener does not answer the peer's RST_STREAM", async () => { + const [near, far] = memoryPair(); + const raw = RawH2Server.over(near); + const client = http2.connect("http://localhost", { createConnection: () => far }); + try { + // The request body stays open, so the peer's reset is an abort. + const req = await openRequest(raw, client, { ":method": "POST", ":path": "/" }); + const onAborted = jest.fn(() => req.close(http2.constants.NGHTTP2_CANCEL)); + req.on("aborted", onAborted); + const closed = closeOf(req); + raw.sendFrame(FrameType.RST_STREAM, 0, 1, u32(ErrorCode.CANCEL)); + // A second read in the same turn: stream 1 loses its table entry before close()'s reset runs. + raw.sendFrame(FrameType.WINDOW_UPDATE, 0, 0, u32(1)); + await closed; + expect(onAborted).toHaveBeenCalledTimes(1); + expect(await resetsWritten(raw, 1)).toEqual([]); + } finally { + client.destroy(); + raw.close(); + } + }); + + test("destroy(err) on a request stream that closed cleanly sends no RST_STREAM", async () => { + const raw = await RawH2Server.listen(); + const client = http2.connect(`http://127.0.0.1:${raw.port}`); + try { + const req = await openRequest(raw, client); + // Nothing reads the body, so the JS stream outlives the native one. + raw.sendFrame(FrameType.HEADERS, END_HEADERS, 1, Buffer.from([0x88])); + raw.sendFrame(FrameType.DATA, 0x1 /* END_STREAM */, 1, Buffer.from("hello")); + await pingRoundTrip(raw); + // The second PING is a read of its own, after the one that carried END_STREAM. + await pingRoundTrip(raw); + expect({ closed: req.closed, destroyed: req.destroyed }).toEqual({ closed: true, destroyed: false }); + const closed = closeOf(req); + req.destroy(new Error("late")); + await closed; + expect(await resetsWritten(raw, 1)).toEqual([]); + } finally { + client.destroy(); + raw.close(); + } + }); + + test("a server does not answer the peer's RST_STREAM(CANCEL) when another read follows it", async () => { + const server = http2.createServer(); + const [near, far] = memoryPair(); + const c = new RawH2(near); + server.emit("connection", far); + try { + const { closed } = await openServerStream(server, c); + c.sendFrame(FrameType.RST_STREAM, 0, 1, u32(ErrorCode.CANCEL)); + // A second read in the same turn: stream 1 loses its table entry before it is destroyed. + c.sendFrame(FrameType.WINDOW_UPDATE, 0, 0, u32(1)); + expect(await closed).toBe(ErrorCode.CANCEL); + expect(await resetsWritten(c, 1)).toEqual([]); + } finally { + c.destroy(); + } + }); + + // The reset a server stream submits with REFUSED_STREAM is counted as a stream this side + // rejected. At a budget of 1 that count ends the session with GOAWAY(ENHANCE_YOUR_CALM). + test("a peer's RST_STREAM(REFUSED_STREAM) does not count against maxSessionRejectedStreams", async () => { + const server = http2.createServer({ maxSessionRejectedStreams: 1 }); + server.listen(0); + await once(server, "listening"); + const c = await RawH2.connect((server.address() as net.AddressInfo).port); + try { + const { closed } = await openServerStream(server, c); + c.sendFrame(FrameType.RST_STREAM, 0, 1, u32(ErrorCode.REFUSED_STREAM)); + expect(await closed).toBe(ErrorCode.REFUSED_STREAM); + await new Promise(resolve => setImmediate(resolve)); + c.sendFrame(FrameType.PING, 0, 0, Buffer.alloc(8)); + const answer = await c.waitFor( + f => f.type === FrameType.GOAWAY || (f.type === FrameType.PING && (f.flags & 0x1) !== 0), + ); + expect(answer.type).toBe(FrameType.PING); + } finally { + c.destroy(); + server.close(); + } + }); +}); + // A DATA frame that cannot be written right away (the peer's flow-control window is used up, the // socket has backpressure, or another stream on the session already has frames waiting) is put on // the session's outbound queue and written later, when a WINDOW_UPDATE or a writable socket drains diff --git a/test/js/node/http2/h2-push-refusal-staged.test.ts b/test/js/node/http2/h2-push-refusal-staged.test.ts index a28ceda1fd07..ab0f48b539d3 100644 --- a/test/js/node/http2/h2-push-refusal-staged.test.ts +++ b/test/js/node/http2/h2-push-refusal-staged.test.ts @@ -2,10 +2,11 @@ import { expect, test } from "bun:test"; import http2 from "node:http2"; import net from "node:net"; -// Staged twin of h2-conformance's "DATA on a promised stream before its -// response HEADERS is refused, not delivered", which times out waiting for a -// frame on the darwin agents. This variant tapes the client session's events -// and every frame the raw server receives, and names the stage that stalled. +// Staged twin of h2-conformance's "a whole DATA frame on a promised stream +// before its response HEADERS is a connection error", which timed out waiting +// for a frame on the darwin agents. This variant tapes the client session's +// events and every frame the raw server receives, and names the stage that +// stalled. function frame(len: number, type: number, flags: number, id: number, payload = Buffer.alloc(0)) { const h = Buffer.alloc(9); h.writeUIntBE(len, 0, 3); @@ -25,7 +26,7 @@ const TYPE_NAME: Record = { 8: "WINDOW_UPDATE", }; -test("DATA on a reserved push stream is refused with RST(STREAM_CLOSED) (event-taped)", async () => { +test("DATA on a reserved push stream is a connection error, GOAWAY(PROTOCOL_ERROR) (event-taped)", async () => { const tape: string[] = []; const t = (name: string) => tape.push(name); const frames: Array<{ type: number; flags: number; id: number; payload: Buffer }> = []; @@ -102,10 +103,11 @@ test("DATA on a reserved push stream is refused with RST(STREAM_CLOSED) (event-t socket.write(frame(1, 0, 0, 2, Buffer.from("x"))); t("sent-data-on-reserved"); - const rst = await waitFor(f => f.type === 3 && f.id === 2, "RST on stream 2"); - expect(rst.payload.readUInt32BE(0)).toBe(5 /* STREAM_CLOSED */); - socket.write(frame(8, 6, 0, 0, Buffer.alloc(8))); - await waitFor(f => f.type === 6 && (f.flags & 0x1) !== 0, "PING ACK"); + // RFC 9113 §5.1, reserved (remote): a connection error of type PROTOCOL_ERROR, as in nghttp2. + const goaway = await waitFor(f => f.type === 7, "GOAWAY"); + expect(goaway.payload.readUInt32BE(4)).toBe(1 /* PROTOCOL_ERROR */); + // A RST_STREAM written ahead of the GOAWAY would already be here. + expect(frames.filter(f => f.type === 3)).toEqual([]); expect(Buffer.concat(pushedData).length).toBe(0); } finally { client.destroy();