diff --git a/src/runtime/api/bun/h2/connection.rs b/src/runtime/api/bun/h2/connection.rs index f3c6232711ce..5635bb61d240 100644 --- a/src/runtime/api/bun/h2/connection.rs +++ b/src/runtime/api/bun/h2/connection.rs @@ -72,6 +72,17 @@ pub(crate) enum WriteResult { Sent = 1, } +/// What `Sink::credit_send_window` did with a WINDOW_UPDATE increment. +#[derive(Clone, Copy, PartialEq, Eq, Debug)] +pub(crate) enum SendCredit { + /// The embedder keeps no send window for this id: the engine's own window takes the increment. + NotOwned, + /// The embedder added the increment to the window it sends against. + Applied, + /// §6.9.1: the increment takes that window past 2^31-1. Nothing was added. + Overflow, +} + /// Outcome of feeding bytes: how many were consumed, and whether the connection is now closing. #[derive(Clone, Copy, Debug)] enum StreamedDataStart { @@ -199,8 +210,10 @@ pub(crate) trait Sink { fn on_ping(&self, payload: &[u8], is_ack: bool); /// `code` is the raw u32 from the wire so unknown error codes survive to JS (node parity). fn on_go_away(&self, code: u32, last_stream_id: u32, debug: &[u8]); - /// After a WINDOW_UPDATE has been applied (for resuming sends). - fn on_window_update(&self, stream_id: u32, increment: u32); + /// A WINDOW_UPDATE with a non-zero increment for `stream_id` (0 = the connection). + fn credit_send_window(&self, _stream_id: u32, _increment: u32) -> SendCredit { + SendCredit::NotOwned + } /// The embedder cannot take further callbacks in this batch (its VM has an exception pending /// from an earlier one): stop before the next frame; the unconsumed bytes stay queued. @@ -632,6 +645,15 @@ impl Connection { } } + /// INITIAL_WINDOW_SIZE as last sent. setLocalWindowSize() raises `local_settings` past it. + fn advertised_initial_window(&self) -> u32 { + let sent = match self.pending_local_settings_acks.back() { + Some(pending) => pending.settings.initial_window_size, + None => self.acked_local_initial_window, + }; + sent.min(self.local_settings.initial_window_size) + } + /// Send WINDOW_UPDATE for every receive window that has consumed at least half its size. fn replenish_windows(&mut self, sink: &impl Sink) { if self.recv_window.needs_update() { @@ -642,9 +664,10 @@ impl Connection { } let mut buf = std::mem::take(&mut self.replenish_buf); buf.clear(); + let advertised = self.advertised_initial_window(); for (id, s) in self.streams.iter_mut() { if s.state != State::Closed - && s.recv_window.needs_update() + && s.recv_window.needs_update_within(advertised) && sink.is_stream_reading(*id) { let inc = s.recv_window.take_update(); @@ -971,28 +994,38 @@ impl Connection { sink.on_stream_reset(hdr.stream_id, ErrorCode::ProtocolError.as_u32()); return false; } + let overflow = match sink.credit_send_window(hdr.stream_id, increment) { + SendCredit::Applied => false, + SendCredit::Overflow => true, + SendCredit::NotOwned => { + let window = if hdr.stream_id == 0 { + Some(&mut self.send_window) + } else { + self.streams + .get_mut(&hdr.stream_id) + .map(|s| &mut s.send_window) + }; + window.is_some_and(|w| w.increase(increment).is_err()) + } + }; + if !overflow { + return false; + } if hdr.stream_id == 0 { // 6.9.1: the connection window must not exceed 2^31-1. - if self.send_window.increase(increment).is_err() { - self.send_go_away( - sink, - ErrorCode::FlowControlError, - b"connection flow-control window overflow", - ); - return true; - } - } else if let Some(s) = self.streams.get_mut(&hdr.stream_id) { - // 6.9.1: a per-stream overflow is a stream error, not a connection error. - if s.send_window.increase(increment).is_err() { - self.send_rst_stream(sink, hdr.stream_id, ErrorCode::FlowControlError); - if let Some(s) = self.streams.get_mut(&hdr.stream_id) { - s.state = State::Closed; - } - sink.on_stream_reset(hdr.stream_id, ErrorCode::FlowControlError.as_u32()); - return false; - } + self.send_go_away( + sink, + ErrorCode::FlowControlError, + b"connection flow-control window overflow", + ); + return true; + } + // 6.9.1: a per-stream overflow is a stream error, not a connection error. + self.send_rst_stream(sink, hdr.stream_id, ErrorCode::FlowControlError); + if let Some(s) = self.streams.get_mut(&hdr.stream_id) { + s.state = State::Closed; } - sink.on_window_update(hdr.stream_id, increment); + sink.on_stream_reset(hdr.stream_id, ErrorCode::FlowControlError.as_u32()); false } @@ -1472,8 +1505,7 @@ impl Connection { let data_total = hdr.length as usize - off - pad; // §6.9: the whole declared frame counts against the connection recv window on receipt. - self.recv_window.on_data(hdr.length as i64); - if self.recv_window.is_overflowed() { + if self.recv_window.on_data(hdr.length as i64) && self.recv_window.is_overflowed() { self.send_go_away( sink, ErrorCode::FlowControlError, @@ -1508,8 +1540,9 @@ impl Connection { sink.on_stream_reset(hdr.stream_id, ErrorCode::StreamClosed.as_u32()); discard = true; } else { - st.recv_window.on_data(hdr.length as i64); - if st.recv_window.is_overflowed_with(recv_limit) { + if st.recv_window.on_data(hdr.length as i64) + && st.recv_window.is_overflowed_with(recv_limit) + { // nghttp2 (nghttp2_session_update_recv_stream_window_size): a peer that // violates a stream's flow-control window terminates the whole session // with FLOW_CONTROL_ERROR; node surfaces it as NghttpError "Protocol @@ -1590,8 +1623,7 @@ impl Connection { let consumed = payload.len() as i64; // full frame counts against flow control, incl. padding // §6.9: the whole frame counts against the connection recv window. - self.recv_window.on_data(consumed); - if self.recv_window.is_overflowed() { + if self.recv_window.on_data(consumed) && self.recv_window.is_overflowed() { self.send_go_away( sink, ErrorCode::FlowControlError, @@ -1639,8 +1671,9 @@ impl Connection { if !stream::can_receive_data(s.state) { DataDecision::Rst(ErrorCode::StreamClosed) } else { - s.recv_window.on_data(consumed); - if s.recv_window.is_overflowed_with(recv_limit) { + if s.recv_window.on_data(consumed) + && s.recv_window.is_overflowed_with(recv_limit) + { DataDecision::FlowControlViolation } else { s.recv_body_bytes = s.recv_body_bytes.saturating_add((end - off) as u64); @@ -2054,8 +2087,11 @@ impl Connection { /// pause). Without this, a peer stalled on a zero stream window would only be released by the /// next inbound batch — which may never come, since the peer is the one waiting. pub(crate) fn replenish_stream(&mut self, sink: &impl Sink, stream_id: u32) { + let advertised = self.advertised_initial_window(); let inc = match self.streams.get_mut(&stream_id) { - Some(s) if s.state != State::Closed && s.recv_window.needs_update() => { + Some(s) + if s.state != State::Closed && s.recv_window.needs_update_within(advertised) => + { s.recv_window.take_update() } _ => return, @@ -2161,7 +2197,6 @@ mod tests { fn on_go_away(&self, c: u32, l: u32, _d: &[u8]) { self.goaway.set(Some((c, l))); } - fn on_window_update(&self, _id: u32, _inc: u32) {} fn on_stream_open(&self, id: u32) { self.opens.borrow_mut().push(id); } diff --git a/src/runtime/api/bun/h2/flow_control.rs b/src/runtime/api/bun/h2/flow_control.rs index fbfc2629e904..888ca54b6bf9 100644 --- a/src/runtime/api/bun/h2/flow_control.rs +++ b/src/runtime/api/bun/h2/flow_control.rs @@ -92,9 +92,12 @@ impl RecvWindow { } } + /// Counts a DATA frame of `n` bytes. False for an empty frame: it takes no window (§6.9.1). #[inline] - pub(crate) fn on_data(&mut self, n: i64) { + #[must_use] + pub(crate) fn on_data(&mut self, n: i64) -> bool { self.consumed += n; + n > 0 } /// Whether the peer exceeded our advertised window (a FLOW_CONTROL_ERROR, §6.9.1). @@ -116,6 +119,12 @@ impl RecvWindow { self.consumed > 0 && self.consumed >= self.size / 2 } + /// `needs_update` for a stream, capped at `advertised`: the local INITIAL_WINDOW_SIZE (§6.9.2). + #[inline] + pub(crate) fn needs_update_within(&self, advertised: u32) -> bool { + self.consumed > 0 && self.consumed >= self.size.min(advertised as i64) / 2 + } + /// Take the pending WINDOW_UPDATE increment and reset the consumed counter (0 if none). pub(crate) fn take_update(&mut self) -> u32 { if self.consumed <= 0 { @@ -145,7 +154,7 @@ mod tests { #[test] fn recv_window_replenish() { let mut w = RecvWindow::new(100); - w.on_data(60); + assert!(w.on_data(60)); assert!(w.needs_update()); assert_eq!(w.take_update(), 60); assert_eq!(w.consumed, 0); diff --git a/src/runtime/api/bun/h2_frame_parser.rs b/src/runtime/api/bun/h2_frame_parser.rs index 46a7ea31dd49..5408f1cd6d68 100644 --- a/src/runtime/api/bun/h2_frame_parser.rs +++ b/src/runtime/api/bun/h2_frame_parser.rs @@ -1089,14 +1089,6 @@ pub(crate) struct H2FrameParser { /// Receive-window growth requested by setLocalWindowSize() while a dispatch held the engine /// borrow; applied by rewrite_read() on its next pass. pending_recv_window_growth: Cell, - /// Bridge: outbound DATA bytes the legacy encoder wrote since the engine last ran, applied to - /// the engine's connection-level send window in rewrite_read (the engine cell may be borrowed - /// when the DATA goes out). Without this the engine's window only ever grows and a compliant - /// peer's cumulative WINDOW_UPDATEs would eventually trip the §6.9.1 overflow error. - pending_send_window_consumed: Cell, - /// Same bridge per stream: (stream id, bytes) pairs drained into the engine's per-stream send - /// windows in rewrite_read. - pending_stream_send_consumed: JsCell>, /// Local-settings snapshot of each SETTINGS frame the legacy encoder sent, in send order, /// drained into the engine's per-SETTINGS ack queue in rewrite_read (§6.5.3: ACKs apply in /// order). @@ -1572,7 +1564,6 @@ impl Stream { client .remote_used_window_size .set(client.remote_used_window_size.get() + payload_size as u64); - client.note_engine_send_consumed(self.id, payload_size as u64); let mut flags: u8 = 0; // we ignore end_stream for now because we know we have more data to send if padding != 0 { @@ -1621,7 +1612,6 @@ impl Stream { client .remote_used_window_size .set(client.remote_used_window_size.get() + payload_size as u64); - client.note_engine_send_consumed(self.id, payload_size as u64); let mut flags: u8 = if frame.end_stream && !self.wait_for_trailers { DataFrameFlags::END_STREAM as u8 } else { @@ -3521,26 +3511,6 @@ impl H2FrameParser { JSValue::UNDEFINED } - /// Record outbound DATA the legacy encoder wrote so the engine's send windows track reality. - /// Buffered in cells and applied in rewrite_read: inbound WINDOW_UPDATE handling always goes - /// through rewrite_read first, so the windows are in sync before any overflow check runs. - pub(crate) fn note_engine_send_consumed(&self, stream_id: u32, n: u64) { - if n == 0 { - return; - } - self.pending_send_window_consumed - .set(self.pending_send_window_consumed.get() + n); - self.pending_stream_send_consumed.with_mut(|v| { - if let Some(last) = v.last_mut() - && last.0 == stream_id - { - last.1 += n; - } else { - v.push((stream_id, n)); - } - }); - } - /// Mirror the engine's frame counters into plain Cells so getFrameCounters() never /// contends with the engine borrow (destroy can run inside a dispatch). fn sync_engine_frame_counters(&self) { @@ -3588,29 +3558,6 @@ impl H2FrameParser { if pending > 0 { engine.recv_window.grow(pending); } - // Apply outbound DATA the legacy encoder wrote since the last batch, so the engine's - // send windows reflect what is actually in flight (§6.9.1 overflow stays peer-error - // only). - let sent = self.pending_send_window_consumed.replace(0); - if sent > 0 { - engine.send_window.consume(sent as i64); - } - self.pending_stream_send_consumed.with_mut(|v| { - // A client-initiated stream has no engine entry until its first inbound - // frame; dropping its consume here would leave the engine's send window - // permanently wider than the peer's view. Keep unmatched entries queued — - // but only while the legacy stream is still alive: a pushed stream the - // client never sends frames on would otherwise park its entry forever - // (and get scanned on every read). - v.retain(|&(id, n)| { - if let Some(s) = engine.streams.get_mut(&id) { - s.send_window.consume(n as i64); - false - } else { - self.streams.get().contains_key(&id) - } - }); - }); // Register SETTINGS submissions the legacy encoder sent since the last batch, so the // engine attributes each inbound ACK to the right submission (§6.5.3). self.pending_settings_window_submissions.with_mut(|v| { @@ -3822,21 +3769,24 @@ impl crate::api::h2::connection::Sink for H2FrameParser { max_header_list_size: settings.max_header_list_size, enable_connect_protocol: settings.enable_connect_protocol, }; + // handle_received_stream_id opens every stream with this value. + let old_initial_window = self + .remote_settings + .get() + .map(|s| s.initial_window_size) + .unwrap_or(DEFAULT_WINDOW_SIZE as u32); self.remote_settings.set(Some(fp)); - // §6.9.2 (mirrors the legacy inbound): when the peer's INITIAL_WINDOW_SIZE grows, raise the - // send window of streams opened before its SETTINGS arrived (a client's first request is - // typically sent before the server's SETTINGS lands), then resume queued sends. - let mut window_grew = false; - for (_, item) in self.streams.get().iter() { - // SAFETY: item is &*mut Stream from streams.iter(); the boxed Stream outlives the iteration - let stream = unsafe { &mut **item }; - if (settings.initial_window_size as u64) > stream.remote_window_size { - stream.remote_window_size = settings.initial_window_size as u64; - window_grew = true; + // §6.9.2. remote_window_size is cumulative: a decrease can leave it below the used count. + let delta = settings.initial_window_size as i64 - old_initial_window as i64; + if delta != 0 { + for (_, item) in self.streams.get().iter() { + // SAFETY: item is &*mut Stream from streams.iter(); the boxed Stream outlives the iteration + let stream = unsafe { &mut **item }; + stream.remote_window_size = stream.remote_window_size.saturating_add_signed(delta); } } // Resume queued sends only when a window actually grew; there is nothing to flush otherwise. - if window_grew { + if delta > 0 && !self.streams.get().is_empty() { let _ = self.flush(); } let g = self.global(); @@ -3910,22 +3860,39 @@ impl crate::api::h2::connection::Sink for H2FrameParser { ); } - fn on_window_update(&self, stream_id: u32, increment: u32) { + fn credit_send_window( + &self, + stream_id: u32, + increment: u32, + ) -> crate::api::h2::connection::SendCredit { + use crate::api::h2::connection::SendCredit; bun_output::scoped_log!( H2FrameParser, "engine WU received stream={} inc={}", stream_id, increment ); - // Bridge: the legacy outbound reads its own window cells to decide how much DATA to send. + // `granted` is below `used` after the peer lowered INITIAL_WINDOW_SIZE that far (§6.9.2). + let overflows = + |granted: u64, used: u64| granted + increment as u64 > used + MAX_WINDOW_SIZE as u64; if stream_id == 0 { - self.remote_window_size - .set(self.remote_window_size.get() + increment as u64); + let granted = self.remote_window_size.get(); + if overflows(granted, self.remote_used_window_size.get()) { + return SendCredit::Overflow; + } + self.remote_window_size.set(granted + increment as u64); } else if let Some(stream) = self.streams.get().get(&stream_id).copied() { // SAFETY: stream is *mut Stream from self.streams; valid while the map entry exists - unsafe { (*stream).remote_window_size += increment as u64 }; + let stream = unsafe { &mut *stream }; + if overflows(stream.remote_window_size, stream.remote_used_window_size) { + return SendCredit::Overflow; + } + stream.remote_window_size += increment as u64; + } else { + return SendCredit::NotOwned; } let _ = self.flush(); + SendCredit::Applied } fn on_altsvc(&self, stream_id: u32, origin: &[u8], value: &[u8]) { @@ -5248,7 +5215,6 @@ impl H2FrameParser { stream.remote_used_window_size += payload_size as u64; self.remote_used_window_size .set(self.remote_used_window_size.get() + payload_size as u64); - self.note_engine_send_consumed(stream_id, payload_size as u64); let mut flags: u8 = if end_stream { DataFrameFlags::END_STREAM as u8 } else { @@ -7461,8 +7427,6 @@ impl H2FrameParser { max_header_list_pairs: Cell::new(128), max_settings: Cell::new(32), pending_recv_window_growth: Cell::new(0), - pending_send_window_consumed: Cell::new(0), - pending_stream_send_consumed: JsCell::new(Vec::new()), pending_engine_stream_closes: JsCell::new(Vec::new()), dispatch_depth: Cell::new(0), pending_settings_window_submissions: JsCell::new(Vec::new()), diff --git a/test/js/node/http2/h2-conformance.test.ts b/test/js/node/http2/h2-conformance.test.ts index b121dbb06374..2981d001288c 100644 --- a/test/js/node/http2/h2-conformance.test.ts +++ b/test/js/node/http2/h2-conformance.test.ts @@ -727,8 +727,11 @@ class RawH2Server { } } + send(buf: Buffer) { + this.socket!.write(buf); + } sendFrame(type: number, flags: number, streamId: number, payload?: Buffer) { - this.socket!.write(encodeFrame(type, flags, streamId, payload)); + this.send(encodeFrame(type, flags, streamId, payload)); } waitFor(pred: (f: Frame) => boolean, timeoutMs = 2000): Promise { @@ -929,6 +932,564 @@ describe("SETTINGS ack ordering (RFC 9113 §6.5.3)", () => { }); }); +// The window a side sends DATA against is what the peer granted minus what the side sent. A +// WINDOW_UPDATE is checked against that window (§6.9.1), and a change of the peer's +// SETTINGS_INITIAL_WINDOW_SIZE moves it by the difference (§6.9.2). The side that lowered its own +// value grants window against the new size, or a sender that obeys it waits forever. +describe("flow-control windows after WINDOW_UPDATE and SETTINGS (RFC 9113 §6.9.1, §6.9.2)", () => { + const MAX_WINDOW = 0x7fffffff; + const DEFAULT_WINDOW = 65535; + const STREAM_1_OVERFLOW = [[1, ErrorCode.FLOW_CONTROL_ERROR]]; + const settingsAck = encodeFrame(FrameType.SETTINGS, 0x1, 0); + // :status 200 (static table index 8), END_HEADERS. + const response = encodeFrame(FrameType.HEADERS, 0x4, 1, Buffer.from([0x88])); + + function windowUpdate(streamId: number, increment: number) { + const payload = Buffer.alloc(4); + payload.writeUInt32BE(increment); + return encodeFrame(FrameType.WINDOW_UPDATE, 0, streamId, payload); + } + + function initialWindowSize(value: number) { + const payload = Buffer.alloc(6); + payload.writeUInt16BE(0x4, 0); + payload.writeUInt32BE(value, 2); + return encodeFrame(FrameType.SETTINGS, 0, 0, payload); + } + + type Peer = Pick; + + const dataBytes = (peer: Peer) => + peer.frames.reduce((n, f) => (f.type === FrameType.DATA && f.streamId === 1 ? n + f.length : n), 0); + + let pings = 0; + /** + * Sends `bytes` and a PING. The other side answers frames in order, so its PING ACK comes after + * everything that `bytes` made it send. A GOAWAY ends the wait too: a session that failed sends + * no PING ACK. Returns the errors among the answers, as [stream id, error code]. + */ + async function errorsAfter(peer: Peer, ...bytes: Buffer[]) { + const seen = peer.frames.length; + const payload = Buffer.alloc(8); + payload.writeUInt32BE(++pings, 4); + peer.send(Buffer.concat([...bytes, encodeFrame(FrameType.PING, 0, 0, payload)])); + const ends = (f: Frame) => + f.type === FrameType.GOAWAY || (f.type === FrameType.PING && (f.flags & 0x1) !== 0 && f.payload.equals(payload)); + await peer.waitFor(f => ends(f) && peer.frames.lastIndexOf(f) >= seen); + return peer.frames + .slice(seen) + .filter(f => f.type === FrameType.RST_STREAM || f.type === FrameType.GOAWAY) + .map(f => [f.streamId, f.type === FrameType.GOAWAY ? goawayErrorCode(f) : f.payload.readUInt32BE(0)]); + } + + /** + * A bun client with a POST open on stream 1 of a raw peer that sends no response. The peer + * grants `credit` bytes on the stream, and the client uploads `upload` bytes. `window` is what + * the peer then counts as the send window of the stream. + */ + async function openUpload({ upload = 0, credit = 0, connectionCredit = 16 << 20 } = {}) { + const raw = await RawH2Server.listen(); + const client = http2.connect(`http://127.0.0.1:${raw.port}`); + client.on("error", () => {}); + const req = client.request({ ":method": "POST", ":path": "/" }); + req.on("error", () => {}); + await raw.waitFor(f => f.type === FrameType.HEADERS && f.streamId === 1); + const first = [encodeFrame(FrameType.SETTINGS, 0, 0), settingsAck]; + if (connectionCredit > 0) first.push(windowUpdate(0, connectionCredit)); + if (credit > 0) first.push(windowUpdate(1, credit)); + raw.send(Buffer.concat(first)); + if (upload > 0) { + req.write(Buffer.alloc(upload, 0x61)); + await raw.waitFor(() => dataBytes(raw) >= upload); + } + return { + raw, + req, + window: DEFAULT_WINDOW + credit - upload, + [Symbol.dispose]() { + client.destroy(); + raw.close(); + }, + }; + } + + /** + * A bun server with a response open on stream 1 of a raw client, which grants `credit` bytes on + * that stream. `stream` is the server's side of it. + */ + async function openResponse({ credit = 0, connectionCredit = 16 << 20 } = {}) { + const opened = Promise.withResolvers(); + const server = http2.createServer(); + server.on("stream", stream => { + stream.on("error", () => {}); + stream.respond({ ":status": 200 }); + opened.resolve(stream); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const raw = await RawH2.connect((server.address() as net.AddressInfo).port); + raw.sendPreface(); + raw.sendEmptySettings(); + await raw.waitFor(f => f.type === FrameType.SETTINGS && (f.flags & 0x1) === 0); + raw.send(Buffer.concat([settingsAck, encodeFrame(FrameType.HEADERS, 0x5, 1, requestHeaderBlock("GET"))])); + const stream = await opened.promise; + const credits: Buffer[] = []; + if (connectionCredit > 0) credits.push(windowUpdate(0, connectionCredit)); + if (credit > 0) credits.push(windowUpdate(1, credit)); + await errorsAfter(raw, ...credits); + return { + raw, + stream, + [Symbol.dispose]() { + raw.destroy(); + server.close(); + }, + }; + } + + /** A bun endpoint that sends DATA on stream 1 to a raw peer. */ + const senders = { + async client(options?: { credit?: number }) { + const opened = await openUpload(options); + return { + raw: opened.raw as Peer, + end: (body: Buffer) => void opened.req.end(body), + [Symbol.dispose]: () => opened[Symbol.dispose](), + }; + }, + async server(options?: { credit?: number }) { + const opened = await openResponse(options); + return { + raw: opened.raw as Peer, + end: (body: Buffer) => void opened.stream.end(body), + [Symbol.dispose]: () => opened[Symbol.dispose](), + }; + }, + }; + + // The response HEADERS are the first frame the inbound engine sees for a stream that the client + // opened. A window count that starts there misses the bytes the stream sent before. + test.each([1000, 60000])( + "a grant to exactly 2^31-1 in the same read as the response is accepted after a %d-byte upload", + async upload => { + using opened = await openUpload({ upload }); + const { raw, req, window } = opened; + const responded = once(req, "response"); + const granted = await errorsAfter(raw, response, windowUpdate(1, MAX_WINDOW - window)); + const status = (await responded)[0][":status"]; + expect({ granted, status, oneMore: await errorsAfter(raw, windowUpdate(1, 1)) }).toEqual({ + granted: [], + status: 200, + oneMore: STREAM_1_OVERFLOW, + }); + }, + ); + + test.each([ + ["before the upload", {}], + ["after a 1000-byte upload", { upload: 1000 }], + ["after a 1 MiB upload", { upload: 1 << 20, credit: 1 << 20 }], + ])("the send window of a request with no response yet is capped at 2^31-1 %s", async (_, options) => { + using opened = await openUpload(options); + const { raw, window } = opened; + const granted = await errorsAfter(raw, windowUpdate(1, MAX_WINDOW - window)); + expect({ granted, oneMore: await errorsAfter(raw, windowUpdate(1, 1)) }).toEqual({ + granted: [], + oneMore: STREAM_1_OVERFLOW, + }); + }); + + test("the send window of a response is capped at 2^31-1", async () => { + const sent = 1000; + const server = http2.createServer(); + server.on("stream", stream => { + stream.on("error", () => {}); + stream.respond({ ":status": 200 }); + stream.write(Buffer.alloc(sent, 0x61)); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const c = await RawH2.connect((server.address() as net.AddressInfo).port); + try { + c.sendPreface(); + c.sendEmptySettings(); + await c.waitFor(f => f.type === FrameType.SETTINGS && (f.flags & 0x1) === 0); + c.sendSettingsAck(); + c.sendFrame(FrameType.HEADERS, 0x5 /* END_STREAM | END_HEADERS */, 1, requestHeaderBlock("GET")); + await c.waitFor(() => dataBytes(c) >= sent); + const granted = await errorsAfter(c, windowUpdate(1, MAX_WINDOW - (DEFAULT_WINDOW - sent))); + expect({ granted, oneMore: await errorsAfter(c, windowUpdate(1, 1)) }).toEqual({ + granted: [], + oneMore: STREAM_1_OVERFLOW, + }); + } finally { + c.destroy(); + server.close(); + } + }); + + test("the connection send window is capped at 2^31-1", async () => { + const upload = 1000; + using opened = await openUpload({ upload, connectionCredit: 0 }); + const { raw } = opened; + expect(await errorsAfter(raw, windowUpdate(0, MAX_WINDOW - (DEFAULT_WINDOW - upload)))).toEqual([]); + raw.send(windowUpdate(0, 1)); + expect(goawayErrorCode(await raw.waitFor(f => f.type === FrameType.GOAWAY))).toBe(ErrorCode.FLOW_CONTROL_ERROR); + }); + + test("a grant to exactly 2^31-1 is accepted after the peer lowered INITIAL_WINDOW_SIZE", async () => { + using opened = await openUpload({ upload: 1000 }); + const { raw, window } = opened; + // The change moves the window by 1 - 65535, to 999 bytes below zero. + const lowered = window + 1 - DEFAULT_WINDOW; + const granted = await errorsAfter( + raw, + initialWindowSize(1), + windowUpdate(1, MAX_WINDOW), + windowUpdate(1, -lowered), + ); + expect({ lowered, granted, oneMore: await errorsAfter(raw, windowUpdate(1, 1)) }).toEqual({ + lowered: -999, + granted: [], + oneMore: STREAM_1_OVERFLOW, + }); + }); + + // The client uploads to a raw server. The server responds to a raw client. + describe.each(["client", "server"] as const)("a bun %s that sends DATA", role => { + test("sends none past a window that the peer lowered", async () => { + using sender = await senders[role](); + const { raw } = sender; + await errorsAfter(raw, initialWindowSize(0)); + sender.end(Buffer.alloc(100_000, 0x61)); + await errorsAfter(raw); + const sentIntoNoWindow = dataBytes(raw); + await errorsAfter(raw, windowUpdate(1, 100)); + expect({ sentIntoNoWindow, sentInto100: dataBytes(raw) }).toEqual({ sentIntoNoWindow: 0, sentInto100: 100 }); + }); + + test("keeps its WINDOW_UPDATE credit when the peer raises INITIAL_WINDOW_SIZE", async () => { + const credit = 10_000; + const raised = 100_000; + using sender = await senders[role]({ credit }); + const { raw } = sender; + sender.end(Buffer.alloc(256 * 1024, 0x61)); + await raw.waitFor(() => dataBytes(raw) >= DEFAULT_WINDOW + credit); + await errorsAfter(raw); + const sentBefore = dataBytes(raw); + await errorsAfter(raw, initialWindowSize(raised)); + expect({ sentBefore, sentAfter: dataBytes(raw) }).toEqual({ + sentBefore: DEFAULT_WINDOW + credit, + sentAfter: raised + credit, + }); + }); + + // The second change starts from 100000, not from 65535. + test("moves its window by the difference at every change of INITIAL_WINDOW_SIZE", async () => { + using sender = await senders[role](); + const { raw } = sender; + const sentAfter = async (...bytes: Buffer[]) => { + await errorsAfter(raw, ...bytes); + return dataBytes(raw); + }; + sender.end(Buffer.alloc(256 * 1024, 0x61)); + await raw.waitFor(() => dataBytes(raw) >= DEFAULT_WINDOW); + expect({ + initial: await sentAfter(), + raisedTo100000: await sentAfter(initialWindowSize(100_000)), + loweredTo80000: await sentAfter(initialWindowSize(80_000)), + granted30000: await sentAfter(windowUpdate(1, 30_000)), + raisedTo120000: await sentAfter(initialWindowSize(120_000)), + }).toEqual({ + initial: DEFAULT_WINDOW, + raisedTo100000: 100_000, + loweredTo80000: 100_000, + granted30000: 110_000, + raisedTo120000: 150_000, + }); + }); + }); + + /** + * A server that lowers its initialWindowSize to `lowered` at the first DATA of stream 1, and + * then calls `andThen`. + */ + async function serverThatLowersItsWindow( + lowered: number, + { + onEnd = () => {}, + andThen = () => {}, + }: { onEnd?: (bytes: number, id: number) => void; andThen?: (session: http2.Http2Session) => void } = {}, + ) { + const server = http2.createServer(); + server.on("stream", stream => { + let bytes = 0; + const session = stream.session!; + stream.on("error", () => {}); + if (stream.id === 1) { + stream.once("data", () => { + session.settings({ initialWindowSize: lowered }); + andThen(session); + }); + } + stream.on("data", (chunk: Buffer) => (bytes += chunk.length)); + stream.on("end", () => { + onEnd(bytes, stream.id!); + if (stream.destroyed) return; + stream.respond({ ":status": 200 }); + stream.end(); + }); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + return server; + } + + /** Whether `f` is a SETTINGS frame that sets INITIAL_WINDOW_SIZE to `value`. */ + function setsInitialWindowSize(f: Frame, value: number) { + if (f.type !== FrameType.SETTINGS || (f.flags & 0x1) !== 0) return false; + for (let entry = 0; entry + 6 <= f.payload.length; entry += 6) { + if (f.payload.readUInt16BE(entry) === 0x4 && f.payload.readUInt32BE(entry + 2) === value) return true; + } + return false; + } + + // setLocalWindowSize() raises the initialWindowSize that the session keeps, and sends no SETTINGS. + test.each([ + ["", (_: http2.Http2Session) => {}], + [" and then calls setLocalWindowSize()", (session: http2.Http2Session) => session.setLocalWindowSize(1 << 20)], + ])("a server that lowers initialWindowSize%s grants window against the new size", async (_, andThen) => { + const lowered = 1024; + const server = await serverThatLowersItsWindow(lowered, { andThen }); + const c = await RawH2.connect((server.address() as net.AddressInfo).port); + try { + c.sendPreface(); + c.sendEmptySettings(); + await c.waitFor(f => f.type === FrameType.SETTINGS && (f.flags & 0x1) === 0); + c.sendSettingsAck(); + c.sendFrame(FrameType.HEADERS, 0x4 /* END_HEADERS */, 1, requestHeaderBlock("POST")); + c.sendFrame(FrameType.DATA, 0, 1, Buffer.alloc(1000, 0x61)); + await c.waitFor(f => setsInitialWindowSize(f, lowered)); + // The sender now has 65535 - 1000 + (1024 - 65535) = 24 bytes of window. Half of the old + // size is more than it can ever send, so the 1000 bytes come back when the ACK arrives. + c.sendSettingsAck(); + const update = await c.waitFor(f => f.type === FrameType.WINDOW_UPDATE && f.streamId === 1); + expect(update.payload.readUInt32BE(0)).toBe(1000); + } finally { + c.destroy(); + server.close(); + } + }); + + test("a client that lowers initialWindowSize grants window against the new size", async () => { + const lowered = 1024; + const raw = await RawH2Server.listen(); + const client = http2.connect(`http://127.0.0.1:${raw.port}`); + client.on("error", () => {}); + try { + const req = client.request({ ":path": "/" }); + req.on("error", () => {}); + req.once("data", () => client.settings({ initialWindowSize: lowered })); + await raw.waitFor(f => f.type === FrameType.HEADERS && f.streamId === 1); + raw.send( + Buffer.concat([ + encodeFrame(FrameType.SETTINGS, 0, 0), + settingsAck, + response, + encodeFrame(FrameType.DATA, 0, 1, Buffer.alloc(1000, 0x61)), + ]), + ); + await raw.waitFor(f => setsInitialWindowSize(f, lowered)); + raw.send(settingsAck); + const update = await raw.waitFor(f => f.type === FrameType.WINDOW_UPDATE && f.streamId === 1); + expect(update.payload.readUInt32BE(0)).toBe(1000); + } finally { + client.destroy(); + raw.close(); + } + }); + + // The peer still sends against 65535 per stream: setLocalWindowSize() is about the connection. + test("a stream gets its window back after setLocalWindowSize() raised the connection window", async () => { + const raw = await RawH2Server.listen(); + const client = http2.connect(`http://127.0.0.1:${raw.port}`); + client.on("error", () => {}); + try { + await once(client, "connect"); + client.setLocalWindowSize(1 << 20); + const req = client.request({ ":path": "/" }); + req.on("error", () => {}); + req.resume(); + await raw.waitFor(f => f.type === FrameType.HEADERS && f.streamId === 1); + const half = encodeFrame(FrameType.DATA, 0, 1, Buffer.alloc(16384, 0x61)); + raw.send(Buffer.concat([encodeFrame(FrameType.SETTINGS, 0, 0), settingsAck, response, half, half])); + const update = await raw.waitFor(f => f.type === FrameType.WINDOW_UPDATE && f.streamId === 1); + expect(update.payload.readUInt32BE(0)).toBe(2 * 16384); + } finally { + client.destroy(); + raw.close(); + } + }); + + /** + * A bun server stream that the handler does not read, and the raw client that uploads on it. + * The client fills the readable buffer to one byte below its high-water mark (16 KiB on Windows, + * 64 KiB elsewhere), so the `more` bytes behind that pause the stream. `sent` is what the client + * sent, and `counted` is the part that the server has not returned as window. `onSession` runs + * before the first stream. + */ + async function pausedUpload(more: number, onSession: (session: http2.Http2Session) => void = () => {}) { + const opened = Promise.withResolvers(); + const server = http2.createServer(); + server.on("session", onSession); + server.on("stream", stream => { + stream.on("error", () => {}); + opened.resolve(stream); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const raw = await RawH2.connect((server.address() as net.AddressInfo).port); + raw.sendPreface(); + raw.sendEmptySettings(); + await raw.waitFor(f => f.type === FrameType.SETTINGS && (f.flags & 0x1) === 0); + raw.send(Buffer.concat([settingsAck, encodeFrame(FrameType.HEADERS, 0x4, 1, requestHeaderBlock("POST"))])); + const stream = await opened.promise; + const data = (bytes: number) => encodeFrame(FrameType.DATA, 0, 1, Buffer.alloc(bytes, 0x61)); + const fill = stream.readableHighWaterMark - 1; + const filler: Buffer[] = []; + for (let left = fill; left > 0; left -= 16384) filler.push(data(Math.min(left, 16384))); + await errorsAfter(raw, ...filler); + await errorsAfter(raw, data(more)); + await errorsAfter(raw); + const returned = raw.frames.reduce( + (n, f) => (f.type === FrameType.WINDOW_UPDATE && f.streamId === 1 ? n + f.payload.readUInt32BE(0) : n), + 0, + ); + return { + raw, + stream, + sent: fill + more, + counted: fill + more - returned, + [Symbol.dispose]() { + raw.destroy(); + server.close(); + }, + }; + } + + // A resume is the only place that returns window to a stream that was paused. + test("a paused stream gets its bytes back at resume, against a lowered initialWindowSize", async () => { + using paused = await pausedUpload(1000); + const { raw, stream, counted } = paused; + stream.session!.settings({ initialWindowSize: 1024 }); + await raw.waitFor(f => setsInitialWindowSize(f, 1024)); + await errorsAfter(raw, settingsAck); + await errorsAfter(raw); + const seen = raw.frames.length; + stream.resume(); + const update = await raw.waitFor( + f => f.type === FrameType.WINDOW_UPDATE && f.streamId === 1 && raw.frames.lastIndexOf(f) >= seen, + ); + expect(update.payload.readUInt32BE(0)).toBe(counted); + }); + + // The stream opens after the settings() call, so its window is 1 from the start. Its bytes are + // in flight until the client sends the ACK, and a paused stream returns none of them. + test("an empty END_STREAM frame is accepted on a paused stream that holds more than a lowered initialWindowSize", async () => { + using paused = await pausedUpload(1000, session => session.settings({ initialWindowSize: 1 })); + const { raw, stream, sent, counted } = paused; + const errors = [ + ...(await errorsAfter(raw, settingsAck)), + ...(await errorsAfter(raw, encodeFrame(FrameType.DATA, 0x1 /* END_STREAM */, 1))), + ]; + expect({ counted, errors }).toEqual({ counted: 1000, errors: [] }); + let bytes = 0; + const ended = Promise.withResolvers(); + stream.on("data", (chunk: Buffer) => (bytes += chunk.length)); + stream.on("end", () => ended.resolve(bytes)); + expect(await ended.promise).toBe(sent); + }); + + // The client cannot know of the lower value when it sends the 2000 bytes of stream 3. Its ACK + // makes that value the limit, and the empty frame behind it takes no window (§6.9.1). + test("an empty END_STREAM frame is accepted when the bytes in flight exceed a lowered initialWindowSize", async () => { + const ended = Promise.withResolvers(); + const server = await serverThatLowersItsWindow(1, { onEnd: (bytes, id) => id === 3 && ended.resolve(bytes) }); + const c = await RawH2.connect((server.address() as net.AddressInfo).port); + try { + c.sendPreface(); + c.sendEmptySettings(); + await c.waitFor(f => f.type === FrameType.SETTINGS && (f.flags & 0x1) === 0); + c.sendSettingsAck(); + c.sendFrame(FrameType.HEADERS, 0x4 /* END_HEADERS */, 1, requestHeaderBlock("POST")); + c.sendFrame(FrameType.DATA, 0, 1, Buffer.alloc(100, 0x61)); + await c.waitFor(f => setsInitialWindowSize(f, 1)); + c.send( + Buffer.concat([ + encodeFrame(FrameType.HEADERS, 0x4 /* END_HEADERS */, 3, requestHeaderBlock("POST")), + encodeFrame(FrameType.DATA, 0, 3, Buffer.alloc(2000, 0x62)), + settingsAck, + encodeFrame(FrameType.DATA, 0x1 /* END_STREAM */, 3), + ]), + ); + const answer = await c.waitFor( + f => f.type === FrameType.GOAWAY || (f.type === FrameType.HEADERS && f.streamId === 3), + ); + expect(answer.type === FrameType.GOAWAY ? goawayErrorCode(answer) : "response").toBe("response"); + expect(await ended.promise).toBe(2000); + } finally { + c.destroy(); + server.close(); + } + }); + + test("an upload finishes when the server lowers initialWindowSize while it runs", async () => { + const total = 128 * 1024; + const received = Promise.withResolvers(); + const server = await serverThatLowersItsWindow(1024, { onEnd: received.resolve }); + const client = http2.connect(`http://127.0.0.1:${(server.address() as net.AddressInfo).port}`); + client.on("error", received.reject); + try { + const req = client.request({ ":method": "POST", ":path": "/" }); + req.on("error", received.reject); + req.resume(); + req.end(Buffer.alloc(total, 0x61)); + expect(await received.promise).toBe(total); + } finally { + client.destroy(); + server.close(); + } + }); + + test("a download finishes when the client lowers initialWindowSize while it runs", async () => { + const total = 128 * 1024; + const server = http2.createServer(); + server.on("stream", stream => { + stream.on("error", () => {}); + stream.respond({ ":status": 200 }); + stream.end(Buffer.alloc(total, 0x61)); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const received = Promise.withResolvers(); + const client = http2.connect(`http://127.0.0.1:${(server.address() as net.AddressInfo).port}`); + client.on("error", received.reject); + try { + let bytes = 0; + const req = client.request({ ":path": "/" }); + req.on("error", received.reject); + req.once("data", () => client.settings({ initialWindowSize: 1024 })); + req.on("data", (chunk: Buffer) => (bytes += chunk.length)); + req.on("end", () => received.resolve(bytes)); + req.end(); + expect(await received.promise).toBe(total); + } finally { + client.destroy(); + server.close(); + } + }); +}); + function requestHeaderBlock(method: "GET" | "POST", extra: Buffer = Buffer.alloc(0)): Buffer { return Buffer.concat([ Buffer.from([method === "POST" ? 0x83 : 0x82, 0x86, 0x84, 0x01]),