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
99 changes: 67 additions & 32 deletions src/runtime/api/bun/h2/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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)
Comment thread
robobun marked this conversation as resolved.
}

/// 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() {
Expand All @@ -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();
Expand Down Expand Up @@ -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());
Comment thread
robobun marked this conversation as resolved.
false
}

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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);
}
Expand Down
13 changes: 11 additions & 2 deletions src/runtime/api/bun/h2/flow_control.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand All @@ -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 {
Expand Down Expand Up @@ -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);
Expand Down
Loading
Loading