Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 7 additions & 5 deletions src/runtime/api/bun/h2_frame_parser.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

/// 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.
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Comment thread
robobun marked this conversation as resolved.
// Preserve wire order: anything already batched goes out before the
// queued remainder is flushed later by the drain path.
self.flush_batch_buffer();
Expand Down
20 changes: 11 additions & 9 deletions test/js/node/http2/h2-conformance.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down
91 changes: 91 additions & 0 deletions test/js/node/http2/node-http2.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading