diff --git a/Cargo.lock b/Cargo.lock index de3da95d..3a5b6c33 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5038,7 +5038,7 @@ dependencies = [ [[package]] name = "s2-cli" -version = "0.43.0" +version = "0.43.1" dependencies = [ "assert_cmd", "async-stream", @@ -5123,7 +5123,7 @@ dependencies = [ [[package]] name = "s2-lite" -version = "0.43.0" +version = "0.43.1" dependencies = [ "async-stream", "async-trait", @@ -5183,7 +5183,7 @@ dependencies = [ [[package]] name = "s2-sdk" -version = "0.35.0" +version = "0.35.1" dependencies = [ "assert_matches", "async-compression", @@ -5244,7 +5244,7 @@ dependencies = [ [[package]] name = "s2-testcontainers" -version = "0.43.0" +version = "0.43.1" dependencies = [ "reqwest", "s2-sdk", diff --git a/charts/s2-lite-helm/Chart.yaml b/charts/s2-lite-helm/Chart.yaml index 953a1674..c243fb56 100644 --- a/charts/s2-lite-helm/Chart.yaml +++ b/charts/s2-lite-helm/Chart.yaml @@ -2,8 +2,8 @@ apiVersion: v2 name: s2-lite-helm description: Self-hostable S2 streaming datastore using SlateDB on object storage type: application -version: 0.1.71 -appVersion: "0.43.0" +version: 0.1.72 +appVersion: "0.43.1" keywords: - s2 - streaming diff --git a/cli/CHANGELOG.md b/cli/CHANGELOG.md index 9004262c..6d7e4fcb 100644 --- a/cli/CHANGELOG.md +++ b/cli/CHANGELOG.md @@ -2,6 +2,10 @@ All notable changes to this project will be documented in this file. +## [0.43.1] - 2026-09-28 + + + ## [0.43.0] - 2026-09-25 ### Features diff --git a/cli/Cargo.toml b/cli/Cargo.toml index b5284ce4..f3af2ba8 100644 --- a/cli/Cargo.toml +++ b/cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "s2-cli" -version = "0.43.0" +version = "0.43.1" description = "CLI for S2" edition.workspace = true license.workspace = true diff --git a/lite/CHANGELOG.md b/lite/CHANGELOG.md index 4f80ead3..833fb00f 100644 --- a/lite/CHANGELOG.md +++ b/lite/CHANGELOG.md @@ -2,6 +2,14 @@ All notable changes to this project will be documented in this file. +## [0.43.1] - 2026-09-28 + +### Bug Fixes + +- Resolve AWS region from profile chain for static credentials ([#780](https://github.com/s2-streamstore/s2/issues/780)) + + + ## [0.43.0] - 2026-09-25 ### Features diff --git a/lite/Cargo.toml b/lite/Cargo.toml index 8f0d8988..0600233d 100644 --- a/lite/Cargo.toml +++ b/lite/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "s2-lite" -version = "0.43.0" +version = "0.43.1" description = "Lightweight server implementation of S2, the durable streams API, backed by object storage" edition.workspace = true license.workspace = true diff --git a/lite/src/server.rs b/lite/src/server.rs index e0d62153..d78f71f7 100644 --- a/lite/src/server.rs +++ b/lite/src/server.rs @@ -400,6 +400,12 @@ async fn s3_builder() -> object_store::aws::AmazonS3Builder { (Some(key_id), Some(secret_key)) => { info!(key_id, "using static credentials from env vars"); + let aws_config = aws_config::load_defaults(aws_config::BehaviorVersion::latest()).await; + if let Some(region) = aws_config.region() { + info!(region = region.as_ref()); + builder = builder.with_region(region.to_string()); + } + let token = std::env::var_os("AWS_SESSION_TOKEN").and_then(|s| s.into_string().ok()); builder = builder.with_credentials(Arc::new( object_store::StaticCredentialProvider::new(object_store::aws::AwsCredential { diff --git a/sdk/CHANGELOG.md b/sdk/CHANGELOG.md index ed7c6a3b..2268ccd6 100644 --- a/sdk/CHANGELOG.md +++ b/sdk/CHANGELOG.md @@ -2,6 +2,14 @@ All notable changes to this project will be documented in this file. +## [0.35.1] - 2026-09-28 + +### Features + +- Surface http2 knobs relevant to flow control ([#776](https://github.com/s2-streamstore/s2/issues/776)) + + + ## [0.35.0] - 2026-09-25 ### Features diff --git a/sdk/Cargo.toml b/sdk/Cargo.toml index db739f37..194465ff 100644 --- a/sdk/Cargo.toml +++ b/sdk/Cargo.toml @@ -1,7 +1,7 @@ [package] name = "s2-sdk" description = "Rust SDK for S2" -version = "0.35.0" +version = "0.35.1" edition.workspace = true license.workspace = true repository = "https://github.com/s2-streamstore/s2/tree/main/sdk" 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..781844a2 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: Http2Config, 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: Http2Config) -> 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_initial_stream_window_size(self.http2.stream_receive_window) + .http2_initial_connection_window_size(self.http2.connection_receive_window); + let max_concurrent_requests = self + .http2 + .max_concurrent_requests + .unwrap_or(MAX_CONCURRENT_REQUESTS_PER_CLIENT); + PooledClient::new( + builder.build(self.connector.clone()), + max_concurrent_requests, + ) } fn checkout(&self) -> (Arc>, RequestPermit, ConnectionId) { @@ -905,6 +921,7 @@ where } pub struct Pool { + http2: Http2Config, 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: Http2Config) -> 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(), Http2Config::new()) } #[test] @@ -1140,6 +1158,22 @@ mod tests { assert_eq!(host_client_count(&pool, TEST_HOST).await, 2); } + #[tokio::test] + async fn custom_request_cap_bounds_client() { + let http2 = Http2Config::new().with_max_concurrent_requests(2).unwrap(); + let pool = Pool::new(HttpConnector::new(), http2); + let mut permits = Vec::new(); + for _ in 0..2 { + let (_client, permit, _) = pool.checkout(TEST_HOST).await; + permits.push(permit); + } + assert_eq!(host_client_count(&pool, TEST_HOST).await, 1); + + let (_client, permit, _) = pool.checkout(TEST_HOST).await; + permits.push(permit); + assert_eq!(host_client_count(&pool, TEST_HOST).await, 2); + } + #[tokio::test] async fn permit_drop_frees_capacity() { let pool = test_pool(); diff --git a/sdk/src/types.rs b/sdk/src/types.rs index 099048d2..a5f22cc7 100644 --- a/sdk/src/types.rs +++ b/sdk/src/types.rs @@ -502,12 +502,115 @@ impl RetryConfig { } } +/// Overrides for the HTTP/2 transport used by pooled connections. +/// +/// Unset knobs keep the SDK defaults: 90 concurrent requests per connection and +/// Hyper's receive windows (2 MiB per stream, 5 MiB per connection). Receive +/// windows bound unconsumed server-to-client DATA only; they do not affect +/// append throughput, which is governed by the server's receive windows. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +#[non_exhaustive] +pub struct Http2Config { + pub(crate) max_concurrent_requests: Option, + pub(crate) stream_receive_window: Option, + pub(crate) connection_receive_window: Option, +} + +impl Http2Config { + /// Stream limit advertised by S2 servers per HTTP/2 connection. + pub const MAX_CONCURRENT_REQUESTS: usize = 100; + + /// Maximum HTTP/2 flow-control window in bytes (`2^31 - 1`). + pub const MAX_RECEIVE_WINDOW: u32 = i32::MAX as u32; + + /// Minimum HTTP/2 connection receive window in bytes. + /// + /// The connection window cannot be shrunk below the protocol default. + pub const MIN_CONNECTION_RECEIVE_WINDOW: u32 = 65_535; + + /// Start from the SDK defaults. + pub fn new() -> Self { + Self::default() + } + + /// Cap concurrent requests per pooled connection. + /// + /// The cap includes unary requests and streaming sessions until their + /// bodies finish or are dropped. Additional requests use another pooled + /// connection. + /// + /// # Errors + /// + /// Must be between 1 and [`Self::MAX_CONCURRENT_REQUESTS`]. Requests beyond + /// the server's stream limit would queue on the connection instead of + /// spilling over to another one. + pub fn with_max_concurrent_requests( + self, + max_concurrent_requests: usize, + ) -> Result { + if !(1..=Self::MAX_CONCURRENT_REQUESTS).contains(&max_concurrent_requests) { + return Err(format!( + "HTTP/2 concurrent request limit must be between 1 and {}", + Self::MAX_CONCURRENT_REQUESTS + ) + .into()); + } + Ok(Self { + max_concurrent_requests: Some(max_concurrent_requests), + ..self + }) + } + + /// Set the initial per-stream receive window in bytes. + /// + /// # Errors + /// + /// Must be between 1 and [`Self::MAX_RECEIVE_WINDOW`]. + pub fn with_stream_receive_window( + self, + stream_receive_window: u32, + ) -> Result { + if !(1..=Self::MAX_RECEIVE_WINDOW).contains(&stream_receive_window) { + return Err("HTTP/2 stream receive window must be between 1 and 2^31 - 1 bytes".into()); + } + Ok(Self { + stream_receive_window: Some(stream_receive_window), + ..self + }) + } + + /// Set the connection-level receive window in bytes, shared by all streams + /// on the connection. + /// + /// # Errors + /// + /// Must be between [`Self::MIN_CONNECTION_RECEIVE_WINDOW`] and + /// [`Self::MAX_RECEIVE_WINDOW`]. + pub fn with_connection_receive_window( + self, + connection_receive_window: u32, + ) -> Result { + if !(Self::MIN_CONNECTION_RECEIVE_WINDOW..=Self::MAX_RECEIVE_WINDOW) + .contains(&connection_receive_window) + { + return Err( + "HTTP/2 connection receive window must be between 65,535 and 2^31 - 1 bytes".into(), + ); + } + Ok(Self { + connection_receive_window: Some(connection_receive_window), + ..self + }) + } +} + #[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: Http2Config, pub(crate) connection_timeout: Duration, pub(crate) request_timeout: Duration, pub(crate) retry: RetryConfig, @@ -524,6 +627,7 @@ impl S2Config { Self { access_token: AccessToken::Static(access_token.into().into()), endpoints: S2Endpoints::for_cloud(), + http2: Http2Config::new(), connection_timeout: Duration::from_secs(3), request_timeout: Duration::from_secs(5), retry: RetryConfig::new(), @@ -599,6 +703,13 @@ impl S2Config { }) } + /// Override HTTP/2 transport settings for pooled connections. + /// + /// Defaults to [`Http2Config::new()`]. + pub fn with_http2(self, http2: Http2Config) -> Self { + Self { http2, ..self } + } + /// Set the timeout for establishing a connection to the server. /// /// Defaults to `3s`. @@ -4557,4 +4668,54 @@ mod tests { assert_eq!(record.headers[0].value.as_ref(), b"v"); assert_eq!(record.timestamp, 1234); } + + #[test] + fn http2_config_bounds() { + let config = Http2Config::new() + .with_max_concurrent_requests(Http2Config::MAX_CONCURRENT_REQUESTS) + .unwrap() + .with_stream_receive_window(Http2Config::MAX_RECEIVE_WINDOW) + .unwrap() + .with_connection_receive_window(Http2Config::MIN_CONNECTION_RECEIVE_WINDOW) + .unwrap(); + assert_eq!( + config.max_concurrent_requests, + Some(Http2Config::MAX_CONCURRENT_REQUESTS) + ); + assert_eq!( + config.stream_receive_window, + Some(Http2Config::MAX_RECEIVE_WINDOW) + ); + assert_eq!( + config.connection_receive_window, + Some(Http2Config::MIN_CONNECTION_RECEIVE_WINDOW) + ); + + let partial = Http2Config::new().with_stream_receive_window(1).unwrap(); + assert_eq!(partial.max_concurrent_requests, None); + assert_eq!(partial.connection_receive_window, None); + + assert!(Http2Config::new().with_max_concurrent_requests(0).is_err()); + assert!( + Http2Config::new() + .with_max_concurrent_requests(Http2Config::MAX_CONCURRENT_REQUESTS + 1) + .is_err() + ); + assert!(Http2Config::new().with_stream_receive_window(0).is_err()); + assert!( + Http2Config::new() + .with_stream_receive_window(Http2Config::MAX_RECEIVE_WINDOW + 1) + .is_err() + ); + assert!( + Http2Config::new() + .with_connection_receive_window(Http2Config::MIN_CONNECTION_RECEIVE_WINDOW - 1) + .is_err() + ); + assert!( + Http2Config::new() + .with_connection_receive_window(u32::MAX) + .is_err() + ); + } } diff --git a/testcontainers/CHANGELOG.md b/testcontainers/CHANGELOG.md index 6d8b7d6d..ff6e0b55 100644 --- a/testcontainers/CHANGELOG.md +++ b/testcontainers/CHANGELOG.md @@ -2,6 +2,10 @@ All notable changes to this project will be documented in this file. +## [0.43.1] - 2026-09-28 + + + ## [0.43.0] - 2026-09-25 diff --git a/testcontainers/Cargo.toml b/testcontainers/Cargo.toml index 23647de1..9d2ed490 100644 --- a/testcontainers/Cargo.toml +++ b/testcontainers/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "s2-testcontainers" -version = "0.43.0" +version = "0.43.1" description = "Testcontainers helpers for the S2 Docker image" edition.workspace = true license.workspace = true