You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
hyper 1.10.1, 1.11.1 and master (e60932d); want 0.3.1; hyper-util 0.1.20 and 0.1.21
Platform
Linux 6.16 x86_64 (not platform-specific)
Summary
On an HTTP/1 client connection, SendRequest::is_ready() can keep returning true after the connection has taken a request, until that request's response is complete. hyper-util's legacy pool trusts is_ready() when a response head arrives, so it pools a connection that is still streaming a response body. The next request to check out that connection is written only after the whole previous response has been read.
This isn't HTTP/1.1 pipelining. It happens with default settings, and no option turns it off.
How readiness works
SendRequest and the connection task share a want flag. The connection task signals want when its request queue is empty, send_request clears it (Giver::give) when it queues a request, and is_ready() reports it.
sequenceDiagram
participant S as SendRequest
participant F as want flag
participant T as connection task
T->>T: poll_recv: queue empty, Pending
T->>F: want(): Idle → Want
S->>F: is_ready()? true
S->>F: send_request: give(), Want → Idle
S->>T: request queued
T->>T: poll_recv: Ready(request)
Note over T: writes the request and reads the response.<br/>It won't poll the queue again until the connection is idle.
S->>F: is_ready()? false
Loading
How the flag goes stale
The want can be set again after give() cleared it, with the request already queued. There are two ways.
A race. The connection task finds the queue empty, but before it calls taker.want(), send_request on another thread runs give() and queues a request. The task's want() lands after the give().
sequenceDiagram
participant S as SendRequest (thread A)
participant F as want flag
participant T as connection task (thread B)
Note over F: Want (connection idle)
T->>T: poll_recv: queue empty
rect rgb(255, 228, 228)
S->>F: give(): Want → Idle
S->>T: request queued
T->>F: want(): Idle → Want (stale)
end
T->>T: poll_recv: Ready(request), flag stays Want
Loading
tokio's coop budget. Once the task's budget is used up, mpsc::UnboundedReceiver::poll_recv returns Pending even though a request is queued, so the task signals want anyway.
sequenceDiagram
participant S as SendRequest
participant F as want flag
participant T as connection task
Note over F: Want (connection idle)
S->>F: give(): Want → Idle
S->>T: request queued
rect rgb(255, 228, 228)
T->>T: poll_recv: budget used up, Pending
T->>F: want(): Idle → Want (stale)
end
T->>T: next poll: Ready(request), flag stays Want
Loading
Either way the dispatcher takes the request with want still set. While the request is in flight the dispatcher doesn't poll the queue, so nothing clears the flag, and is_ready() stays true until the response completes.
What that does to a pool
hyper-util's legacy client returns a connection to the pool when the response head arrives, if is_ready() is true (client/legacy/client.rs: if pooled.is_http2() || !pooled.is_pool_enabled() || pooled.is_ready()).
sequenceDiagram
participant A as request A (long stream)
participant P as hyper-util pool
participant C as connection
participant B as request B (short)
A->>C: sent, response head arrives
P->>C: is_ready()? true (stale)
P->>P: pools C while A's body is still streaming
B->>P: checkout
P->>B: C
B->>C: queued behind A
Note over B,C: B is written only after A's whole body has been read
Loading
We hit this in a reverse proxy that sends short GET probes and long streaming POSTs through one reqwest client. Probes regularly waited behind streams for hundreds of milliseconds to seconds. With the fix below, requests that waited behind a stream went from hundreds per 30-second run to zero in a standalone repro, and the proxy's tail latency dropped in production.
Code Sample
This test drives path 2 deterministically, using only the public API. On master it fails at the final assert!(!sender.is_ready()). Run it with cargo test --features full --test h1_ready_in_flight.
tests/h1_ready_in_flight.rs
// 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.structSilentIo;implReadforSilentIo{fnpoll_read(self:Pin<&mutSelf>,
_:&mutContext<'_>,
_:ReadBufCursor<'_>,) -> Poll<io::Result<()>>{Poll::Pending}}implWriteforSilentIo{fnpoll_write(self:Pin<&mutSelf>,
_:&mutContext<'_>,buf:&[u8],) -> Poll<io::Result<usize>>{Poll::Ready(Ok(buf.len()))}fnpoll_flush(self:Pin<&mutSelf>, _:&mutContext<'_>) -> Poll<io::Result<()>>{Poll::Ready(Ok(()))}fnpoll_shutdown(self:Pin<&mutSelf>, _:&mutContext<'_>) -> Poll<io::Result<()>>{Poll::Ready(Ok(()))}}/// Runs the connection task once; it never finishes, as the peer never answers.fnpoll_conn(conn:&mutConnection<SilentIo,Empty<Bytes>>,cx:&mutContext<'_>) -> Poll<()>{assert!(Pin::new(conn).poll(cx).is_pending());Poll::Ready(())}#[tokio::test]asyncfnh1_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| {whileletPoll::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");}
Expected Behavior
Once the connection takes a request, is_ready() returns false until the connection is idle and can take another request.
Actual Behavior
is_ready() returns true while the request is in flight. hyper-util's pool then hands the busy connection to the next request, which waits for the whole previous response.
Additional Context
Proposed fix
When Receiver::poll_recv takes a request, withdraw any outstanding want with a new want::Taker::unwant(): a compare-exchange from Want to Idle, the taker-side counterpart of Giver::give. The connection task signals want again only when it is idle.
flowchart LR
R["Receiver::poll_recv"] -->|Pending| W["taker.want()<br/>signal ready"]
R -->|"Ready(request)"| U["taker.unwant() (new)<br/>withdraw a stale want"]
U --> D["dispatcher sends the request"]
Loading
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")))
}
With both changes the test passes, and cargo test --features full passes on master. HTTP/2 is unaffected: its is_ready() only checks whether the connection is closed.
A sender-side reorder (queue the request first, then give()) isn't enough. It leaves both paths open, and in our repro a few requests per run still waited behind a stream.
Version
hyper 1.10.1, 1.11.1 and master (e60932d); want 0.3.1; hyper-util 0.1.20 and 0.1.21
Platform
Linux 6.16 x86_64 (not platform-specific)
Summary
On an HTTP/1 client connection,
SendRequest::is_ready()can keep returningtrueafter the connection has taken a request, until that request's response is complete. hyper-util's legacy pool trustsis_ready()when a response head arrives, so it pools a connection that is still streaming a response body. The next request to check out that connection is written only after the whole previous response has been read.This isn't HTTP/1.1 pipelining. It happens with default settings, and no option turns it off.
How readiness works
SendRequestand the connection task share awantflag. The connection task signals want when its request queue is empty,send_requestclears it (Giver::give) when it queues a request, andis_ready()reports it.sequenceDiagram participant S as SendRequest participant F as want flag participant T as connection task T->>T: poll_recv: queue empty, Pending T->>F: want(): Idle → Want S->>F: is_ready()? true S->>F: send_request: give(), Want → Idle S->>T: request queued T->>T: poll_recv: Ready(request) Note over T: writes the request and reads the response.<br/>It won't poll the queue again until the connection is idle. S->>F: is_ready()? falseHow the flag goes stale
The want can be set again after
give()cleared it, with the request already queued. There are two ways.taker.want(),send_requeston another thread runsgive()and queues a request. The task'swant()lands after thegive().sequenceDiagram participant S as SendRequest (thread A) participant F as want flag participant T as connection task (thread B) Note over F: Want (connection idle) T->>T: poll_recv: queue empty rect rgb(255, 228, 228) S->>F: give(): Want → Idle S->>T: request queued T->>F: want(): Idle → Want (stale) end T->>T: poll_recv: Ready(request), flag stays Wantmpsc::UnboundedReceiver::poll_recvreturnsPendingeven though a request is queued, so the task signals want anyway.sequenceDiagram participant S as SendRequest participant F as want flag participant T as connection task Note over F: Want (connection idle) S->>F: give(): Want → Idle S->>T: request queued rect rgb(255, 228, 228) T->>T: poll_recv: budget used up, Pending T->>F: want(): Idle → Want (stale) end T->>T: next poll: Ready(request), flag stays WantEither way the dispatcher takes the request with want still set. While the request is in flight the dispatcher doesn't poll the queue, so nothing clears the flag, and
is_ready()staystrueuntil the response completes.What that does to a pool
hyper-util's legacy client returns a connection to the pool when the response head arrives, if
is_ready()is true (client/legacy/client.rs:if pooled.is_http2() || !pooled.is_pool_enabled() || pooled.is_ready()).sequenceDiagram participant A as request A (long stream) participant P as hyper-util pool participant C as connection participant B as request B (short) A->>C: sent, response head arrives P->>C: is_ready()? true (stale) P->>P: pools C while A's body is still streaming B->>P: checkout P->>B: C B->>C: queued behind A Note over B,C: B is written only after A's whole body has been readWe hit this in a reverse proxy that sends short GET probes and long streaming POSTs through one reqwest client. Probes regularly waited behind streams for hundreds of milliseconds to seconds. With the fix below, requests that waited behind a stream went from hundreds per 30-second run to zero in a standalone repro, and the proxy's tail latency dropped in production.
Code Sample
This test drives path 2 deterministically, using only the public API. On master it fails at the final
assert!(!sender.is_ready()). Run it withcargo test --features full --test h1_ready_in_flight.tests/h1_ready_in_flight.rsExpected Behavior
Once the connection takes a request,
is_ready()returnsfalseuntil the connection is idle and can take another request.Actual Behavior
is_ready()returnstruewhile the request is in flight. hyper-util's pool then hands the busy connection to the next request, which waits for the whole previous response.Additional Context
Proposed fix
When
Receiver::poll_recvtakes a request, withdraw any outstanding want with a newwant::Taker::unwant(): a compare-exchange from Want to Idle, the taker-side counterpart ofGiver::give. The connection task signals want again only when it is idle.flowchart LR R["Receiver::poll_recv"] -->|Pending| W["taker.want()<br/>signal ready"] R -->|"Ready(request)"| U["taker.unwant() (new)<br/>withdraw a stale want"] U --> D["dispatcher sends the request"]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"))) }Taker::unwant.Receiver::poll_recvand adds the test above. It's a draft, building against the want PR branch, until want has a release.With both changes the test passes, and
cargo test --features fullpasses on master. HTTP/2 is unaffected: itsis_ready()only checks whether the connection is closed.A sender-side reorder (queue the request first, then
give()) isn't enough. It leaves both paths open, and in our repro a few requests per run still waited behind a stream.