From 7dcccb42419961e0dd5fbc3151e0dec424bad051 Mon Sep 17 00:00:00 2001 From: Claudia Zhu Date: Wed, 5 Aug 2026 23:52:04 +0000 Subject: [PATCH 1/2] Honor server-advertised MAX_CONCURRENT_STREAMS in h2 stream admission MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit spawn_stream() admitted streams against the configured max_h2_streams only, snapshotted before the server's initial SETTINGS frame is processed. A server advertising a lower limit (e.g. the worker HTTP relay's per-connection stream cap) caused excess requests to queue silently inside h2 — parked in ready() behind long-lived streams until the downstream deadline killed them — instead of signaling the caller to dial another connection. - Re-read current_max_send_streams() on every admission and admit against min(configured, advertised). - Bound the ready() wait: a connection at its real capacity reports "no free stream" so the caller dials instead of queueing. - Fix a counter leak when spawn_stream() is cancelled mid-wait. Co-Authored-By: Claude Fable 5 --- pingora-core/src/connectors/http/v2.rs | 124 +++++++++++++++++++++++-- 1 file changed, 118 insertions(+), 6 deletions(-) diff --git a/pingora-core/src/connectors/http/v2.rs b/pingora-core/src/connectors/http/v2.rs index c8e804d46..f671f8492 100644 --- a/pingora-core/src/connectors/http/v2.rs +++ b/pingora-core/src/connectors/http/v2.rs @@ -44,6 +44,40 @@ impl Stub { } } +// How long to wait for h2 to accept a new stream on a connection that passed +// the admission check. ready() resolves ~immediately below the server's +// advertised stream limit; a longer wait means the connection is actually at +// capacity (possible before the server's initial SETTINGS frame is processed, +// when the advertised limit isn't known yet), and the caller is better served +// by dialing another connection than by queueing behind long-lived streams. +const H2_STREAM_READY_TIMEOUT: Duration = Duration::from_millis(500); + +// Decrements a stream counter on drop unless disarmed. spawn_stream() can be +// cancelled while parked in ready() (callers race it against the downstream +// request lifetime); without this guard the increment would leak and the +// connection would permanently appear busier than it is. +struct StreamCounterGuard<'a>(Option<&'a AtomicUsize>); + +impl<'a> StreamCounterGuard<'a> { + fn new(counter: &'a AtomicUsize) -> Self { + Self(Some(counter)) + } + + // The stream was created: its slot is now owned by the Http2Session and + // released via release_stream(). + fn disarm(mut self) { + self.0.take(); + } +} + +impl Drop for StreamCounterGuard<'_> { + fn drop(&mut self) { + if let Some(counter) = self.0 { + counter.fetch_sub(1, Ordering::SeqCst); + } + } +} + pub(crate) struct ConnectionRefInner { connection_stub: Stub, closed: watch::Receiver, @@ -130,20 +164,46 @@ impl ConnectionRef { // spawn a stream if more stream is allowed, otherwise return Ok(None) pub async fn spawn_stream(&self) -> Result> { + // Admit against the smaller of our configured limit and the limit the + // server currently advertises via SETTINGS_MAX_CONCURRENT_STREAMS. + // The advertised value is only known once the server's initial + // SETTINGS frame has been processed (until then h2 reports the + // client-side initial value), so it must be re-read on every + // admission rather than snapshotted at handshake: a server that + // advertises fewer streams than max_streams would otherwise have + // excess requests silently queued inside h2 — parked until a + // long-lived stream finishes — instead of signaling the caller to + // dial another connection. + let max_streams = self + .0 + .max_streams + .min(self.0.connection_stub.0.current_max_send_streams()); // Atomically check if the current_stream is over the limit // load(), compare and then fetch_add() cannot guarantee the same let current_streams = self.0.current_streams.fetch_add(1, Ordering::SeqCst); - if current_streams >= self.0.max_streams { + if current_streams >= max_streams { // already over the limit, reset the counter to the previous value self.0.current_streams.fetch_sub(1, Ordering::SeqCst); return Ok(None); } - match self.0.connection_stub.new_stream().await { - Ok(send_req) => Ok(Some(Http2Session::new(send_req, self.clone()))), - Err(e) => { - // fail to create the stream, reset the counter - self.0.current_streams.fetch_sub(1, Ordering::SeqCst); + // Undo the increment above if stream creation fails, times out, or + // this future is cancelled while waiting. + let guard = StreamCounterGuard::new(&self.0.current_streams); + + match pingora_timeout::timeout(H2_STREAM_READY_TIMEOUT, self.0.connection_stub.new_stream()) + .await + { + // Connection is at its real capacity even though the admission + // check passed (e.g. the server's initial SETTINGS had not been + // processed yet). Report no free stream so the caller dials a + // new connection instead of queueing invisibly. + Err(_elapsed) => Ok(None), + Ok(Ok(send_req)) => { + guard.disarm(); + Ok(Some(Http2Session::new(send_req, self.clone()))) + } + Ok(Err(e)) => { // Remote sends GOAWAY(NO_ERROR): graceful shutdown: this connection no longer // accepts new streams. We can still try to create new connection. if e.root_cause() @@ -593,6 +653,58 @@ mod tests { } } + #[tokio::test] + async fn test_h2_admission_respects_server_advertised_max_streams() { + use tokio::net::TcpListener; + + // An h2c server that advertises SETTINGS_MAX_CONCURRENT_STREAMS=1, + // like a gateway that wants at most one stream per connection. + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + loop { + let (socket, _) = listener.accept().await.unwrap(); + tokio::spawn(async move { + let mut conn = h2::server::Builder::new() + .max_concurrent_streams(1) + .handshake::<_, Bytes>(socket) + .await + .unwrap(); + // Keep the connection alive without answering anything so + // the first stream's slot stays occupied. + while (conn.accept().await).is_some() {} + }); + } + }); + + let connector = Connector::new(None); + let mut peer = HttpPeer::new(addr, false, "".into()); + peer.options.set_http_version(2, 2); + peer.options.max_h2_streams = 1024; // far above the server's limit + + // The first stream is admitted before the server's SETTINGS is known. + let h2 = connector + .new_http_session::(&peer) + .await + .unwrap(); + let _h2_1 = match h2 { + HttpSession::H2(h2_stream) => h2_stream, + _ => panic!("expect h2"), + }; + + // Let the connection task process the server's initial SETTINGS. + tokio::time::sleep(Duration::from_millis(100)).await; + + // 1023 configured slots remain, but the server allows one concurrent + // stream: the connector must report "no free stream" so the caller + // dials a new connection, not queue the request inside h2 behind the + // live stream. + let start = std::time::Instant::now(); + let reused = connector.reused_http_session(&peer).await.unwrap(); + assert!(reused.is_none()); + assert!(start.elapsed() < H2_STREAM_READY_TIMEOUT); + } + #[tokio::test] #[cfg(feature = "any_tls")] async fn test_h2_single_stream() { From ecdee5e9802b8cd8962e17940f4042457a9cae19 Mon Sep 17 00:00:00 2001 From: molocule <34072934+molocule@users.noreply.github.com> Date: Wed, 5 Aug 2026 18:35:33 -0600 Subject: [PATCH 2/2] Update v2.rs --- pingora-core/src/connectors/http/v2.rs | 49 ++++++++++---------------- 1 file changed, 19 insertions(+), 30 deletions(-) diff --git a/pingora-core/src/connectors/http/v2.rs b/pingora-core/src/connectors/http/v2.rs index f671f8492..c9d32f39c 100644 --- a/pingora-core/src/connectors/http/v2.rs +++ b/pingora-core/src/connectors/http/v2.rs @@ -44,18 +44,9 @@ impl Stub { } } -// How long to wait for h2 to accept a new stream on a connection that passed -// the admission check. ready() resolves ~immediately below the server's -// advertised stream limit; a longer wait means the connection is actually at -// capacity (possible before the server's initial SETTINGS frame is processed, -// when the advertised limit isn't known yet), and the caller is better served -// by dialing another connection than by queueing behind long-lived streams. -const H2_STREAM_READY_TIMEOUT: Duration = Duration::from_millis(500); - // Decrements a stream counter on drop unless disarmed. spawn_stream() can be -// cancelled while parked in ready() (callers race it against the downstream -// request lifetime); without this guard the increment would leak and the -// connection would permanently appear busier than it is. +// cancelled while creating a stream; without this guard the increment would +// leak and the connection would permanently appear busier than it is. struct StreamCounterGuard<'a>(Option<&'a AtomicUsize>); impl<'a> StreamCounterGuard<'a> { @@ -187,23 +178,16 @@ impl ConnectionRef { return Ok(None); } - // Undo the increment above if stream creation fails, times out, or - // this future is cancelled while waiting. + // Undo the increment above if stream creation fails or this future is + // cancelled while waiting. let guard = StreamCounterGuard::new(&self.0.current_streams); - match pingora_timeout::timeout(H2_STREAM_READY_TIMEOUT, self.0.connection_stub.new_stream()) - .await - { - // Connection is at its real capacity even though the admission - // check passed (e.g. the server's initial SETTINGS had not been - // processed yet). Report no free stream so the caller dials a - // new connection instead of queueing invisibly. - Err(_elapsed) => Ok(None), - Ok(Ok(send_req)) => { + match self.0.connection_stub.new_stream().await { + Ok(send_req) => { guard.disarm(); Ok(Some(Http2Session::new(send_req, self.clone()))) } - Ok(Err(e)) => { + Err(e) => { // Remote sends GOAWAY(NO_ERROR): graceful shutdown: this connection no longer // accepts new streams. We can still try to create new connection. if e.root_cause() @@ -670,8 +654,8 @@ mod tests { .handshake::<_, Bytes>(socket) .await .unwrap(); - // Keep the connection alive without answering anything so - // the first stream's slot stays occupied. + // Keep the connection alive while the first Pingora + // session remains active. while (conn.accept().await).is_some() {} }); } @@ -687,22 +671,27 @@ mod tests { .new_http_session::(&peer) .await .unwrap(); - let _h2_1 = match h2 { + let h2_1 = match h2 { HttpSession::H2(h2_stream) => h2_stream, _ => panic!("expect h2"), }; - // Let the connection task process the server's initial SETTINGS. - tokio::time::sleep(Duration::from_millis(100)).await; + // Wait until the connection task has processed the server's initial + // SETTINGS rather than relying on scheduler timing. + tokio::time::timeout(Duration::from_secs(1), async { + while h2_1.conn.0.connection_stub.0.current_max_send_streams() != 1 { + tokio::task::yield_now().await; + } + }) + .await + .expect("server SETTINGS was not processed"); // 1023 configured slots remain, but the server allows one concurrent // stream: the connector must report "no free stream" so the caller // dials a new connection, not queue the request inside h2 behind the // live stream. - let start = std::time::Instant::now(); let reused = connector.reused_http_session(&peer).await.unwrap(); assert!(reused.is_none()); - assert!(start.elapsed() < H2_STREAM_READY_TIMEOUT); } #[tokio::test]