From c271f6873293c010bd05b070f067c83ba8fe17e1 Mon Sep 17 00:00:00 2001 From: Jeremy Justus Date: Wed, 30 Sep 2026 00:45:22 -0400 Subject: [PATCH 1/2] fix(http2): wait for useful send capacity before handing a chunk to h2 A body chunk was handed to h2 once the stream held any capacity. h2 cuts DATA frames from the capacity a stream holds, so on a connection whose window was nearly spent a chunk left as a 1-byte frame followed by more slivers. h2 0.4.16+ servers charge small non-final DATA frames to a per-connection budget and GOAWAY the connection with ENHANCE_YOUR_CALM when it runs out. Reserve min(len, 1024) instead of 1 and wait until that much is assigned. The claim stays small, so #4003 still holds, but no sub-1 KiB first frame is cut from a larger chunk. Closes #4211 --- src/proto/h2/mod.rs | 33 +++++++++-- tests/client.rs | 133 ++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 161 insertions(+), 5 deletions(-) diff --git a/src/proto/h2/mod.rs b/src/proto/h2/mod.rs index 393d4179ab..909c918db9 100644 --- a/src/proto/h2/mod.rs +++ b/src/proto/h2/mod.rs @@ -113,6 +113,27 @@ struct Peeked { is_eos: bool, } +/// Smallest amount of assigned stream capacity worth handing a chunk to h2 +/// with. +/// +/// h2 cuts a DATA frame from whatever capacity a stream holds when the frame +/// is written. With only the minimal claim assigned, a chunk sent on a +/// connection whose window is nearly spent leaves as a run of tiny frames +/// (silly-window syndrome): a 1-byte frame, then more slivers as capacity +/// trickles in. Each costs a full frame header, and h2 0.4.16+ servers charge +/// every non-final DATA frame under 256 bytes to a small per-connection budget +/// and answer its exhaustion with `GOAWAY(ENHANCE_YOUR_CALM)`, failing every +/// stream on the connection. +/// +/// Waiting for this much keeps the claim small (the reasoning above still +/// holds) while ruling out sliver frames. 1 KiB is far below any stream window +/// a real peer advertises, so it cannot hold a chunk back indefinitely. +const MIN_DATA_FRAME_CAPACITY: usize = 1024; + +fn min_data_frame_capacity(len: usize) -> usize { + len.min(MIN_DATA_FRAME_CAPACITY) +} + impl PipeToSendStream where S: Body, @@ -157,11 +178,13 @@ where // If a previously-polled chunk is still waiting for stream-level // send capacity, drive that to completion before touching the // body again. - if me.buffered_data.is_some() { - while me.body_tx.capacity() == 0 { + if let Some(peeked) = me.buffered_data.as_ref() { + // Wait for enough capacity to cut a useful first frame, not + // just any capacity: see `MIN_DATA_FRAME_CAPACITY`. + let needed = min_data_frame_capacity(peeked.data.remaining()); + while me.body_tx.capacity() < needed { match ready!(me.body_tx.poll_capacity(cx)) { - Some(Ok(0)) => {} - Some(Ok(_)) => break, + Some(Ok(_)) => {} Some(Err(e)) => return Poll::Ready(Err(crate::Error::new_body_write(e))), None => { // None means the stream is no longer in a @@ -230,7 +253,7 @@ where // chunk in `self` so it survives the upcoming // `poll_capacity` wait even if it returns // `Poll::Pending`. - me.body_tx.reserve_capacity(1); + me.body_tx.reserve_capacity(min_data_frame_capacity(len)); *me.buffered_data = Some(Peeked { data: chunk, is_eos, diff --git a/tests/client.rs b/tests/client.rs index e4c8d12cf6..fbf2cd1f3d 100644 --- a/tests/client.rs +++ b/tests/client.rs @@ -3345,6 +3345,139 @@ mod conn { let _ = tokio::time::timeout(Duration::from_secs(5), a_handle).await; } + // A chunk must not be handed to h2 while the stream holds only a sliver + // of capacity. h2 cuts a DATA frame from whatever capacity a stream holds, + // so on a connection whose window is nearly spent a chunk would otherwise + // leave as a 1-byte frame followed by more slivers. h2 0.4.16+ servers + // charge every small non-final DATA frame to a per-connection budget and + // answer its exhaustion with `GOAWAY(ENHANCE_YOUR_CALM)`. + #[tokio::test] + async fn h2_chunk_waits_for_useful_capacity_instead_of_sliver_frames() { + use std::sync::{Arc, Mutex}; + + // Leave exactly one byte of the 65535-byte initial connection window. + const STREAM_A_LEN: usize = 65534; + const STREAM_B_LEN: usize = 10_000; + + let (client_io, server_io, _) = setup_duplex_test_server(); + let (stream_a_full_tx, stream_a_full_rx) = oneshot::channel::<()>(); + let (release_a_tx, release_a_rx) = oneshot::channel::<()>(); + let (stream_b_first_tx, stream_b_first_rx) = oneshot::channel::(); + + let stream_a_full_tx = Arc::new(Mutex::new(Some(stream_a_full_tx))); + let release_a_rx = Arc::new(Mutex::new(Some(release_a_rx))); + let stream_b_first_tx = Arc::new(Mutex::new(Some(stream_b_first_tx))); + tokio::spawn(async move { + let mut h2 = h2::server::handshake(server_io).await.unwrap(); + let mut seen = 0u32; + while let Some(result) = h2.accept().await { + let (req, mut respond) = result.unwrap(); + seen += 1; + let which = seen; + let stream_a_full_tx = stream_a_full_tx.clone(); + let release_a_rx = release_a_rx.clone(); + let stream_b_first_tx = stream_b_first_tx.clone(); + tokio::spawn(async move { + let mut body = req.into_body(); + if which == 1 { + // Stream A: take the burst without releasing capacity + // until the test says so, then release all of it. + let mut received = 0usize; + while received < STREAM_A_LEN { + match body.data().await { + Some(Ok(f)) => received += f.len(), + _ => return, + } + } + if let Some(tx) = stream_a_full_tx.lock().unwrap().take() { + let _ = tx.send(()); + } + let release = release_a_rx.lock().unwrap().take(); + if let Some(release) = release { + let _ = release.await; + } + let _ = body.flow_control().release_capacity(received); + // Hold the stream open; the test ends before it matters. + std::future::pending::<()>().await; + } else { + // Stream B: record the size of its FIRST data frame. + let first = match body.data().await { + Some(Ok(f)) => { + let len = f.len(); + let _ = body.flow_control().release_capacity(len); + len + } + _ => 0, + }; + if let Some(tx) = stream_b_first_tx.lock().unwrap().take() { + let _ = tx.send(first); + } + while let Some(Ok(f)) = body.data().await { + let _ = body.flow_control().release_capacity(f.len()); + } + let mut send = respond.send_response(Response::new(()), false).unwrap(); + let _ = send.send_data(Bytes::from_static(b"ok"), true); + } + }); + } + }); + + let io = TokioIo::new(client_io); + let (mut client, conn) = conn::http2::Builder::new(TokioExecutor) + .handshake::<_, BoxBody>>(io) + .await + .expect("http handshake"); + tokio::spawn(async move { + let _ = conn.await; + }); + + // Request A fills the connection window down to its last byte. + let (mut tx_a, rx_a) = + mpsc::channel::, Box>>(4); + let body_a: BoxBody> = + BodyExt::boxed(StreamBody::new(rx_a)); + let req_a = Request::post("http://localhost/a").body(body_a).unwrap(); + let mut client_a = client.clone(); + let _a_handle = tokio::spawn(async move { client_a.send_request(req_a).await }); + use futures_util::SinkExt; + let mut remaining = STREAM_A_LEN; + while remaining > 0 { + let take = remaining.min(16_384); + tx_a.send(Ok(Frame::data(Bytes::from(vec![b'A'; take])))) + .await + .expect("stream A channel send"); + remaining -= take; + } + tokio::time::timeout(Duration::from_secs(5), stream_a_full_rx) + .await + .expect("server should receive full stream A body in time") + .expect("stream_a_full_rx"); + + // Request B: a 10 KB chunk while only one byte of window is left. + let body_b: BoxBody> = BodyExt::boxed( + http_body_util::Full::new(Bytes::from(vec![b'B'; STREAM_B_LEN])) + .map_err(|never: std::convert::Infallible| match never {}), + ); + let req_b = Request::post("http://localhost/b").body(body_b).unwrap(); + let b_fut = tokio::spawn(client.send_request(req_b)); + + // Let stream B's pipe poll its chunk against the one-byte window, then + // free the window by releasing stream A's capacity. + tokio::time::sleep(Duration::from_millis(50)).await; + let _ = release_a_tx.send(()); + + let first_b = tokio::time::timeout(Duration::from_secs(5), stream_b_first_rx) + .await + .expect("stream B must reach the server once the window is released") + .expect("stream_b_first_rx"); + assert!( + first_b >= 1024, + "stream B's first DATA frame was {first_b} bytes: the chunk was cut from a sliver of capacity" + ); + let _ = tokio::time::timeout(Duration::from_secs(5), b_fut).await; + drop(tx_a); + } + // https://github.com/hyperium/hyper/issues/4003, for HTTP/2 CONNECT // // Like `h2_idle_stream_does_not_pin_connection_window`, but the idle From 21ff6ea30f1e017062b7db3aa392fc1f7cb86302 Mon Sep 17 00:00:00 2001 From: Jeremy Justus Date: Sun, 4 Oct 2026 04:15:49 -0400 Subject: [PATCH 2/2] test(http2): cover small flow-control windows Remove the fixed 1 KiB body-capacity gate because peers may legally advertise a smaller stream window. Add regressions showing that a 512-byte peer window and the final byte of connection capacity both make request-body progress. --- src/proto/h2/mod.rs | 33 ++----- tests/client.rs | 204 +++++++++++++++++++++----------------------- 2 files changed, 101 insertions(+), 136 deletions(-) diff --git a/src/proto/h2/mod.rs b/src/proto/h2/mod.rs index 909c918db9..393d4179ab 100644 --- a/src/proto/h2/mod.rs +++ b/src/proto/h2/mod.rs @@ -113,27 +113,6 @@ struct Peeked { is_eos: bool, } -/// Smallest amount of assigned stream capacity worth handing a chunk to h2 -/// with. -/// -/// h2 cuts a DATA frame from whatever capacity a stream holds when the frame -/// is written. With only the minimal claim assigned, a chunk sent on a -/// connection whose window is nearly spent leaves as a run of tiny frames -/// (silly-window syndrome): a 1-byte frame, then more slivers as capacity -/// trickles in. Each costs a full frame header, and h2 0.4.16+ servers charge -/// every non-final DATA frame under 256 bytes to a small per-connection budget -/// and answer its exhaustion with `GOAWAY(ENHANCE_YOUR_CALM)`, failing every -/// stream on the connection. -/// -/// Waiting for this much keeps the claim small (the reasoning above still -/// holds) while ruling out sliver frames. 1 KiB is far below any stream window -/// a real peer advertises, so it cannot hold a chunk back indefinitely. -const MIN_DATA_FRAME_CAPACITY: usize = 1024; - -fn min_data_frame_capacity(len: usize) -> usize { - len.min(MIN_DATA_FRAME_CAPACITY) -} - impl PipeToSendStream where S: Body, @@ -178,13 +157,11 @@ where // If a previously-polled chunk is still waiting for stream-level // send capacity, drive that to completion before touching the // body again. - if let Some(peeked) = me.buffered_data.as_ref() { - // Wait for enough capacity to cut a useful first frame, not - // just any capacity: see `MIN_DATA_FRAME_CAPACITY`. - let needed = min_data_frame_capacity(peeked.data.remaining()); - while me.body_tx.capacity() < needed { + if me.buffered_data.is_some() { + while me.body_tx.capacity() == 0 { match ready!(me.body_tx.poll_capacity(cx)) { - Some(Ok(_)) => {} + Some(Ok(0)) => {} + Some(Ok(_)) => break, Some(Err(e)) => return Poll::Ready(Err(crate::Error::new_body_write(e))), None => { // None means the stream is no longer in a @@ -253,7 +230,7 @@ where // chunk in `self` so it survives the upcoming // `poll_capacity` wait even if it returns // `Poll::Pending`. - me.body_tx.reserve_capacity(min_data_frame_capacity(len)); + me.body_tx.reserve_capacity(1); *me.buffered_data = Some(Peeked { data: chunk, is_eos, diff --git a/tests/client.rs b/tests/client.rs index fbf2cd1f3d..73ab121ff0 100644 --- a/tests/client.rs +++ b/tests/client.rs @@ -3345,137 +3345,125 @@ mod conn { let _ = tokio::time::timeout(Duration::from_secs(5), a_handle).await; } - // A chunk must not be handed to h2 while the stream holds only a sliver - // of capacity. h2 cuts a DATA frame from whatever capacity a stream holds, - // so on a connection whose window is nearly spent a chunk would otherwise - // leave as a 1-byte frame followed by more slivers. h2 0.4.16+ servers - // charge every small non-final DATA frame to a per-connection budget and - // answer its exhaustion with `GOAWAY(ENHANCE_YOUR_CALM)`. #[tokio::test] - async fn h2_chunk_waits_for_useful_capacity_instead_of_sliver_frames() { - use std::sync::{Arc, Mutex}; + async fn h2_body_progresses_with_peer_window_below_one_kibibyte() { + let (client_io, server_io, _) = setup_duplex_test_server(); + let (headers_tx, headers_rx) = oneshot::channel(); + let (first_tx, first_rx) = oneshot::channel(); + tokio::spawn(async move { + let mut builder = h2::server::Builder::new(); + builder.initial_window_size(512); + let mut server = builder.handshake::<_, Bytes>(server_io).await.unwrap(); + let (request, _respond) = server.accept().await.unwrap().unwrap(); + tokio::spawn(async move { + let _ = poll_fn(|cx| server.poll_closed(cx)).await; + }); + headers_tx.send(()).unwrap(); + let mut body = request.into_body(); + let first = body.data().await.unwrap().unwrap(); + first_tx.send(first.len()).unwrap(); + }); + + let io = TokioIo::new(client_io); + let (mut client, connection) = conn::http2::Builder::new(TokioExecutor) + .handshake::<_, BoxBody>>(io) + .await + .expect("http handshake"); + tokio::spawn(async move { + let _ = connection.await; + }); + + let (mut tx, rx) = mpsc::channel::, Box>>(1); + let request = Request::post("http://localhost/small-window") + .body(BodyExt::boxed(StreamBody::new(rx))) + .unwrap(); + tokio::spawn(client.send_request(request)); + headers_rx.await.unwrap(); + + use futures_util::SinkExt; + tx.send(Ok(Frame::data(Bytes::from(vec![0; 2048])))) + .await + .unwrap(); + let first = tokio::time::timeout(Duration::from_secs(5), first_rx) + .await + .expect("a legal 512-byte stream window must permit DATA progress") + .unwrap(); + assert_eq!(first, 512); + } - // Leave exactly one byte of the 65535-byte initial connection window. + #[tokio::test] + async fn h2_body_uses_the_last_byte_of_connection_capacity() { const STREAM_A_LEN: usize = 65534; - const STREAM_B_LEN: usize = 10_000; let (client_io, server_io, _) = setup_duplex_test_server(); - let (stream_a_full_tx, stream_a_full_rx) = oneshot::channel::<()>(); - let (release_a_tx, release_a_rx) = oneshot::channel::<()>(); - let (stream_b_first_tx, stream_b_first_rx) = oneshot::channel::(); - - let stream_a_full_tx = Arc::new(Mutex::new(Some(stream_a_full_tx))); - let release_a_rx = Arc::new(Mutex::new(Some(release_a_rx))); - let stream_b_first_tx = Arc::new(Mutex::new(Some(stream_b_first_tx))); + let (a_full_tx, a_full_rx) = oneshot::channel(); + let (b_first_tx, b_first_rx) = oneshot::channel(); tokio::spawn(async move { - let mut h2 = h2::server::handshake(server_io).await.unwrap(); - let mut seen = 0u32; - while let Some(result) = h2.accept().await { - let (req, mut respond) = result.unwrap(); - seen += 1; - let which = seen; - let stream_a_full_tx = stream_a_full_tx.clone(); - let release_a_rx = release_a_rx.clone(); - let stream_b_first_tx = stream_b_first_tx.clone(); - tokio::spawn(async move { - let mut body = req.into_body(); - if which == 1 { - // Stream A: take the burst without releasing capacity - // until the test says so, then release all of it. - let mut received = 0usize; - while received < STREAM_A_LEN { - match body.data().await { - Some(Ok(f)) => received += f.len(), - _ => return, - } - } - if let Some(tx) = stream_a_full_tx.lock().unwrap().take() { - let _ = tx.send(()); - } - let release = release_a_rx.lock().unwrap().take(); - if let Some(release) = release { - let _ = release.await; - } - let _ = body.flow_control().release_capacity(received); - // Hold the stream open; the test ends before it matters. - std::future::pending::<()>().await; - } else { - // Stream B: record the size of its FIRST data frame. - let first = match body.data().await { - Some(Ok(f)) => { - let len = f.len(); - let _ = body.flow_control().release_capacity(len); - len - } - _ => 0, - }; - if let Some(tx) = stream_b_first_tx.lock().unwrap().take() { - let _ = tx.send(first); - } - while let Some(Ok(f)) = body.data().await { - let _ = body.flow_control().release_capacity(f.len()); - } - let mut send = respond.send_response(Response::new(()), false).unwrap(); - let _ = send.send_data(Bytes::from_static(b"ok"), true); - } - }); - } + let mut server = h2::server::Builder::new() + .handshake::<_, Bytes>(server_io) + .await + .unwrap(); + let (a, respond_a) = server.accept().await.unwrap().unwrap(); + tokio::spawn(async move { + let _respond_a = respond_a; + let mut body = a.into_body(); + let mut received = 0; + while received < STREAM_A_LEN { + received += body.data().await.unwrap().unwrap().len(); + } + a_full_tx.send(()).unwrap(); + std::future::pending::<()>().await; + }); + let (b, _respond_b) = server.accept().await.unwrap().unwrap(); + tokio::spawn(async move { + let _ = poll_fn(|cx| server.poll_closed(cx)).await; + }); + let mut body = b.into_body(); + let first = body.data().await.unwrap().unwrap(); + b_first_tx.send(first.len()).unwrap(); }); let io = TokioIo::new(client_io); - let (mut client, conn) = conn::http2::Builder::new(TokioExecutor) + let (mut client, connection) = conn::http2::Builder::new(TokioExecutor) .handshake::<_, BoxBody>>(io) .await .expect("http handshake"); tokio::spawn(async move { - let _ = conn.await; + let _ = connection.await; }); - // Request A fills the connection window down to its last byte. - let (mut tx_a, rx_a) = + let (mut a_tx, a_rx) = mpsc::channel::, Box>>(4); - let body_a: BoxBody> = - BodyExt::boxed(StreamBody::new(rx_a)); - let req_a = Request::post("http://localhost/a").body(body_a).unwrap(); let mut client_a = client.clone(); - let _a_handle = tokio::spawn(async move { client_a.send_request(req_a).await }); + tokio::spawn(async move { + let request = Request::post("http://localhost/a") + .body(BodyExt::boxed(StreamBody::new(a_rx))) + .unwrap(); + let _ = client_a.send_request(request).await; + }); use futures_util::SinkExt; let mut remaining = STREAM_A_LEN; - while remaining > 0 { - let take = remaining.min(16_384); - tx_a.send(Ok(Frame::data(Bytes::from(vec![b'A'; take])))) + while remaining != 0 { + let len = remaining.min(16_384); + a_tx.send(Ok(Frame::data(Bytes::from(vec![0; len])))) .await - .expect("stream A channel send"); - remaining -= take; + .unwrap(); + remaining -= len; } - tokio::time::timeout(Duration::from_secs(5), stream_a_full_rx) - .await - .expect("server should receive full stream A body in time") - .expect("stream_a_full_rx"); + a_full_rx.await.unwrap(); - // Request B: a 10 KB chunk while only one byte of window is left. - let body_b: BoxBody> = BodyExt::boxed( - http_body_util::Full::new(Bytes::from(vec![b'B'; STREAM_B_LEN])) - .map_err(|never: std::convert::Infallible| match never {}), - ); - let req_b = Request::post("http://localhost/b").body(body_b).unwrap(); - let b_fut = tokio::spawn(client.send_request(req_b)); - - // Let stream B's pipe poll its chunk against the one-byte window, then - // free the window by releasing stream A's capacity. - tokio::time::sleep(Duration::from_millis(50)).await; - let _ = release_a_tx.send(()); - - let first_b = tokio::time::timeout(Duration::from_secs(5), stream_b_first_rx) + let request = Request::post("http://localhost/b") + .body(BodyExt::boxed( + Full::new(Bytes::from(vec![0; 10_000])) + .map_err(|never: std::convert::Infallible| match never {}), + )) + .unwrap(); + tokio::spawn(client.send_request(request)); + let first = tokio::time::timeout(Duration::from_secs(5), b_first_rx) .await - .expect("stream B must reach the server once the window is released") - .expect("stream_b_first_rx"); - assert!( - first_b >= 1024, - "stream B's first DATA frame was {first_b} bytes: the chunk was cut from a sliver of capacity" - ); - let _ = tokio::time::timeout(Duration::from_secs(5), b_fut).await; - drop(tx_a); + .expect("the final connection byte must be usable without a WINDOW_UPDATE") + .unwrap(); + assert_eq!(first, 1); } // https://github.com/hyperium/hyper/issues/4003, for HTTP/2 CONNECT