diff --git a/src/runtime/api/bun/h2_frame_parser.rs b/src/runtime/api/bun/h2_frame_parser.rs index d0e9f4267aee..bff378303e42 100644 --- a/src/runtime/api/bun/h2_frame_parser.rs +++ b/src/runtime/api/bun/h2_frame_parser.rs @@ -2821,6 +2821,11 @@ impl H2FrameParser { self.write_buffer.get().len_u32() > 0 || self.has_nonnative_backpressure.get() } + /// Per stream, not per session: another stream's queue can wait on its own window forever. + fn must_queue_data(&self, stream: &Stream) -> bool { + self.has_backpressure() || !stream.data_frame_queue.is_empty() + } + /// Whether a write to this session's transport synchronously runs user JS: a JS-backed /// socket's onWrite is the user's Duplex, and a socket upgraded from a JS Duplex /// (`tls.connect({ socket })`) writes its records through that Duplex. @@ -5157,7 +5162,7 @@ impl H2FrameParser { stream_identifier: stream_id, length: 0, }; - if self.has_backpressure() || self.outbound_queue_size.get() > 0 { + if self.must_queue_data(stream) { enqueued = true; stream.queue_frame(self, b"", callback, close); } else { @@ -5195,10 +5200,7 @@ impl H2FrameParser { offset += size; let end_stream = offset >= payload.len() && can_close; - if self.has_backpressure() - || self.outbound_queue_size.get() > 0 - || is_flow_control_limited - { + if self.must_queue_data(stream) || is_flow_control_limited { // Preserve wire order: anything already batched goes out before the // queued remainder is flushed later by the drain path. self.flush_batch_buffer(); diff --git a/test/js/node/http2/h2-conformance.test.ts b/test/js/node/http2/h2-conformance.test.ts index 414bf6c3b68b..d469d3ff0298 100644 --- a/test/js/node/http2/h2-conformance.test.ts +++ b/test/js/node/http2/h2-conformance.test.ts @@ -1719,15 +1719,17 @@ describe("inbound stream lifecycle", () => { }); // 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 -// the queue. Writing the last queued frame of a stream whose peer half is already closed completes -// the stream, and, exactly like the direct-write path, has to release it (the JS stream object and +// socket has backpressure, or the same stream already has frames waiting) is put on the stream's +// outbound queue and written later, when a WINDOW_UPDATE or a writable socket drains the queue. +// Writing the last queued frame of a stream whose peer half is already closed completes the +// stream, and, exactly like the direct-write path, has to release it (the JS stream object and // the native entry) while the session lives on. Node releases these streams too; a stream that is -// only released at session teardown is a per-request leak on a long-lived session. END_STREAM can -// ride on the queued frame that carries the last of the body or on an empty frame queued by itself -// (end() without a body, or the empty frame that follows a body once no trailers are coming); the -// queue writes the two through different branches, so both shapes are covered below. +// only released at session teardown is a per-request leak on a long-lived session. The first and +// the last case below reach the queue through the stream's own window. The cases that answer +// behind another stream's stalled response do not: a stream waits only behind its own queue, so +// those responses are written directly. They cover the same release on a session that carries a +// stalled stream, with END_STREAM on the body's frame and on an empty frame of its own (end() +// without a body, or the empty frame that follows a body once no trailers are coming). describe("stream release after a queued END_STREAM", () => { // Well above GC_STRAGGLERS: without the release every one of these survives. const STREAMS = 16; @@ -1850,7 +1852,7 @@ describe("stream release after a queued END_STREAM", () => { refs.push(new WeakRef(stream)); stream.resume(); // Answer once the request's END_STREAM has been processed, like a handler that consumes the - // request body does: the queued response is then the only thing the stream still waits for. + // request body does: the response is then the only thing the stream still waits for. stream.on("end", () => { stream.respond({ ":status": 200 }); finish(stream); diff --git a/test/js/node/http2/node-http2.test.js b/test/js/node/http2/node-http2.test.js index 7cf5c3e7ab3a..57e28d46ed2a 100644 --- a/test/js/node/http2/node-http2.test.js +++ b/test/js/node/http2/node-http2.test.js @@ -5628,6 +5628,97 @@ it("Http2Stream pull-mode read() after pause() replenishes the receive window", } }); +// DATA that waits for stream A's own send window must not hold back stream B: B's window and +// the connection window still have credit. The raw client never sends WINDOW_UPDATE, so after +// a write on B the server receives nothing that would flush a queued frame. The PING is the +// first inbound traffic after that write: a frame that needed it shows up after the PING ACK. +it("http2 sends DATA on one stream while another stream's DATA waits for its own send window", async () => { + const STREAM_WINDOW = 16384; // the connection window stays at the default 65535 + const server = http2.createServer(); + let socket; + try { + const serverStreams = []; + server.on("stream", stream => { + stream.on("error", () => {}); + if (serverStreams.push(stream) === 1) { + // A: more than its send window. The rest stays queued for good. + stream.respond({ ":status": 200 }); + stream.write(Buffer.alloc(4 * STREAM_WINDOW, "a")); + } else { + stream.respond({ ":status": 200 }, { waitForTrailers: true }); + } + }); + await new Promise(resolve => server.listen(0, "127.0.0.1", resolve)); + + // HPACK: :method GET, :scheme http, :path /, :authority localhost (literal, 7-bit length). + const requestBlock = Buffer.concat([Buffer.from([0x82, 0x86, 0x84, 0x01, 9]), Buffer.from("localhost")]); + const settings = http2.getPackedSettings({ initialWindowSize: STREAM_WINDOW }); + socket = net.connect(server.address().port, "127.0.0.1", () => { + socket.write(http2utils.kClientMagic); + socket.write(Buffer.concat([new http2utils.Frame(settings.length, 4, 0, 0).data, settings])); + socket.write(new http2utils.HeadersFrame(1, requestBlock, 0, true, true).data); // A + socket.write(new http2utils.HeadersFrame(3, requestBlock, 0, true, true).data); // B + }); + + const wireOrder = []; // B's DATA frames and PING ACKs + let aBytes = 0; + let bResponded = false; + let closed; + let progress = Promise.withResolvers(); + const notify = () => { + progress.resolve(); + progress = Promise.withResolvers(); + }; + // Waits until the frames received so far satisfy `condition`. + const until = async condition => { + while (!condition()) { + if (closed) throw closed; + await progress.promise; + } + }; + socket.on("error", err => ((closed = err), notify())); + socket.on("close", () => ((closed ??= new Error("the server closed the connection")), notify())); + let received = Buffer.alloc(0); + socket.on("data", chunk => { + received = Buffer.concat([received, chunk]); + while (received.length >= 9) { + const length = received.readUIntBE(0, 3); + if (received.length < 9 + length) break; + const type = received[3]; + const endFlag = (received[4] & 1) !== 0; // END_STREAM on DATA, ACK on PING + const streamId = received.readUInt32BE(5) & 0x7fffffff; + const payload = received.subarray(9, 9 + length).toString(); + received = received.subarray(9 + length); + if (type === 0 && streamId === 1) aBytes += length; + if (type === 1 && streamId === 3) bResponded = true; + if (type === 0 && streamId === 3) wireOrder.push({ data: payload, endStream: endFlag }); + if (type === 6 && endFlag) wireOrder.push("PING ACK"); + } + notify(); + }); + + await until(() => aBytes === STREAM_WINDOW && bResponded); + const b = serverStreams[1]; + b.write("hello from b"); + socket.write(new http2utils.PingFrame(false).data); + await until(() => wireOrder.includes("PING ACK")); + expect(wireOrder).toEqual([{ data: "hello from b", endStream: false }, "PING ACK"]); + + // sendTrailers({}) ends B with an empty DATA frame, which has its own path in the writer. + // Node sends that frame from a later event loop phase, so no PING here: wait for the frame. + const wantTrailers = new Promise(resolve => b.once("wantTrailers", resolve)); + b.end(); + await wantTrailers; + b.sendTrailers({}); + await until(() => wireOrder.length === 3); + expect(wireOrder[2]).toEqual({ data: "", endStream: true }); + expect(aBytes).toBe(STREAM_WINDOW); + } finally { + socket?.destroy(); + server.close(); + } +}); + // The outbound cork buffer is thread-local across every Http2Session. Interleaving // respond()/write() across two sessions used to let the second session's corked // HEADERS be prepended to the first session's multi-frame DATA batch and sent to