From b83ce0034726af1ac59215f2f49a4163f895bc66 Mon Sep 17 00:00:00 2001 From: Stephen Balogh Date: Mon, 14 Sep 2026 21:01:08 -0700 Subject: [PATCH 1/2] flow control --- sdk/src/api.rs | 2 +- sdk/src/client.rs | 42 +++++++++++++++------- sdk/src/types.rs | 91 +++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 122 insertions(+), 13 deletions(-) diff --git a/sdk/src/api.rs b/sdk/src/api.rs index 84175479..92d442f7 100644 --- a/sdk/src/api.rs +++ b/sdk/src/api.rs @@ -859,7 +859,7 @@ impl BaseClient { Compression::None => {} } - let client = client::Pool::new(connector); + let client = client::Pool::new(connector, config.http2); Ok(Self { client: Arc::new(client), diff --git a/sdk/src/client.rs b/sdk/src/client.rs index abf927ab..be63c012 100644 --- a/sdk/src/client.rs +++ b/sdk/src/client.rs @@ -39,7 +39,10 @@ use tokio::{ }; use tokio_util::task::AbortOnDropHandle; -use crate::frame_signal::{FrameSignal, RequestFrameMonitorBody}; +use crate::{ + frame_signal::{FrameSignal, RequestFrameMonitorBody}, + types::Http2Config, +}; const APPLICATION_JSON: HeaderValue = HeaderValue::from_static("application/json"); const MAX_CONCURRENT_REQUESTS_PER_CLIENT: usize = 90; @@ -786,6 +789,7 @@ impl Drop for RequestPermit { } struct PooledClient { + max_concurrent_requests: usize, id: ConnectionId, client: Arc>, active_requests: Arc, @@ -793,8 +797,9 @@ struct PooledClient { } impl PooledClient { - fn new(client: HyperClient) -> Self { + fn new(client: HyperClient, max_concurrent_requests: usize) -> Self { Self { + max_concurrent_requests, id: ConnectionId::next(), client: Arc::new(client), active_requests: Arc::new(AtomicUsize::new(0)), @@ -805,7 +810,7 @@ impl PooledClient { fn request_permit(&self) -> Option { self.active_requests .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |ar| { - (ar < MAX_CONCURRENT_REQUESTS_PER_CLIENT).then_some(ar + 1) + (ar < self.max_concurrent_requests).then_some(ar + 1) }) .ok()?; *self.idle_since.lock().unwrap() = None; @@ -827,6 +832,7 @@ impl PooledClient { } struct HostPool { + http2: Option, clients: StdRwLock>>, connector: C, } @@ -835,21 +841,31 @@ impl HostPool where C: Connect + Clone + Send + Sync + 'static, { - fn new(connector: C) -> Self { + fn new(connector: C, http2: Option) -> Self { Self { + http2, clients: StdRwLock::new(Vec::new()), connector, } } fn create_client(&self) -> PooledClient { - let client = HyperClient::builder(TokioExecutor::new()) + let mut builder = HyperClient::builder(TokioExecutor::new()); + builder .timer(TokioTimer::new()) .http2_only(true) .http2_keep_alive_interval(Duration::from_secs(20)) - .http2_keep_alive_timeout(Duration::from_secs(10)) - .build(self.connector.clone()); - PooledClient::new(client) + .http2_keep_alive_timeout(Duration::from_secs(10)); + let max_requests = if let Some(config) = self.http2 { + builder + .http2_initial_stream_window_size(config.stream_receive_window) + .http2_initial_connection_window_size(config.connection_receive_window) + .http2_adaptive_window(false); + config.max_concurrent_requests + } else { + MAX_CONCURRENT_REQUESTS_PER_CLIENT + }; + PooledClient::new(builder.build(self.connector.clone()), max_requests) } fn checkout(&self) -> (Arc>, RequestPermit, ConnectionId) { @@ -905,6 +921,7 @@ where } pub struct Pool { + http2: Option, hosts: Arc>>>>, connector: C, _reaper: AbortOnDropHandle<()>, @@ -914,7 +931,7 @@ impl Pool where C: Connect + Clone + Send + Sync + 'static, { - pub fn new(connector: C) -> Self { + pub fn new(connector: C, http2: Option) -> Self { let hosts = Arc::new(RwLock::new(HashMap::new())); let _reaper = AbortOnDropHandle::new(tokio::spawn({ @@ -929,6 +946,7 @@ where })); Self { + http2, hosts, connector, _reaper, @@ -945,7 +963,7 @@ where let mut hosts = self.hosts.write().await; hosts .entry(host.to_owned()) - .or_insert_with(|| Arc::new(HostPool::new(self.connector.clone()))) + .or_insert_with(|| Arc::new(HostPool::new(self.connector.clone(), self.http2))) .clone() } @@ -1008,7 +1026,7 @@ mod tests { const TEST_HOST: &str = "localhost:8080"; fn test_pool() -> Pool { - Pool::new(HttpConnector::new()) + Pool::new(HttpConnector::new(), None) } #[test] @@ -1113,7 +1131,7 @@ mod tests { .supported_schemes() ); } - + #[tokio::test] async fn checkout_within_capacity() { let pool = test_pool(); diff --git a/sdk/src/types.rs b/sdk/src/types.rs index 3e312852..bb121ffe 100644 --- a/sdk/src/types.rs +++ b/sdk/src/types.rs @@ -501,12 +501,69 @@ impl RetryConfig { } } +/// HTTP/2 receive-window sizes and per-connection request concurrency. +/// +/// The connection window must cover the combined stream windows of all +/// concurrent requests. These limits apply to unconsumed HTTP/2 DATA, not +/// total memory usage or decoded message size. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct Http2Config { + pub(crate) max_concurrent_requests: usize, + pub(crate) stream_receive_window: u32, + pub(crate) connection_receive_window: u32, +} + +impl Http2Config { + /// Configure the request cap per connection and receive windows in bytes. + /// + /// The request cap includes unary requests and streaming responses until + /// their bodies finish or are dropped. Additional requests use another + /// pooled connection. Adaptive receive windows are disabled. + /// + /// # Errors + /// + /// The request cap and stream window must be nonzero. Windows cannot exceed + /// 2^31 - 1 bytes, the connection window must be at least 65,535 bytes, and + /// it must cover `max_concurrent_requests * stream_receive_window`. + pub fn new( + max_concurrent_requests: usize, + stream_receive_window: u32, + connection_receive_window: u32, + ) -> Result { + if max_concurrent_requests == 0 { + return Err("HTTP/2 concurrent request limit must be nonzero".into()); + } + if !(1..=i32::MAX as u32).contains(&stream_receive_window) { + return Err("HTTP/2 stream receive window must be between 1 and 2^31 - 1 bytes".into()); + } + if !(65_535..=i32::MAX as u32).contains(&connection_receive_window) { + return Err( + "HTTP/2 connection receive window must be between 65,535 and 2^31 - 1 bytes".into(), + ); + } + if max_concurrent_requests + .checked_mul(stream_receive_window as usize) + .is_none_or(|total| total > connection_receive_window as usize) + { + return Err( + "HTTP/2 connection receive window must cover every concurrent stream window".into(), + ); + } + Ok(Self { + max_concurrent_requests, + stream_receive_window, + connection_receive_window, + }) + } +} + #[derive(Debug, Clone)] #[non_exhaustive] /// Configuration for [`S2`](crate::S2). pub struct S2Config { pub(crate) access_token: AccessToken, pub(crate) endpoints: S2Endpoints, + pub(crate) http2: Option, pub(crate) connection_timeout: Duration, pub(crate) request_timeout: Duration, pub(crate) retry: RetryConfig, @@ -523,6 +580,7 @@ impl S2Config { Self { access_token: AccessToken::Static(access_token.into().into()), endpoints: S2Endpoints::for_cloud(), + http2: None, connection_timeout: Duration::from_secs(3), request_timeout: Duration::from_secs(5), retry: RetryConfig::new(), @@ -598,6 +656,17 @@ impl S2Config { }) } + /// Set fixed HTTP/2 receive windows and the concurrent request cap. + /// + /// Without an override, the SDK uses its default transport configuration: + /// 90 concurrent requests per pooled client and Hyper's receive windows. + pub fn with_http2(self, http2: Http2Config) -> Self { + Self { + http2: Some(http2), + ..self + } + } + /// Set the timeout for establishing a connection to the server. /// /// Defaults to `3s`. @@ -4593,3 +4662,25 @@ mod tests { assert_eq!(record.timestamp, 1234); } } + +#[cfg(test)] +mod http2_config_tests { + use super::Http2Config; + + #[test] + fn receive_limits_cover_every_stream_without_overflow() { + assert!(Http2Config::new(16, 128 * 1024, 4 * 1024 * 1024).is_ok()); + assert!(Http2Config::new(16, 128 * 1024, 2 * 1024 * 1024).is_ok()); + for (requests, stream, connection) in [ + (0, 128 * 1024, 4 * 1024 * 1024), + (16, 0, 4 * 1024 * 1024), + (1, u32::MAX, u32::MAX), + (1, 1, 65_534), + (1, 1, u32::MAX), + (90, 2 * 1024 * 1024, 5 * 1024 * 1024), + (usize::MAX, 128 * 1024, i32::MAX as u32), + ] { + assert!(Http2Config::new(requests, stream, connection).is_err()); + } + } +} From 8599a20de911736552cab99cc92cf46b6a1c0c13 Mon Sep 17 00:00:00 2001 From: Stephen Balogh Date: Wed, 23 Sep 2026 17:21:36 -0700 Subject: [PATCH 2/2] fmt --- sdk/src/client.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/src/client.rs b/sdk/src/client.rs index be63c012..f28649a5 100644 --- a/sdk/src/client.rs +++ b/sdk/src/client.rs @@ -1131,7 +1131,7 @@ mod tests { .supported_schemes() ); } - + #[tokio::test] async fn checkout_within_capacity() { let pool = test_pool();