Skip to content
Draft
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
9 changes: 9 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -352,3 +352,12 @@ required-features = ["full"]
name = "h1_shutdown_while_buffered"
path = "tests/h1_shutdown_while_buffered.rs"
required-features = ["full"]

[[test]]
name = "h1_ready_in_flight"
path = "tests/h1_ready_in_flight.rs"
required-features = ["full"]

[patch.crates-io]
# Temporary, until want releases Taker::unwant (seanmonstar/want#6).
want = { git = "https://github.com/shodoco/want", branch = "taker-unwant" }
6 changes: 6 additions & 0 deletions src/client/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,12 @@ impl<T, U> Receiver<T, U> {
pub(crate) fn poll_recv(&mut self, cx: &mut Context<'_>) -> Poll<Option<(T, Callback<T, U>)>> {
match self.inner.poll_recv(cx) {
Poll::Ready(item) => {
// A want signaled after finding the queue empty can land after
// the Sender's `give()` for the message taken here, as can one
// signaled on a coop-budget Pending with a message queued.
// Withdraw it, or the Sender reports this connection as ready
// while it is still serving this message.
self.taker.unwant();
Poll::Ready(item.map(|mut env| env.0.take().expect("envelope not dropped")))
}
Poll::Pending => {
Expand Down
95 changes: 95 additions & 0 deletions tests/h1_ready_in_flight.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
// Test: an HTTP/1 client connection must not report ready while a request is
// in flight on it.
//
// The dispatch `Receiver` signals want when its queue reports `Pending`, and
// `SendRequest::is_ready` reports that want. The signal can go stale: the
// connection task can find its queue empty and a request land before it
// signals, or tokio's coop budget can make the queue report `Pending` with a
// request already in it. The request is then taken with want still set, so
// `is_ready` stays true until the response is complete. Pools that return a
// connection on `is_ready`, such as hyper-util's legacy client, then give that
// connection to the next request, which waits behind the whole response.
//
// The coop path is deterministic, so the test drives that one.

use std::future::Future;
use std::io;
use std::pin::Pin;
use std::task::{Context, Poll};

use bytes::Bytes;
use futures_util::future::poll_fn;
use http_body_util::Empty;
use hyper::client::conn::http1::{self, Connection};
use hyper::rt::{Read, ReadBufCursor, Write};
use hyper::Request;

/// Accepts every write and never answers, so a sent request stays in flight.
struct SilentIo;

impl Read for SilentIo {
fn poll_read(
self: Pin<&mut Self>,
_: &mut Context<'_>,
_: ReadBufCursor<'_>,
) -> Poll<io::Result<()>> {
Poll::Pending
}
}

impl Write for SilentIo {
fn poll_write(
self: Pin<&mut Self>,
_: &mut Context<'_>,
buf: &[u8],
) -> Poll<io::Result<usize>> {
Poll::Ready(Ok(buf.len()))
}

fn poll_flush(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<io::Result<()>> {
Poll::Ready(Ok(()))
}

fn poll_shutdown(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<io::Result<()>> {
Poll::Ready(Ok(()))
}
}

/// Runs the connection task once; it never finishes, as the peer never answers.
fn poll_conn(conn: &mut Connection<SilentIo, Empty<Bytes>>, cx: &mut Context<'_>) -> Poll<()> {
assert!(Pin::new(conn).poll(cx).is_pending());
Poll::Ready(())
}

#[tokio::test]
async fn h1_connection_is_not_ready_while_a_request_is_in_flight() {
let (mut sender, mut conn) = http1::handshake(SilentIo).await.unwrap();

// Idle, the connection finds its queue empty and signals it is ready.
poll_fn(|cx| poll_conn(&mut conn, cx)).await;
assert!(sender.is_ready());

// Kept alive: dropping the response future cancels the request.
let _response = sender.send_request(Request::new(Empty::new()));
assert!(!sender.is_ready());

// Out of coop budget, the queue reports Pending with the request in it, so
// the connection signals ready again. A connection task preempted between
// finding its queue empty and signaling leaves the same stale signal.
poll_fn(|cx| {
while let Poll::Ready(restore) = tokio::task::coop::poll_proceed(cx) {
restore.made_progress();
}
poll_conn(&mut conn, cx)
})
.await;
tokio::task::yield_now().await;

// With a fresh budget it takes and writes the request, which then waits
// for a response that never comes.
poll_fn(|cx| poll_conn(&mut conn, cx)).await;
assert!(
!sender.is_ready(),
"connection reports ready with a request in flight"
);
}