diff --git a/CHANGELOG.md b/CHANGELOG.md index 83831de..01067f5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,10 @@ ## Unreleased +### Added + +- Add `--http-pool-size`, `--http-pool-max-idle-time`, and `--http-pool-max-reuse` CLI flags. + ### Changed - Replace `uri` in `sandhole_http_elapsed_time` for `status_code`. diff --git a/book/src/cli.md b/book/src/cli.md index e8da017..3da592b 100644 --- a/book/src/cli.md +++ b/book/src/cli.md @@ -324,6 +324,30 @@ Expose HTTP/SSH/TCP services through SSH port forwarding. By default, connections are immediately timed out when the pool is exhausted. + --http-pool-size <SIZE> + Maximum size for each HTTP connection pool (per unique client). + + Controls how many keep-alive HTTP connections are cached per unique + client. Lower values reduce memory usage, higher values improve + performance. + + [default: 16] + + --http-pool-max-idle-time <DURATION> + Maximum idle time for HTTP pooled connections. + + Connections idle longer than this will be evicted from the pool. + + [default: 60s] + + --http-pool-max-reuse <COUNT> + Maximum number of times an HTTP connection can be reused. + + Prevents connections with accumulated state from persisting + indefinitely. + + [default: 1000] + --max-simultaneous-connections-per-ip <SIZE> Maximum number of simultaneous connections per IP to a proxied service. The maximum is 65535. diff --git a/src/config.rs b/src/config.rs index 25763a1..d34eec8 100644 --- a/src/config.rs +++ b/src/config.rs @@ -429,6 +429,28 @@ pub struct ApplicationConfig { #[arg(long, value_parser = validate_duration, value_name = "DURATION")] pub pool_timeout: Option, + /// Maximum size for each HTTP connection pool (per unique client). + /// + /// Controls how many keep-alive HTTP connections are cached + /// per unique client. Lower values reduce memory usage, + /// higher values improve performance. + #[arg(long, default_value_t = 16, value_name = "SIZE")] + pub http_pool_size: usize, + + /// Maximum idle time for HTTP pooled connections. + /// + /// Connections idle longer than this will be evicted + /// from the pool. + #[arg(long, default_value = "60s", value_parser = validate_duration, value_name = "DURATION")] + pub http_pool_max_idle_time: Duration, + + /// Maximum number of times an HTTP connection can be reused. + /// + /// Prevents connections with accumulated state from + /// persisting indefinitely. + #[arg(long, default_value_t = 1_000, value_name = "COUNT")] + pub http_pool_max_reuse: usize, + /// Maximum number of simultaneous connections per IP to a proxied service. /// The maximum is 65535. /// @@ -607,6 +629,9 @@ mod application_config_tests { buffer_size: 32_768, pool_size: 65_536, pool_timeout: None, + http_pool_size: 16, + http_pool_max_idle_time: Duration::from_secs(60), + http_pool_max_reuse: 1_000, max_simultaneous_connections_per_ip: 64, ssh_keepalive_interval: Duration::from_secs(15), ssh_keepalive_max: 3, @@ -671,6 +696,9 @@ mod application_config_tests { "--buffer-size=4KB", "--pool-size=2048", "--pool-timeout=9s", + "--http-pool-size=32", + "--http-pool-max-idle-time=120s", + "--http-pool-max-reuse=500", "--max-simultaneous-connections-per-ip=32", "--ssh-keepalive-interval=10s", "--ssh-keepalive-max=2", @@ -737,6 +765,9 @@ mod application_config_tests { buffer_size: 4_000, pool_size: 2_048, pool_timeout: Some(Duration::from_secs(9)), + http_pool_size: 32, + http_pool_max_idle_time: Duration::from_secs(120), + http_pool_max_reuse: 500, max_simultaneous_connections_per_ip: 32, ssh_keepalive_interval: Duration::from_secs(10), ssh_keepalive_max: 2, diff --git a/src/entrypoint.rs b/src/entrypoint.rs index 7d43a51..b46c716 100644 --- a/src/entrypoint.rs +++ b/src/entrypoint.rs @@ -668,6 +668,9 @@ pub async fn entrypoint(config: ApplicationConfig) -> color_eyre::Result<()> { // Always use aliasing channels instead of tunneling channels. .proxy_type(ProxyType::Aliasing) .has_pool_queue(config.pool_timeout.is_some()) + .http_pool_size(config.http_pool_size) + .http_pool_max_idle_time(config.http_pool_max_idle_time) + .http_pool_max_reuse(config.http_pool_max_reuse) .buffer_size(buffer_size) .maybe_http_request_timeout(http_request_timeout) .maybe_websocket_timeout(tcp_connection_timeout) @@ -774,6 +777,9 @@ pub async fn entrypoint(config: ApplicationConfig) -> color_eyre::Result<()> { // Always use tunneling channels. .proxy_type(ProxyType::Tunneling) .has_pool_queue(config.pool_timeout.is_some()) + .http_pool_size(config.http_pool_size) + .http_pool_max_idle_time(config.http_pool_max_idle_time) + .http_pool_max_reuse(config.http_pool_max_reuse) .buffer_size(buffer_size) .maybe_http_request_timeout(http_request_timeout) .maybe_websocket_timeout(tcp_connection_timeout) @@ -854,6 +860,9 @@ pub async fn entrypoint(config: ApplicationConfig) -> color_eyre::Result<()> { // Always use tunneling channels. .proxy_type(ProxyType::Tunneling) .has_pool_queue(config.pool_timeout.is_some()) + .http_pool_size(config.http_pool_size) + .http_pool_max_idle_time(config.http_pool_max_idle_time) + .http_pool_max_reuse(config.http_pool_max_reuse) .buffer_size(buffer_size) .maybe_http_request_timeout(http_request_timeout) .maybe_websocket_timeout(tcp_connection_timeout) diff --git a/src/http/http11.rs b/src/http/http11.rs index 67b4961..98ac881 100644 --- a/src/http/http11.rs +++ b/src/http/http11.rs @@ -6,8 +6,8 @@ use crate::{ connection_handler::ConnectionHandler, connections::ConnectionGetByHttpHost, http::{ - ArcProxyData, HttpError, HttpLog, ProxyData, ProxyResponse, ProxyType, TimedResponse, - http_log, + ArcProxyData, HttpError, HttpLog, PooledConnection, ProxyData, ProxyResponse, ProxyType, + TimedResponse, http_log, }, keepalive::KeepaliveAlias, }; @@ -226,9 +226,13 @@ where tokio::select! { result = recv => { match result { - // Return pool item - Ok(tuple) if tuple.0.is_ready() => break (tuple.0, tuple.1, guard), - // Connection is closed; discard pool item + // Return pool item if connection is ready and not expired + Ok(mut pooled) if pooled.connection.is_ready() + && !pooled.is_expired(guard.max_idle_time, guard.max_reuse) => { + pooled.touch(); + break (pooled.connection, pooled.log_sender, guard) + }, + // Connection is closed or expired; discard pool item Ok(_) => { recv = guard.pool.recv(); continue; @@ -256,9 +260,12 @@ where 'sender: loop { match proxy_data.get_http11_pool_guard(key.clone()) { Some(guard) => { - while let Ok(sender) = guard.pool.try_recv() { - if sender.0.is_ready() { - break 'sender (sender.0, sender.1, guard); + while let Ok(mut pooled) = guard.pool.try_recv() { + if pooled.connection.is_ready() + && !pooled.is_expired(guard.max_idle_time, guard.max_reuse) + { + pooled.touch(); + break 'sender (pooled.connection, pooled.log_sender, guard); } } } @@ -296,11 +303,12 @@ where // Create entry for pool let key_clone = key.clone(); + let pool_size = proxy_data.http_pool_size; let pool = { let pool_ref = proxy_data .keepalive_http11_pool_map .entry(key_clone) - .or_insert_with(|| Arc::new(async_channel::unbounded())) + .or_insert_with(|| Arc::new(async_channel::bounded(pool_size))) .downgrade(); pool_ref.0.clone() }; @@ -365,7 +373,8 @@ where ); // Return sender to pool tokio::spawn(async move { - let _ = pool.send((sender, tx)).await; + let pooled = PooledConnection::new(sender, tx); + let _ = pool.send(pooled).await; drop(_guard); }); })), diff --git a/src/http/http2.rs b/src/http/http2.rs index 0912fae..e1124a8 100644 --- a/src/http/http2.rs +++ b/src/http/http2.rs @@ -4,8 +4,8 @@ use crate::{ connection_handler::ConnectionHandler, connections::ConnectionGetByHttpHost, http::{ - ArcProxyData, HttpError, HttpLog, ProxyData, ProxyResponse, ProxyType, TimedResponse, - http_log, + ArcProxyData, HttpError, HttpLog, PooledConnection, ProxyData, ProxyResponse, ProxyType, + TimedResponse, http_log, }, keepalive::KeepaliveAlias, }; @@ -66,9 +66,13 @@ where tokio::select! { result = recv => { match result { - // Return pool item - Ok(tuple) if tuple.0.is_ready() => break (tuple.0, tuple.1, guard), - // Connection is closed; discard pool item + // Return pool item if connection is ready and not expired + Ok(mut pooled) if pooled.connection.is_ready() + && !pooled.is_expired(guard.max_idle_time, guard.max_reuse) => { + pooled.touch(); + break (pooled.connection, pooled.log_sender, guard) + }, + // Connection is closed or expired; discard pool item Ok(_) => { recv = guard.pool.recv(); continue; @@ -97,9 +101,12 @@ where 'sender: loop { match proxy_data.get_http2_pool_guard(key.clone()) { Some(guard) => { - while let Ok(sender) = guard.pool.try_recv() { - if sender.0.is_ready() { - break 'sender (sender.0, sender.1, guard); + while let Ok(mut pooled) = guard.pool.try_recv() { + if pooled.connection.is_ready() + && !pooled.is_expired(guard.max_idle_time, guard.max_reuse) + { + pooled.touch(); + break 'sender (pooled.connection, pooled.log_sender, guard); } } } @@ -140,11 +147,12 @@ where // Create entry for pool let key_clone = key.clone(); + let pool_size = proxy_data.http_pool_size; let pool = { let pool_ref = proxy_data .keepalive_http2_pool_map .entry(key_clone) - .or_insert_with(|| Arc::new(async_channel::unbounded())) + .or_insert_with(|| Arc::new(async_channel::bounded(pool_size))) .downgrade(); pool_ref.0.clone() }; @@ -222,7 +230,8 @@ where ); // Send sender to pool tokio::spawn(async move { - let _ = pool.send((sender, tx)).await; + let pooled = PooledConnection::new(sender, tx); + let _ = pool.send(pooled).await; drop(_guard); }); })), diff --git a/src/http/mod.rs b/src/http/mod.rs index fd834f4..8e2112b 100644 --- a/src/http/mod.rs +++ b/src/http/mod.rs @@ -339,9 +339,37 @@ impl IntoResponse for HttpError { } } +// Wrapper for pooled connections with metadata +pub(crate) struct PooledConnection { + connection: T, + log_sender: ServerHandlerSender, + last_used: Instant, + reuse_count: usize, +} + +impl PooledConnection { + fn new(connection: T, log_sender: ServerHandlerSender) -> Self { + Self { + connection, + log_sender, + last_used: Instant::now(), + reuse_count: 0, + } + } + + fn touch(&mut self) { + self.last_used = Instant::now(); + self.reuse_count += 1; + } + + fn is_expired(&self, max_idle: Duration, max_reuse: usize) -> bool { + self.last_used.elapsed() > max_idle || self.reuse_count >= max_reuse + } +} + type KeepalivePool = Arc<(Sender, Receiver)>; -type Http11KeepalivePool = KeepalivePool<(Http11SendRequest, ServerHandlerSender)>; -type Http2KeepalivePool = KeepalivePool<(Http2SendRequest, ServerHandlerSender)>; +type Http11KeepalivePool = KeepalivePool>>; +type Http2KeepalivePool = KeepalivePool>>; // Data commonly reused between HTTP proxy requests. #[derive(Builder)] @@ -362,6 +390,12 @@ where keepalive_http2_pool_map: Arc, RandomState>>, // Whether the server supports a queue pool for handlers. has_pool_queue: bool, + // Maximum size for each connection pool. + http_pool_size: usize, + // Maximum idle time for pooled connections. + http_pool_max_idle_time: Duration, + // Maximum number of times a connection can be reused. + http_pool_max_reuse: usize, // Connection manager to get handlers from. conn_manager: M, // Tuple containing where to redirect requests from the main domain to. @@ -399,9 +433,11 @@ where } pub(crate) struct Http11PoolGuard { - pool: Receiver<(Http11SendRequest, ServerHandlerSender)>, + pool: Receiver>>, key: KeepaliveAlias, map: Arc, RandomState>>, + max_idle_time: Duration, + max_reuse: usize, } impl Drop for Http11PoolGuard { @@ -413,9 +449,11 @@ impl Drop for Http11PoolGuard { } pub(crate) struct Http2PoolGuard { - pool: Receiver<(Http2SendRequest, ServerHandlerSender)>, + pool: Receiver>>, key: KeepaliveAlias, map: Arc, RandomState>>, + max_idle_time: Duration, + max_reuse: usize, } pub(crate) trait ArcProxyData { @@ -443,11 +481,12 @@ where ::Error: Error + Send + Sync + 'static, { fn create_http11_pool_guard(&self, key: KeepaliveAlias) -> Http11PoolGuard { + let pool_size = self.http_pool_size; let pool = { let map_ref = self .keepalive_http11_pool_map .entry(key.clone()) - .or_insert_with(|| Arc::new(async_channel::unbounded())) + .or_insert_with(|| Arc::new(async_channel::bounded(pool_size))) .downgrade(); map_ref.1.clone() }; @@ -455,6 +494,8 @@ where pool, key, map: Arc::clone(&self.keepalive_http11_pool_map), + max_idle_time: self.http_pool_max_idle_time, + max_reuse: self.http_pool_max_reuse, } } @@ -467,15 +508,18 @@ where pool, key, map: Arc::clone(&self.keepalive_http11_pool_map), + max_idle_time: self.http_pool_max_idle_time, + max_reuse: self.http_pool_max_reuse, }) } fn create_http2_pool_guard(&self, key: KeepaliveAlias) -> Http2PoolGuard { + let pool_size = self.http_pool_size; let pool = { let map_ref = self .keepalive_http2_pool_map .entry(key.clone()) - .or_insert_with(|| Arc::new(async_channel::unbounded())) + .or_insert_with(|| Arc::new(async_channel::bounded(pool_size))) .downgrade(); map_ref.1.clone() }; @@ -483,6 +527,8 @@ where pool, key, map: Arc::clone(&self.keepalive_http2_pool_map), + max_idle_time: self.http_pool_max_idle_time, + max_reuse: self.http_pool_max_reuse, } } @@ -495,6 +541,8 @@ where pool, key, map: Arc::clone(&self.keepalive_http2_pool_map), + max_idle_time: self.http_pool_max_idle_time, + max_reuse: self.http_pool_max_reuse, }) } } @@ -761,6 +809,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -807,6 +858,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -853,6 +907,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -919,6 +976,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -988,6 +1048,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1058,6 +1121,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1127,6 +1193,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1223,6 +1292,9 @@ mod proxy_handler_tests { .http_request_timeout(Duration::from_millis(500)) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1314,6 +1386,9 @@ mod proxy_handler_tests { .http_request_timeout(Duration::from_millis(500)) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1427,6 +1502,9 @@ mod proxy_handler_tests { .http_request_timeout(Duration::from_millis(500)) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1527,6 +1605,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1629,6 +1710,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1730,6 +1814,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1830,6 +1917,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1931,6 +2021,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2027,6 +2120,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2146,6 +2242,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2256,6 +2355,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2383,6 +2485,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2460,6 +2565,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) diff --git a/src/telemetry.rs b/src/telemetry.rs index e8f7289..aa1ba4d 100644 --- a/src/telemetry.rs +++ b/src/telemetry.rs @@ -63,23 +63,17 @@ impl SlidingWindowCounter { fn clean(&self) { let mut history = self.history.lock().expect("not poisoned"); - loop { - let Some(element) = history.front() else { + while let Some(element) = history.front() + && element.0.elapsed() >= self.window + { + // Don't remove the first element if it is the last one. + // Instead, update its instant. + // This ensures that the first call to add is counted. + if history.len() == 1 { + history.front_mut().expect("length check").0 = Instant::now(); break; - }; - // Remove elements at the front if they are too old. - if element.0.elapsed() >= self.window { - // Don't remove the first element if it is the last one. - // Instead, update its instant. - // This ensures that the first call to add is counted. - if history.len() == 1 { - history.front_mut().expect("length check").0 = Instant::now(); - break; - } else { - history.pop_front(); - } } else { - break; + history.pop_front(); } } }
= Arc<(Sender
, Receiver
)>; -type Http11KeepalivePool = KeepalivePool<(Http11SendRequest, ServerHandlerSender)>; -type Http2KeepalivePool = KeepalivePool<(Http2SendRequest, ServerHandlerSender)>; +type Http11KeepalivePool = KeepalivePool>>; +type Http2KeepalivePool = KeepalivePool>>; // Data commonly reused between HTTP proxy requests. #[derive(Builder)] @@ -362,6 +390,12 @@ where keepalive_http2_pool_map: Arc, RandomState>>, // Whether the server supports a queue pool for handlers. has_pool_queue: bool, + // Maximum size for each connection pool. + http_pool_size: usize, + // Maximum idle time for pooled connections. + http_pool_max_idle_time: Duration, + // Maximum number of times a connection can be reused. + http_pool_max_reuse: usize, // Connection manager to get handlers from. conn_manager: M, // Tuple containing where to redirect requests from the main domain to. @@ -399,9 +433,11 @@ where } pub(crate) struct Http11PoolGuard { - pool: Receiver<(Http11SendRequest, ServerHandlerSender)>, + pool: Receiver>>, key: KeepaliveAlias, map: Arc, RandomState>>, + max_idle_time: Duration, + max_reuse: usize, } impl Drop for Http11PoolGuard { @@ -413,9 +449,11 @@ impl Drop for Http11PoolGuard { } pub(crate) struct Http2PoolGuard { - pool: Receiver<(Http2SendRequest, ServerHandlerSender)>, + pool: Receiver>>, key: KeepaliveAlias, map: Arc, RandomState>>, + max_idle_time: Duration, + max_reuse: usize, } pub(crate) trait ArcProxyData { @@ -443,11 +481,12 @@ where ::Error: Error + Send + Sync + 'static, { fn create_http11_pool_guard(&self, key: KeepaliveAlias) -> Http11PoolGuard { + let pool_size = self.http_pool_size; let pool = { let map_ref = self .keepalive_http11_pool_map .entry(key.clone()) - .or_insert_with(|| Arc::new(async_channel::unbounded())) + .or_insert_with(|| Arc::new(async_channel::bounded(pool_size))) .downgrade(); map_ref.1.clone() }; @@ -455,6 +494,8 @@ where pool, key, map: Arc::clone(&self.keepalive_http11_pool_map), + max_idle_time: self.http_pool_max_idle_time, + max_reuse: self.http_pool_max_reuse, } } @@ -467,15 +508,18 @@ where pool, key, map: Arc::clone(&self.keepalive_http11_pool_map), + max_idle_time: self.http_pool_max_idle_time, + max_reuse: self.http_pool_max_reuse, }) } fn create_http2_pool_guard(&self, key: KeepaliveAlias) -> Http2PoolGuard { + let pool_size = self.http_pool_size; let pool = { let map_ref = self .keepalive_http2_pool_map .entry(key.clone()) - .or_insert_with(|| Arc::new(async_channel::unbounded())) + .or_insert_with(|| Arc::new(async_channel::bounded(pool_size))) .downgrade(); map_ref.1.clone() }; @@ -483,6 +527,8 @@ where pool, key, map: Arc::clone(&self.keepalive_http2_pool_map), + max_idle_time: self.http_pool_max_idle_time, + max_reuse: self.http_pool_max_reuse, } } @@ -495,6 +541,8 @@ where pool, key, map: Arc::clone(&self.keepalive_http2_pool_map), + max_idle_time: self.http_pool_max_idle_time, + max_reuse: self.http_pool_max_reuse, }) } } @@ -761,6 +809,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -807,6 +858,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -853,6 +907,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -919,6 +976,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -988,6 +1048,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1058,6 +1121,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1127,6 +1193,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1223,6 +1292,9 @@ mod proxy_handler_tests { .http_request_timeout(Duration::from_millis(500)) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1314,6 +1386,9 @@ mod proxy_handler_tests { .http_request_timeout(Duration::from_millis(500)) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1427,6 +1502,9 @@ mod proxy_handler_tests { .http_request_timeout(Duration::from_millis(500)) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1527,6 +1605,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1629,6 +1710,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1730,6 +1814,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1830,6 +1917,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -1931,6 +2021,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2027,6 +2120,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2146,6 +2242,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2256,6 +2355,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2383,6 +2485,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) @@ -2460,6 +2565,9 @@ mod proxy_handler_tests { .buffer_size(8_000) .disable_http_logs(false) .log_format(LogFormat::Default) + .http_pool_size(16) + .http_pool_max_idle_time(Duration::from_secs(60)) + .http_pool_max_reuse(1_000) .build(), ), ) diff --git a/src/telemetry.rs b/src/telemetry.rs index e8f7289..aa1ba4d 100644 --- a/src/telemetry.rs +++ b/src/telemetry.rs @@ -63,23 +63,17 @@ impl SlidingWindowCounter { fn clean(&self) { let mut history = self.history.lock().expect("not poisoned"); - loop { - let Some(element) = history.front() else { + while let Some(element) = history.front() + && element.0.elapsed() >= self.window + { + // Don't remove the first element if it is the last one. + // Instead, update its instant. + // This ensures that the first call to add is counted. + if history.len() == 1 { + history.front_mut().expect("length check").0 = Instant::now(); break; - }; - // Remove elements at the front if they are too old. - if element.0.elapsed() >= self.window { - // Don't remove the first element if it is the last one. - // Instead, update its instant. - // This ensures that the first call to add is counted. - if history.len() == 1 { - history.front_mut().expect("length check").0 = Instant::now(); - break; - } else { - history.pop_front(); - } } else { - break; + history.pop_front(); } } }