Fix connection lifecycle, pool bounds, and log rotation (review #007)

Critical:
- C1: Set server-side idle/keep-alive timeouts on both hyper builders
  (http2 keep_alive_interval=15s + keep_alive_timeout, http1
  header_read_timeout). Both builders now set TokioTimer (required to
  avoid runtime panic). Prevents FD exhaustion from abandoned TLS
  connections — the root cause of the 2026-07-24 outage.
- C2: Add Semaphore(max_connections) gating the accept loop. Provides
  backpressure via OS TCP backlog when all permits are taken.

Warnings:
- W1: Add SIGUSR1 log-reopen handler. New ReopenableFileWriter
  (Arc<ArcSwap<File>> via custom MakeWriter) atomically swaps the log
  file. Enables postrotate logrotate without copytruncate, which caused
  the 1.15GB sparse file that wedged fail2ban.
- W2: Set pool_max_idle_per_host(10) on both upstream clients, bounding
  idle upstream connections per host.
- W3: Add connection_idle_timeout_secs to StaticConfig (default 60).
- W4: Add max_connections to StaticConfig (default 1024).

Both new fields are validated (> 0) and included in static config drift
detection on reload. Docs (config.md, README, ADR-009) updated.
This commit is contained in:
glm-5.2 committed 2026-07-28 10:16:26 +00:00
1 parent e803817350
commit 0885486028
14 files changed
+407 -41

No files matched your search

+2
View File
@@ -115,6 +115,8 @@ Configuration uses TOML and is split into **static** (requires restart) and
| `health_check_port` | `9900` | Local health check port (`0` to disable) |
| `admin_key_path` | `/etc/reverse-proxy/admin-key` | Path to admin Bearer token file (empty string to disable) |
| `shutdown_timeout_secs` | `30` | Graceful shutdown timeout |
| `connection_idle_timeout_secs` | `60` | Server-side idle timeout for client TLS connections (prevents FD exhaustion from abandoned connections) |
| `max_connections` | `1024` | Max concurrent client TLS connections (backpressure via semaphore) |
| `logging.level` | `"info"` | Log level |
| `logging.format` | `"text"` | Log format (`"text"` or `"json"`) |
| `logging.log_file_path` | (not set) | Path to log file for fail2ban |
+16 -3
View File
@@ -90,6 +90,8 @@ Immutable after startup. Changes require a process restart.
| `health_check_port` | `u16` | Port for local health check endpoint (default: `9900`; set to `0` to disable; bound to `127.0.0.1` only; see ADR-013, ADR-022) |
| `admin_key_path` | `String` | Path to file containing the admin Bearer token (default: `/etc/reverse-proxy/admin-key`; empty string to disable admin endpoints; see ADR-028) |
| `shutdown_timeout_secs` | `u64` | Maximum seconds to wait for in-flight requests during graceful shutdown (default: `30`) |
| `connection_idle_timeout_secs` | `u64` | Server-side idle timeout for client TLS connections. Idle HTTP/2 connections are closed after this duration (with keep-alive pings at 15s intervals to detect dead peers). HTTP/1.1 connections are closed if the client doesn't send a complete request header within this duration. Prevents FD exhaustion from abandoned connections (default: `60`; must be > 0; see review #007 C1) |
| `max_connections` | `usize` | Maximum number of concurrent client TLS connections. When the limit is reached, new connections wait in the OS TCP backlog until a slot frees (default: `1024`; must be > 0; see review #007 C2) |
| `logging` | `LoggingConfig` | Logging configuration (see below) |
**LoggingConfig** (nested in `[logging]` TOML section):
@@ -105,9 +107,11 @@ This is critical for fail2ban regex matching and Docker log output (see ADR-024)
Both text and JSON formats produce plain-text output without color codes.
**Note**: The entire `LoggingConfig` (including `log_file_path`) is static and
requires a process restart to change. Log file path changes require reopening
file handles, which is complex and low-value for Phase 1. Log rotation (Phase 2)
will be handled via signal-based or built-in rotation.
requires a process restart to change. However, the log file can be reopened
without restarting by sending `SIGUSR1` to the process — this closes the
current file handle and opens a new one at the same path. This enables
standard `postrotate` logrotate configs (rename + signal) without the
`copytruncate` workaround that creates sparse files. See review #007 W1.
**ListenerConfig** (per-listener static config):
@@ -180,6 +184,8 @@ Phase 2.
| `health_check_port` | `u16` | `9900` | No |
| `admin_key_path` | `String` | `/etc/reverse-proxy/admin-key` | No |
| `shutdown_timeout_secs` | `u64` | `30` | No |
| `connection_idle_timeout_secs` | `u64` | `60` | No |
| `max_connections` | `usize` | `1024` | No |
| `logging.level` | `String` | `"info"` | No |
| `logging.format` | `String` | `"text"` | No |
| `logging.log_file_path` | `String` | (not set) | No |
@@ -306,6 +312,8 @@ certificate:
# Global settings
health_check_port = 9900 # Local health check (0 to disable)
admin_key_path = "/etc/reverse-proxy/admin-key" # Empty string to disable
# connection_idle_timeout_secs = 60 # Server-side idle timeout (default: 60)
# max_connections = 1024 # Max concurrent TLS connections (default: 1024)
[logging]
level = "info"
@@ -448,6 +456,11 @@ On startup, the config is validated:
20. `admin_key_path` must be either an empty string (disabled) or an absolute
path. Relative paths and paths containing `..` are rejected. This prevents
path traversal attacks on the admin key file.
21. `connection_idle_timeout_secs` must be > 0. A zero value would disable
the server-side idle timeout, reintroducing the FD exhaustion bug from
review #007 C1.
22. `max_connections` must be > 0. A zero value would deadlock the connection
semaphore, preventing any client connection from being accepted.
On SIGHUP reload, the same validation applies. If the new config fails
validation, the reload is rejected and the old config remains active. An error
@@ -10,6 +10,8 @@ The proxy needs to handle Unix signals for:
- **Graceful shutdown**: SIGTERM and SIGINT should stop accepting new
connections, drain in-flight requests, then exit.
- **Config reload**: SIGHUP should trigger a DynamicConfig reload from disk.
- **Log reopen**: SIGUSR1 should close and reopen the log file, enabling
`postrotate` logrotate configs without `copytruncate` (see review #007 W1).
Two approaches for signal handling:
- **`tokio::signal`**: Built into tokio. Handles SIGTERM and SIGINT via
@@ -22,6 +24,7 @@ Two approaches for signal handling:
Use `signal-hook` for all signal handling. Specifically:
- `signal-hook::flag` to set termination flags on SIGTERM/SIGINT
- `signal-hook` to register a SIGHUP handler that triggers config reload
- `signal-hook` to register a SIGUSR1 handler that reopens the log file
`tokio::signal::ctrl_c()` is registered as a secondary shutdown trigger; both
mechanisms converge on the same shutdown path. This is a belt-and-suspenders
@@ -34,6 +37,10 @@ The shutdown sequence:
for in-flight requests to complete, then exit with code 0.
2. On SIGHUP: re-read config file, validate, and swap DynamicConfig if valid.
Log the result.
3. On SIGUSR1: close the current log file handle and open a new one at the
same path. Enables standard `postrotate` logrotate configs (rename + signal)
without `copytruncate`, which creates sparse files when the FD offset is
high. See review #007 W1.
## Rationale
+6
View File
@@ -186,6 +186,12 @@ fn diff_static_config(old: &StaticConfig, new: &StaticConfig) -> Vec<String> {
if old.shutdown_timeout_secs != new.shutdown_timeout_secs {
changes.push("shutdown_timeout_secs".to_string());
}
if old.connection_idle_timeout_secs != new.connection_idle_timeout_secs {
changes.push("connection_idle_timeout_secs".to_string());
}
if old.max_connections != new.max_connections {
changes.push("max_connections".to_string());
}
if old.logging != new.logging {
changes.push("logging".to_string());
}
+6
View File
@@ -49,6 +49,10 @@ pub struct FullConfig {
pub admin_key_path: String,
#[serde(default = "static_config::default_shutdown_timeout_secs")]
pub shutdown_timeout_secs: u64,
#[serde(default = "static_config::default_connection_idle_timeout_secs")]
pub connection_idle_timeout_secs: u64,
#[serde(default = "static_config::default_max_connections")]
pub max_connections: usize,
#[serde(default)]
pub logging: LoggingConfig,
pub rate_limit: RateLimitConfig,
@@ -67,6 +71,8 @@ impl FullConfig {
health_check_port: self.health_check_port,
admin_key_path: self.admin_key_path,
shutdown_timeout_secs: self.shutdown_timeout_secs,
connection_idle_timeout_secs: self.connection_idle_timeout_secs,
max_connections: self.max_connections,
logging: self.logging,
};
let dynamic_config = DynamicConfig::from_sites(
+24
View File
@@ -11,6 +11,10 @@ pub struct StaticConfig {
pub admin_key_path: String,
#[serde(default = "default_shutdown_timeout_secs")]
pub shutdown_timeout_secs: u64,
#[serde(default = "default_connection_idle_timeout_secs")]
pub connection_idle_timeout_secs: u64,
#[serde(default = "default_max_connections")]
pub max_connections: usize,
#[serde(default)]
pub logging: LoggingConfig,
}
@@ -27,6 +31,14 @@ pub fn default_shutdown_timeout_secs() -> u64 {
30
}
pub fn default_connection_idle_timeout_secs() -> u64 {
60
}
pub fn default_max_connections() -> usize {
1024
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
pub struct ListenerConfig {
pub bind_addr: String,
@@ -199,6 +211,10 @@ upstream = "127.0.0.1:8080"
admin_key_path: String,
#[serde(default = "default_shutdown_timeout_secs")]
shutdown_timeout_secs: u64,
#[serde(default = "default_connection_idle_timeout_secs")]
connection_idle_timeout_secs: u64,
#[serde(default = "default_max_connections")]
max_connections: usize,
#[serde(default)]
logging: LoggingConfig,
rate_limit: crate::config::dynamic_config::RateLimitConfig,
@@ -215,6 +231,8 @@ upstream = "127.0.0.1:8080"
assert_eq!(config.health_check_port, 9900);
assert_eq!(config.admin_key_path, "/etc/reverse-proxy/admin-key");
assert_eq!(config.shutdown_timeout_secs, 30);
assert_eq!(config.connection_idle_timeout_secs, 60);
assert_eq!(config.max_connections, 1024);
assert_eq!(config.logging.level, "info");
assert_eq!(config.logging.format, "text");
assert!(config.logging.log_file_path.is_none());
@@ -295,6 +313,10 @@ acme_cache_dir = "/tmp/cache"
admin_key_path: String,
#[serde(default = "default_shutdown_timeout_secs")]
shutdown_timeout_secs: u64,
#[serde(default = "default_connection_idle_timeout_secs")]
connection_idle_timeout_secs: u64,
#[serde(default = "default_max_connections")]
max_connections: usize,
#[serde(default)]
logging: LoggingConfig,
rate_limit: crate::config::dynamic_config::RateLimitConfig,
@@ -307,6 +329,8 @@ acme_cache_dir = "/tmp/cache"
assert_eq!(config.health_check_port, 9900);
assert_eq!(config.admin_key_path, "/etc/reverse-proxy/admin-key");
assert_eq!(config.shutdown_timeout_secs, 30);
assert_eq!(config.connection_idle_timeout_secs, 60);
assert_eq!(config.max_connections, 1024);
assert_eq!(config.logging.level, "info");
assert_eq!(config.logging.format, "text");
assert!(config.logging.log_file_path.is_none());
+2
View File
@@ -22,6 +22,8 @@ pub fn test_static_config() -> StaticConfig {
health_check_port: 9900,
admin_key_path: "/etc/reverse-proxy/admin-key".to_string(),
shutdown_timeout_secs: 30,
connection_idle_timeout_secs: 60,
max_connections: 1024,
logging: LoggingConfig::default(),
}
}
+48
View File
@@ -77,6 +77,10 @@ pub enum ValidationError {
AdminKeyPathNotAbsolute { path: String },
#[error("admin_key_path must not contain '..' path traversal: '{path}'")]
AdminKeyPathTraversal { path: String },
#[error("connection_idle_timeout_secs must be > 0, got {value}")]
ConnectionIdleTimeoutZero { value: u64 },
#[error("max_connections must be > 0, got {value}")]
MaxConnectionsZero { value: usize },
}
pub fn validate(
@@ -282,6 +286,18 @@ pub fn validate(
}
}
if static_config.connection_idle_timeout_secs == 0 {
errors.push(ValidationError::ConnectionIdleTimeoutZero {
value: static_config.connection_idle_timeout_secs,
});
}
if static_config.max_connections == 0 {
errors.push(ValidationError::MaxConnectionsZero {
value: static_config.max_connections,
});
}
if errors.is_empty() {
Ok(())
} else {
@@ -376,6 +392,8 @@ mod tests {
health_check_port: 9900,
admin_key_path: "/etc/reverse-proxy/admin-key".to_string(),
shutdown_timeout_secs: 30,
connection_idle_timeout_secs: 60,
max_connections: 1024,
logging: LoggingConfig::default(),
}
}
@@ -412,6 +430,8 @@ mod tests {
health_check_port: 9900,
admin_key_path: "/etc/reverse-proxy/admin-key".to_string(),
shutdown_timeout_secs: 30,
connection_idle_timeout_secs: 60,
max_connections: 1024,
logging: LoggingConfig::default(),
}
}
@@ -1100,6 +1120,8 @@ mod tests {
health_check_port: 9900,
admin_key_path: "/etc/reverse-proxy/admin-key".to_string(),
shutdown_timeout_secs: 30,
connection_idle_timeout_secs: 60,
max_connections: 1024,
logging: LoggingConfig::default(),
};
let mut dynamic = valid_dynamic_config();
@@ -1304,4 +1326,30 @@ mod tests {
.iter()
.any(|e| matches!(e, ValidationError::AdminKeyPathTraversal { .. })));
}
#[test]
fn rule_connection_idle_timeout_zero_rejected() {
let mut config = valid_static_config();
config.connection_idle_timeout_secs = 0;
let dynamic = valid_dynamic_config();
let result = validate(&config, &dynamic, false);
assert!(result.is_err());
let errors = result.unwrap_err();
assert!(errors
.iter()
.any(|e| matches!(e, ValidationError::ConnectionIdleTimeoutZero { value: 0 })));
}
#[test]
fn rule_max_connections_zero_rejected() {
let mut config = valid_static_config();
config.max_connections = 0;
let dynamic = valid_dynamic_config();
let result = validate(&config, &dynamic, false);
assert!(result.is_err());
let errors = result.unwrap_err();
assert!(errors
.iter()
.any(|e| matches!(e, ValidationError::MaxConnectionsZero { value: 0 })));
}
}
+96 -27
View File
@@ -1,16 +1,19 @@
pub mod format;
pub mod reopen;
use crate::config::static_config::LoggingConfig;
use anyhow::Result;
use std::fs::File;
use std::sync::Arc;
use tracing::Level;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
use tracing_subscriber::EnvFilter;
use tracing_subscriber::Layer;
pub fn init(config: &LoggingConfig) -> Result<()> {
pub struct LogInit {
pub reopen_handle: Option<reopen::LogReopenHandle>,
}
pub fn init(config: &LoggingConfig) -> Result<LogInit> {
let level = config.level.parse::<Level>().unwrap_or(Level::INFO);
let env_filter = make_env_filter(level);
@@ -25,16 +28,16 @@ fn make_env_filter(level: Level) -> EnvFilter {
EnvFilter::from_default_env().add_directive(level.into())
}
fn init_json(env_filter: EnvFilter, log_file_path: &Option<String>, level: Level) -> Result<()> {
fn init_json(
env_filter: EnvFilter,
log_file_path: &Option<String>,
level: Level,
) -> Result<LogInit> {
match log_file_path {
Some(path) => {
if let Some(parent) = std::path::Path::new(path).parent() {
if !parent.as_os_str().is_empty() {
std::fs::create_dir_all(parent)?;
}
}
let file = File::create(path)?;
let file_writer = Arc::new(file);
let path_std = std::path::Path::new(path);
let file_writer = reopen::ReopenableFileWriter::new(path_std)?;
let reopen_handle = file_writer.handle_with_path(path_std.to_path_buf());
let file_env_filter = make_env_filter(level);
let stdout_layer = tracing_subscriber::fmt::layer()
@@ -50,29 +53,37 @@ fn init_json(env_filter: EnvFilter, log_file_path: &Option<String>, level: Level
.with(stdout_layer)
.with(file_layer)
.try_init()?;
Ok(LogInit {
reopen_handle: Some(reopen_handle),
})
}
None => {
let layer = tracing_subscriber::fmt::layer()
.json()
.with_ansi(false)
.with_filter(env_filter);
tracing_subscriber::registry().with(layer).try_init()?;
}
}
tracing_subscriber::registry()
.with(layer)
.try_init()?;
Ok(())
Ok(LogInit {
reopen_handle: None,
})
}
}
}
fn init_text(env_filter: EnvFilter, log_file_path: &Option<String>, level: Level) -> Result<()> {
fn init_text(
env_filter: EnvFilter,
log_file_path: &Option<String>,
level: Level,
) -> Result<LogInit> {
match log_file_path {
Some(path) => {
if let Some(parent) = std::path::Path::new(path).parent() {
if !parent.as_os_str().is_empty() {
std::fs::create_dir_all(parent)?;
}
}
let file = File::create(path)?;
let file_writer = Arc::new(file);
let path_std = std::path::Path::new(path);
let file_writer = reopen::ReopenableFileWriter::new(path_std)?;
let reopen_handle = file_writer.handle_with_path(path_std.to_path_buf());
let file_env_filter = make_env_filter(level);
let stdout_layer = tracing_subscriber::fmt::layer()
@@ -86,16 +97,24 @@ fn init_text(env_filter: EnvFilter, log_file_path: &Option<String>, level: Level
.with(stdout_layer)
.with(file_layer)
.try_init()?;
Ok(LogInit {
reopen_handle: Some(reopen_handle),
})
}
None => {
let layer = tracing_subscriber::fmt::layer()
.with_ansi(false)
.with_filter(env_filter);
tracing_subscriber::registry().with(layer).try_init()?;
}
}
tracing_subscriber::registry()
.with(layer)
.try_init()?;
Ok(())
Ok(LogInit {
reopen_handle: None,
})
}
}
}
#[cfg(test)]
@@ -122,4 +141,54 @@ mod tests {
}
assert!(log_path.exists(), "log file should be created");
}
#[test]
fn init_returns_reopen_handle_when_file_configured() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("access.log");
let config = LoggingConfig {
level: "info".to_string(),
format: "text".to_string(),
log_file_path: Some(log_path.to_string_lossy().to_string()),
};
let result = init(&config);
match result {
Ok(init_result) => {
assert!(init_result.reopen_handle.is_some());
}
Err(e) => {
let msg = format!("{e}");
assert!(
msg.contains("global default trace dispatcher")
|| msg.contains("already been set"),
"unexpected init error: {e}"
);
}
}
}
#[test]
fn init_returns_no_reopen_handle_when_no_file() {
let config = LoggingConfig {
level: "info".to_string(),
format: "text".to_string(),
log_file_path: None,
};
let result = init(&config);
match result {
Ok(init_result) => {
assert!(init_result.reopen_handle.is_none());
}
Err(e) => {
let msg = format!("{e}");
assert!(
msg.contains("global default trace dispatcher")
|| msg.contains("already been set"),
"unexpected init error: {e}"
);
}
}
}
}
+129
View File
@@ -0,0 +1,129 @@
use std::fs::{File, OpenOptions};
use std::io::{self, Write};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use arc_swap::ArcSwap;
use tracing_subscriber::fmt::writer::MakeWriter;
pub struct ReopenableFileWriter {
file: Arc<ArcSwap<File>>,
}
impl ReopenableFileWriter {
pub fn new(path: &Path) -> io::Result<Self> {
if let Some(parent) = path.parent() {
if !parent.as_os_str().is_empty() {
std::fs::create_dir_all(parent)?;
}
}
let file = File::create(path)?;
Ok(Self {
file: Arc::new(ArcSwap::from_pointee(file)),
})
}
pub fn handle(&self) -> LogReopenHandle {
LogReopenHandle {
file: self.file.clone(),
path: PathBuf::new(),
}
}
pub fn handle_with_path(&self, path: PathBuf) -> LogReopenHandle {
LogReopenHandle {
file: self.file.clone(),
path,
}
}
}
impl<'a> MakeWriter<'a> for ReopenableFileWriter {
type Writer = ReopenableFileWriterHandle;
fn make_writer(&'a self) -> Self::Writer {
ReopenableFileWriterHandle {
file: self.file.load_full(),
}
}
}
pub struct ReopenableFileWriterHandle {
file: Arc<File>,
}
impl Write for ReopenableFileWriterHandle {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.file.write(buf)
}
fn flush(&mut self) -> io::Result<()> {
self.file.flush()
}
}
pub struct LogReopenHandle {
file: Arc<ArcSwap<File>>,
path: PathBuf,
}
impl LogReopenHandle {
pub fn reopen(&self) -> io::Result<()> {
if self.path.as_os_str().is_empty() {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"log reopen path not configured",
));
}
let new_file = OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(&self.path)?;
self.file.store(Arc::new(new_file));
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Read;
#[test]
fn reopen_swaps_underlying_file() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("access.log");
let writer = ReopenableFileWriter::new(&log_path).unwrap();
let handle = writer.handle_with_path(log_path.clone());
let mut w = writer.make_writer();
w.write_all(b"first").unwrap();
w.flush().unwrap();
std::fs::write(&log_path, "ROTATED").unwrap();
handle.reopen().unwrap();
let mut w = writer.make_writer();
w.write_all(b"second").unwrap();
w.flush().unwrap();
let mut content = String::new();
File::open(&log_path)
.unwrap()
.read_to_string(&mut content)
.unwrap();
assert_eq!(content, "second");
}
#[test]
fn reopen_without_path_returns_error() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("access.log");
let writer = ReopenableFileWriter::new(&log_path).unwrap();
let handle = writer.handle();
let result = handle.reopen();
assert!(result.is_err());
}
}
+7 -1
View File
@@ -60,7 +60,8 @@ fn main() {
}
async fn run_server(loaded_config: cli::LoadedConfig, config_path: &str) -> Result<()> {
logging::init(&loaded_config.static_config.logging).context("failed to initialize logging")?;
let log_init = logging::init(&loaded_config.static_config.logging)
.context("failed to initialize logging")?;
info!("reverse-proxy starting");
@@ -92,6 +93,7 @@ async fn run_server(loaded_config: cli::LoadedConfig, config_path: &str) -> Resu
shutdown.clone(),
reload_handle.clone(),
config_path.to_string(),
log_init.reopen_handle,
)?;
let admin_auth = if !loaded_config.static_config.admin_key_path.is_empty() {
@@ -229,6 +231,10 @@ async fn run_server(loaded_config: cli::LoadedConfig, config_path: &str) -> Resu
app.clone(),
shutdown_rx,
in_flight.clone(),
std::time::Duration::from_secs(
loaded_config.static_config.connection_idle_timeout_secs,
),
loaded_config.static_config.max_connections,
));
info!(
+3
View File
@@ -228,12 +228,14 @@ fn build_upstream_request(req: Request<Body>, upstream_uri: &Uri) -> anyhow::Res
}
const CONNECT_TIMEOUT_CEILING_SECS: u64 = 30;
const POOL_MAX_IDLE_PER_HOST: usize = 10;
pub fn create_http_client() -> Client<HttpConnector, Body> {
let mut connector = HttpConnector::new();
connector.set_connect_timeout(Some(Duration::from_secs(CONNECT_TIMEOUT_CEILING_SECS)));
Client::builder(TokioExecutor::new())
.pool_idle_timeout(Duration::from_secs(90))
.pool_max_idle_per_host(POOL_MAX_IDLE_PER_HOST)
.build(connector)
}
@@ -254,6 +256,7 @@ pub fn create_https_client() -> Client<hyper_rustls::HttpsConnector<HttpConnecto
Client::builder(TokioExecutor::new())
.pool_idle_timeout(Duration::from_secs(90))
.pool_max_idle_per_host(POOL_MAX_IDLE_PER_HOST)
.build(https_connector)
}
+34 -8
View File
@@ -1,6 +1,7 @@
use std::net::SocketAddr;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Duration;
use axum::extract::ConnectInfo;
use axum::http::Request;
@@ -10,10 +11,13 @@ use hyper::body::Incoming;
use hyper_util::rt::TokioExecutor;
use hyper_util::service::TowerToHyperService;
use tokio::net::TcpListener;
use tokio::sync::Semaphore;
use tokio_rustls::TlsAcceptor;
use tower::Service;
use tracing::{error, info, warn};
const HTTP2_KEEP_ALIVE_INTERVAL_SECS: u64 = 15;
pub struct InFlightCounter {
count: AtomicUsize,
}
@@ -59,8 +63,11 @@ pub async fn serve_https_listener(
router: Router,
mut shutdown_rx: tokio::sync::watch::Receiver<bool>,
in_flight: Arc<InFlightCounter>,
connection_idle_timeout: Duration,
max_connections: usize,
) {
let local_addr = tcp_listener.local_addr();
let conn_sem = Arc::new(Semaphore::new(max_connections));
loop {
tokio::select! {
@@ -76,9 +83,19 @@ pub async fn serve_https_listener(
let tls_acceptor = tls_acceptor.clone();
let router = router.clone();
let in_flight = in_flight.clone();
let conn_sem = conn_sem.clone();
let permit = match conn_sem.acquire_owned().await {
Ok(permit) => permit,
Err(e) => {
error!(error = %e, "connection semaphore closed");
continue;
}
};
tokio::spawn(async move {
let _guard = InFlightGuard::new(in_flight.clone());
let _permit = permit;
let tls_stream = match tls_acceptor.accept(tcp_stream).await {
Ok(stream) => stream,
@@ -101,19 +118,28 @@ pub async fn serve_https_listener(
if is_h2 {
let mut builder = hyper::server::conn::http2::Builder::new(TokioExecutor::new());
if let Err(e) = builder
.enable_connect_protocol()
.serve_connection(io, svc)
.await
builder
.timer(hyper_util::rt::TokioTimer::new())
.keep_alive_interval(Some(Duration::from_secs(HTTP2_KEEP_ALIVE_INTERVAL_SECS)))
.keep_alive_timeout(connection_idle_timeout)
.enable_connect_protocol();
if let Err(e) = builder.serve_connection(io, svc).await
{
error!(error = %e, "HTTPS/2 connection error");
}
} else {
let mut builder = hyper_util::server::conn::auto::Builder::new(TokioExecutor::new());
builder.http2().enable_connect_protocol();
if let Err(e) = builder
.serve_connection_with_upgrades(io, svc)
.await
builder
.http1()
.timer(hyper_util::rt::TokioTimer::new())
.header_read_timeout(Some(connection_idle_timeout));
builder
.http2()
.timer(hyper_util::rt::TokioTimer::new())
.keep_alive_interval(Some(Duration::from_secs(HTTP2_KEEP_ALIVE_INTERVAL_SECS)))
.keep_alive_timeout(connection_idle_timeout)
.enable_connect_protocol();
if let Err(e) = builder.serve_connection_with_upgrades(io, svc).await
{
if let Some(hyper_err) = e.downcast_ref::<hyper::Error>() {
if hyper_err.is_incomplete_message() {
+27 -2
View File
@@ -2,10 +2,12 @@ use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;
use signal_hook::consts::{SIGHUP, SIGINT, SIGTERM};
use signal_hook::consts::{SIGHUP, SIGINT, SIGTERM, SIGUSR1};
use signal_hook::iterator::Signals;
use tokio::sync::watch;
use crate::logging::reopen::LogReopenHandle;
pub struct GracefulShutdown {
shutdown_timeout: Duration,
shutdown_tx: watch::Sender<bool>,
@@ -47,8 +49,9 @@ pub fn register_signal_handlers(
shutdown: Arc<GracefulShutdown>,
reload_handle: Arc<crate::config::ConfigReloadHandle>,
config_path: String,
log_reopen_handle: Option<LogReopenHandle>,
) -> anyhow::Result<()> {
let mut signals = Signals::new([SIGTERM, SIGINT, SIGHUP])?;
let mut signals = Signals::new([SIGTERM, SIGINT, SIGHUP, SIGUSR1])?;
let (tx, mut rx) = tokio::sync::mpsc::channel::<i32>(16);
std::thread::spawn(move || {
@@ -59,6 +62,8 @@ pub fn register_signal_handlers(
}
});
let reopen_handle = log_reopen_handle.map(Arc::new);
tokio::spawn(async move {
while let Some(sig) = rx.recv().await {
match sig {
@@ -71,6 +76,10 @@ pub fn register_signal_handlers(
tracing::info!(event = "SIGNAL", signal = "SIGHUP");
handle_sighup_reload(&reload_handle, &config_path).await;
}
SIGUSR1 => {
tracing::info!(event = "SIGNAL", signal = "SIGUSR1");
handle_sigusr1_log_reopen(reopen_handle.as_ref()).await;
}
_ => {
tracing::debug!(event = "SIGNAL", signal = %sig);
}
@@ -117,6 +126,22 @@ pub async fn handle_sighup_reload(
}
}
pub async fn handle_sigusr1_log_reopen(reopen_handle: Option<&Arc<LogReopenHandle>>) {
match reopen_handle {
Some(handle) => match handle.reopen() {
Ok(()) => {
tracing::info!(event = "LOG_REOPEN", status = "success");
}
Err(e) => {
tracing::error!(event = "LOG_REOPEN", status = "error", error = %e);
}
},
None => {
tracing::warn!(event = "LOG_REOPEN", status = "skipped", reason = "no_log_file");
}
}
}
#[cfg(test)]
mod tests {
use super::*;