diff --git a/tests/client.rs b/tests/client.rs index e4c8d12cf6..73ab121ff0 100644 --- a/tests/client.rs +++ b/tests/client.rs @@ -3345,6 +3345,127 @@ mod conn { let _ = tokio::time::timeout(Duration::from_secs(5), a_handle).await; } + #[tokio::test] + 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); + } + + #[tokio::test] + async fn h2_body_uses_the_last_byte_of_connection_capacity() { + const STREAM_A_LEN: usize = 65534; + + let (client_io, server_io, _) = setup_duplex_test_server(); + let (a_full_tx, a_full_rx) = oneshot::channel(); + let (b_first_tx, b_first_rx) = oneshot::channel(); + tokio::spawn(async move { + 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, connection) = conn::http2::Builder::new(TokioExecutor) + .handshake::<_, BoxBody>>(io) + .await + .expect("http handshake"); + tokio::spawn(async move { + let _ = connection.await; + }); + + let (mut a_tx, a_rx) = + mpsc::channel::, Box>>(4); + let mut client_a = client.clone(); + 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 len = remaining.min(16_384); + a_tx.send(Ok(Frame::data(Bytes::from(vec![0; len])))) + .await + .unwrap(); + remaining -= len; + } + a_full_rx.await.unwrap(); + + 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("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 // // Like `h2_idle_stream_does_not_pin_connection_window`, but the idle