From 3ddade24f2eb7a36375213ac7e6339634a51410b Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 8 Jul 2026 18:42:09 +0800 Subject: [PATCH] fix(obs): validate numeric env settings (#4474) Co-authored-by: heihutu --- crates/obs/src/cleaner/core.rs | 47 ++++++++++++++++++++++- crates/obs/src/telemetry/dial9.rs | 64 ++++++++++++++++++++++++++++--- crates/obs/src/telemetry/local.rs | 33 ++++++++++++++-- crates/obs/src/telemetry/otel.rs | 28 +++++++++++++- crates/utils/src/envs.rs | 31 +++++++++++++-- 5 files changed, 189 insertions(+), 14 deletions(-) diff --git a/crates/obs/src/cleaner/core.rs b/crates/obs/src/cleaner/core.rs index 72577aa6d..fefaf6732 100644 --- a/crates/obs/src/cleaner/core.rs +++ b/crates/obs/src/cleaner/core.rs @@ -40,6 +40,12 @@ use tracing::{debug, error, info, warn}; const LOG_COMPONENT_OBS: &str = "obs"; const LOG_SUBSYSTEM_LOG_CLEANER: &str = "log_cleaner"; const EVENT_LOG_CLEANER_STATE: &str = "log_cleaner_state"; +const SECONDS_PER_DAY: u64 = 24 * 60 * 60; +const MAX_RETENTION_DAYS_BEFORE_SATURATION: u64 = u64::MAX / SECONDS_PER_DAY; + +fn compressed_file_retention_window(days: u64) -> Duration { + Duration::from_secs(days.saturating_mul(SECONDS_PER_DAY)) +} #[derive(Debug)] struct CompressionTaskResult { @@ -231,7 +237,18 @@ impl LogCleaner { /// Select compressed archives whose age exceeds the archive retention window. fn select_expired_compressed(&self, files: &mut [FileInfo]) -> Vec { - let retention = Duration::from_secs(self.compressed_file_retention_days * 24 * 3600); + if self.compressed_file_retention_days > MAX_RETENTION_DAYS_BEFORE_SATURATION { + warn!( + event = EVENT_LOG_CLEANER_STATE, + component = LOG_COMPONENT_OBS, + subsystem = LOG_SUBSYSTEM_LOG_CLEANER, + result = "retention_days_saturated", + configured_days = self.compressed_file_retention_days, + fallback_days = MAX_RETENTION_DAYS_BEFORE_SATURATION, + "log cleaner state changed" + ); + } + let retention = compressed_file_retention_window(self.compressed_file_retention_days); let now = SystemTime::now(); let mut expired = Vec::new(); @@ -785,3 +802,31 @@ impl LogCleanerBuilder { } } } + +#[cfg(test)] +mod tests { + use super::{MAX_RETENTION_DAYS_BEFORE_SATURATION, SECONDS_PER_DAY, compressed_file_retention_window}; + use std::time::Duration; + + #[test] + fn compressed_file_retention_window_scales_days_without_wrap() { + assert_eq!(compressed_file_retention_window(3), Duration::from_secs(3 * SECONDS_PER_DAY)); + } + + #[test] + fn compressed_file_retention_window_saturates_on_large_values() { + assert_eq!(compressed_file_retention_window(u64::MAX), Duration::from_secs(u64::MAX)); + } + + #[test] + fn retention_day_saturation_boundary_is_safe() { + assert_eq!( + compressed_file_retention_window(MAX_RETENTION_DAYS_BEFORE_SATURATION), + Duration::from_secs(MAX_RETENTION_DAYS_BEFORE_SATURATION * SECONDS_PER_DAY) + ); + assert_eq!( + compressed_file_retention_window(MAX_RETENTION_DAYS_BEFORE_SATURATION.saturating_add(1)), + Duration::from_secs(u64::MAX) + ); + } +} diff --git a/crates/obs/src/telemetry/dial9.rs b/crates/obs/src/telemetry/dial9.rs index ddd59b4b3..110a9d684 100644 --- a/crates/obs/src/telemetry/dial9.rs +++ b/crates/obs/src/telemetry/dial9.rs @@ -33,9 +33,9 @@ use rustfs_config::{ use rustfs_utils::get_env_bool; use rustfs_utils::get_env_f64; use rustfs_utils::get_env_opt_str; +use rustfs_utils::get_env_opt_u64; +use rustfs_utils::get_env_opt_usize; use rustfs_utils::get_env_str; -use rustfs_utils::get_env_u64; -use rustfs_utils::get_env_usize; use std::path::PathBuf; use tracing::{info, warn}; @@ -43,6 +43,18 @@ const LOG_COMPONENT_OBS: &str = "obs"; const LOG_SUBSYSTEM_DIAL9: &str = "dial9"; const EVENT_DIAL9_STATE: &str = "dial9_state"; +fn sanitize_dial9_max_file_size(bytes: u64) -> u64 { + bytes.max(1) +} + +fn sanitize_dial9_rotation_count(count: usize) -> usize { + count.max(1) +} + +fn dial9_total_rotation_size(max_file_size: u64, rotation_count: usize) -> u64 { + max_file_size.saturating_mul(u64::try_from(rotation_count).unwrap_or(u64::MAX)) +} + /// Configuration for dial9 Tokio telemetry. #[derive(Debug, Clone)] pub struct Dial9Config { @@ -95,12 +107,37 @@ impl Dial9Config { return Self::default(); } + let raw_max_file_size = get_env_opt_u64(ENV_RUNTIME_DIAL9_MAX_FILE_SIZE); + let raw_rotation_count = get_env_opt_usize(ENV_RUNTIME_DIAL9_ROTATION_COUNT); + let max_file_size = sanitize_dial9_max_file_size(raw_max_file_size.unwrap_or(DEFAULT_RUNTIME_DIAL9_MAX_FILE_SIZE)); + let rotation_count = sanitize_dial9_rotation_count(raw_rotation_count.unwrap_or(DEFAULT_RUNTIME_DIAL9_ROTATION_COUNT)); + if raw_max_file_size == Some(0) { + warn!( + event = EVENT_DIAL9_STATE, + component = LOG_COMPONENT_OBS, + subsystem = LOG_SUBSYSTEM_DIAL9, + result = "invalid_max_file_size", + fallback_bytes = 1_u64, + "dial9 state changed" + ); + } + if raw_rotation_count == Some(0) { + warn!( + event = EVENT_DIAL9_STATE, + component = LOG_COMPONENT_OBS, + subsystem = LOG_SUBSYSTEM_DIAL9, + result = "invalid_rotation_count", + fallback_count = 1_usize, + "dial9 state changed" + ); + } + Self { enabled, output_dir: get_env_str(ENV_RUNTIME_DIAL9_OUTPUT_DIR, DEFAULT_RUNTIME_DIAL9_OUTPUT_DIR), file_prefix: get_env_str(ENV_RUNTIME_DIAL9_FILE_PREFIX, DEFAULT_RUNTIME_DIAL9_FILE_PREFIX), - max_file_size: get_env_u64(ENV_RUNTIME_DIAL9_MAX_FILE_SIZE, DEFAULT_RUNTIME_DIAL9_MAX_FILE_SIZE), - rotation_count: get_env_usize(ENV_RUNTIME_DIAL9_ROTATION_COUNT, DEFAULT_RUNTIME_DIAL9_ROTATION_COUNT), + max_file_size, + rotation_count, s3_bucket: get_env_opt_str(ENV_RUNTIME_DIAL9_S3_BUCKET).filter(|s| !s.is_empty()), s3_prefix: get_env_opt_str(ENV_RUNTIME_DIAL9_S3_PREFIX).filter(|s| !s.is_empty()), sampling_rate: get_env_f64(ENV_RUNTIME_DIAL9_SAMPLING_RATE, DEFAULT_RUNTIME_DIAL9_SAMPLING_RATE).clamp(0.0, 1.0), @@ -282,8 +319,12 @@ pub fn build_traced_runtime( // Create rotating writer (synchronous for runtime building) let base_path = config.base_path(); - let writer = RotatingWriter::new(base_path, config.max_file_size, config.max_file_size * config.rotation_count as u64) - .map_err(|e| TelemetryError::Io(format!("Failed to create RotatingWriter: {}", e)))?; + let writer = RotatingWriter::new( + base_path, + config.max_file_size, + dial9_total_rotation_size(config.max_file_size, config.rotation_count), + ) + .map_err(|e| TelemetryError::Io(format!("Failed to create RotatingWriter: {}", e)))?; // Build traced runtime // Note: sampling_rate and S3 upload settings are reserved for future use @@ -319,6 +360,17 @@ mod tests { assert_eq!(config.base_path(), PathBuf::from("/tmp/telemetry/rustfs")); } + #[test] + fn test_dial9_sanitizers_reject_zero_values() { + assert_eq!(sanitize_dial9_max_file_size(0), 1); + assert_eq!(sanitize_dial9_rotation_count(0), 1); + } + + #[test] + fn test_dial9_total_rotation_size_saturates() { + assert_eq!(dial9_total_rotation_size(u64::MAX, 2), u64::MAX); + } + #[test] fn test_is_enabled_default() { // Skip if environment variable is explicitly set diff --git a/crates/obs/src/telemetry/local.rs b/crates/obs/src/telemetry/local.rs index 28a7a4da3..08a97f39b 100644 --- a/crates/obs/src/telemetry/local.rs +++ b/crates/obs/src/telemetry/local.rs @@ -74,6 +74,25 @@ const STDERR_WARNING_PREFIX: &str = "[WARN]"; const REQUEST_ID_CANONICAL: &str = "request_id"; const REQUEST_ID_COMPAT: &str = "request-id"; +fn resolve_log_cleanup_interval_seconds(config: &OtelConfig) -> u64 { + match config.log_cleanup_interval_seconds { + Some(0) => { + warn!( + event = EVENT_LOG_CLEANER_STATE, + component = LOG_COMPONENT_OBS, + subsystem = LOG_SUBSYSTEM_LOCAL_LOGGING, + result = "invalid_cleanup_interval", + configured_seconds = 0_u64, + fallback_seconds = DEFAULT_OBS_LOG_CLEANUP_INTERVAL_SECONDS, + "log cleaner state changed" + ); + DEFAULT_OBS_LOG_CLEANUP_INTERVAL_SECONDS + } + Some(seconds) => seconds, + None => DEFAULT_OBS_LOG_CLEANUP_INTERVAL_SECONDS, + } +} + #[derive(Clone, Debug)] struct RequestIdJsonFormat { inner: Format, @@ -549,9 +568,7 @@ pub fn spawn_cleanup_task( .log_min_file_age_seconds .unwrap_or(DEFAULT_OBS_LOG_MIN_FILE_AGE_SECONDS); let dry_run = config.log_dry_run.unwrap_or(DEFAULT_OBS_LOG_DRY_RUN); - let cleanup_interval = config - .log_cleanup_interval_seconds - .unwrap_or(DEFAULT_OBS_LOG_CLEANUP_INTERVAL_SECONDS); + let cleanup_interval = resolve_log_cleanup_interval_seconds(config); let cleaner = Arc::new( LogCleaner::builder(log_dir, file_pattern, active_filename) @@ -873,4 +890,14 @@ mod tests { && span.get("request_id").and_then(Value::as_str) == Some("req-parent") })); } + + #[test] + fn test_resolve_log_cleanup_interval_seconds_rejects_zero() { + let config = OtelConfig { + log_cleanup_interval_seconds: Some(0), + ..OtelConfig::default() + }; + + assert_eq!(resolve_log_cleanup_interval_seconds(&config), DEFAULT_OBS_LOG_CLEANUP_INTERVAL_SECONDS); + } } diff --git a/crates/obs/src/telemetry/otel.rs b/crates/obs/src/telemetry/otel.rs index 59373b3a7..5a121570c 100644 --- a/crates/obs/src/telemetry/otel.rs +++ b/crates/obs/src/telemetry/otel.rs @@ -424,7 +424,7 @@ fn build_meter_provider( .build() .map_err(|e| TelemetryError::BuildMetricExporter(e.to_string()))?; - let meter_interval = config.meter_interval.unwrap_or(METER_INTERVAL); + let meter_interval = resolve_meter_interval(config); let (provider, recorder) = Recorder::builder(service_name.to_string()) .with_meter_provider(|b: opentelemetry_sdk::metrics::MeterProviderBuilder| { @@ -608,6 +608,22 @@ fn create_periodic_reader(interval: u64) -> PeriodicReader u64 { + match config.meter_interval { + Some(0) => { + warn!( + result = "invalid_meter_interval", + configured_seconds = 0_u64, + fallback_seconds = METER_INTERVAL, + "Metrics export interval is invalid; using default interval" + ); + METER_INTERVAL + } + Some(interval) => interval, + None => METER_INTERVAL, + } +} + fn resolve_signal_headers(common_headers: Option<&str>, signal_headers: Option<&str>) -> HashMap { let mut headers = HashMap::new(); if let Some(raw_headers) = common_headers { @@ -698,6 +714,16 @@ mod tests { assert_eq!(resolve_signal_timeout(None, Some(0)), None); } + #[test] + fn test_resolve_meter_interval_rejects_zero() { + let config = OtelConfig { + meter_interval: Some(0), + ..OtelConfig::default() + }; + + assert_eq!(resolve_meter_interval(&config), METER_INTERVAL); + } + #[test] fn test_get_object_duration_histogram_metric_match_is_scoped() { assert!(is_get_object_duration_histogram_metric("rustfs_io_get_object_stage_duration_seconds")); diff --git a/crates/utils/src/envs.rs b/crates/utils/src/envs.rs index a1c6645bc..05e09223b 100644 --- a/crates/utils/src/envs.rs +++ b/crates/utils/src/envs.rs @@ -13,6 +13,7 @@ // limitations under the License. use std::{ + any::type_name, collections::BTreeSet, collections::HashSet, env, @@ -264,7 +265,15 @@ fn parse_env_value(key: &str) -> Option where T: std::str::FromStr, { - resolve_env_with_aliases(key, &[]).and_then(|(_, value)| value.parse().ok()) + let (used_key, value) = resolve_env_with_aliases(key, &[])?; + value + .parse::() + .map_err(|_| { + log_once(&format!("env_invalid_value:{used_key}"), || { + format!("Invalid {} value for {used_key}: {value}. Treating as unset.", type_name::()) + }); + }) + .ok() } pub fn get_env_str_with_aliases(key: &str, deprecated: &[&str], default: &str) -> String { @@ -681,8 +690,8 @@ pub fn apply_external_env_compat() -> ExternalEnvCompatReport { #[cfg(test)] mod tests { use super::{ - apply_external_env_compat, build_external_env_compat_report_from_entries, get_env_bool_with_aliases, - get_env_i32_with_aliases, get_env_str, + apply_external_env_compat, build_external_env_compat_report_from_entries, get_env_bool_with_aliases, get_env_f64, + get_env_i32_with_aliases, get_env_opt_f64, get_env_opt_u64, get_env_str, get_env_u64, }; fn source_key(suffix: &str) -> String { @@ -806,6 +815,22 @@ mod tests { }); } + #[test] + fn invalid_u64_value_falls_back_to_default_and_optional_none() { + temp_env::with_var("RUSTFS_TEST_U64", Some("not-a-u64"), || { + assert_eq!(get_env_u64("RUSTFS_TEST_U64", 42), 42); + assert_eq!(get_env_opt_u64("RUSTFS_TEST_U64"), None); + }); + } + + #[test] + fn invalid_f64_value_falls_back_to_default_and_optional_none() { + temp_env::with_var("RUSTFS_TEST_F64", Some("not-a-f64"), || { + assert_eq!(get_env_f64("RUSTFS_TEST_F64", 0.25), 0.25); + assert_eq!(get_env_opt_f64("RUSTFS_TEST_F64"), None); + }); + } + #[test] fn apply_external_env_compat_copies_missing_rustfs_keys() { temp_env::with_var("MINIO_ROOT_USER", Some("compat-admin"), || {