From 05afc96853b8207995aa59f630fc4c2a8183a25a Mon Sep 17 00:00:00 2001 From: Huapeng Zhou Date: Sun, 27 Sep 2026 16:36:51 -0700 Subject: [PATCH 1/2] fix(http1): withdraw a stale want when the client dispatcher takes a request The Receiver signals want when its queue reports Pending, and SendRequest::is_ready reports that want. A want signaled after finding the queue empty can land after the Sender's give() for the request that follows, and tokio's coop budget can make the queue report Pending with a request already queued. Either way the request is taken with want set, and the connection reports ready until the response completes. hyper-util's legacy pool returns connections on is_ready, so it hands a busy connection to the next request, which then waits behind the whole response. Withdraw the want with want::Taker::unwant when a request is taken. Want is signaled again once the dispatcher is idle. Requires want with Taker::unwant (seanmonstar/want#6). Closes #4207 --- Cargo.toml | 5 ++ src/client/dispatch.rs | 6 +++ tests/h1_ready_in_flight.rs | 95 +++++++++++++++++++++++++++++++++++++ 3 files changed, 106 insertions(+) create mode 100644 tests/h1_ready_in_flight.rs diff --git a/Cargo.toml b/Cargo.toml index 2c98449797..66fc1058dc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -352,3 +352,8 @@ 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"] diff --git a/src/client/dispatch.rs b/src/client/dispatch.rs index 0d91a1d92e..37afa9f78d 100644 --- a/src/client/dispatch.rs +++ b/src/client/dispatch.rs @@ -182,6 +182,12 @@ impl Receiver { pub(crate) fn poll_recv(&mut self, cx: &mut Context<'_>) -> Poll)>> { 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 => { diff --git a/tests/h1_ready_in_flight.rs b/tests/h1_ready_in_flight.rs new file mode 100644 index 0000000000..f2317ba308 --- /dev/null +++ b/tests/h1_ready_in_flight.rs @@ -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> { + Poll::Pending + } +} + +impl Write for SilentIo { + fn poll_write( + self: Pin<&mut Self>, + _: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + Poll::Ready(Ok(buf.len())) + } + + fn poll_flush(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + + fn poll_shutdown(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } +} + +/// Runs the connection task once; it never finishes, as the peer never answers. +fn poll_conn(conn: &mut Connection>, 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" + ); +} From 9d2b1f1e47b0c0a9b5dffdd9aa3aa8efa6391706 Mon Sep 17 00:00:00 2001 From: Huapeng Zhou Date: Sun, 27 Sep 2026 23:15:25 -0700 Subject: [PATCH 2/2] chore(dependencies): require want 0.3.2 --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index 66fc1058dc..93b03527fe 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -39,7 +39,7 @@ itoa = { version = "1", optional = true } pin-project-lite = { version = "0.2.4", optional = true } smallvec = { version = "1.12", features = ["const_generics", "const_new"], optional = true } tracing = { version = "0.1", default-features = false, features = ["std"], optional = true } -want = { version = "0.3", optional = true } +want = { version = "0.3.2", optional = true } [dev-dependencies] form_urlencoded = "1"