Skip to content
Closed
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
121 changes: 121 additions & 0 deletions tests/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Bytes, Box<dyn Error + Send + Sync>>>(io)
.await
.expect("http handshake");
tokio::spawn(async move {
let _ = connection.await;
});

let (mut tx, rx) = mpsc::channel::<Result<Frame<Bytes>, Box<dyn Error + Send + Sync>>>(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<Bytes, Box<dyn Error + Send + Sync>>>(io)
.await
.expect("http handshake");
tokio::spawn(async move {
let _ = connection.await;
});

let (mut a_tx, a_rx) =
mpsc::channel::<Result<Frame<Bytes>, Box<dyn Error + Send + Sync>>>(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
Expand Down