From 658b8dea666f0fb9eee17b8450a18833b81f7a08 Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 27 May 2026 02:49:42 +0800 Subject: [PATCH] fix: unify runtime readiness publication and graceful shutdown flow (#3087) * fix(server): unify runtime readiness and shutdown flow * refactor(server): decouple readiness shared types * refactor(server): keep readiness types crate-private * refactor(server): await protocol shutdown handles * fix(server): bypass cached startup readiness --- rustfs/src/admin/handlers/health.rs | 9 +- rustfs/src/app/admin_usecase.rs | 292 +------------------ rustfs/src/embedded.rs | 32 ++- rustfs/src/init.rs | 26 +- rustfs/src/main.rs | 105 +++---- rustfs/src/server/http.rs | 76 ++--- rustfs/src/server/mod.rs | 46 +++ rustfs/src/server/readiness.rs | 416 +++++++++++++++++++++++++++- rustfs/src/server/service_state.rs | 22 ++ 9 files changed, 582 insertions(+), 442 deletions(-) diff --git a/rustfs/src/admin/handlers/health.rs b/rustfs/src/admin/handlers/health.rs index c8d23e698..e4b55f88b 100644 --- a/rustfs/src/admin/handlers/health.rs +++ b/rustfs/src/admin/handlers/health.rs @@ -14,8 +14,10 @@ use super::profile::{TriggerProfileCPU, TriggerProfileMemory}; use crate::admin::router::{AdminOperation, Operation, S3Router}; -use crate::app::admin_usecase::DefaultAdminUsecase; -use crate::server::{HEALTH_PREFIX, HEALTH_READY_PATH, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH}; +use crate::server::{ + HEALTH_PREFIX, HEALTH_READY_PATH, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH, + collect_dependency_readiness as collect_runtime_dependency_readiness, +}; use http::{HeaderMap, HeaderValue}; use hyper::{Method, StatusCode}; use matchit::Params; @@ -56,8 +58,7 @@ pub(crate) enum HealthProbe { } pub(crate) async fn collect_dependency_readiness() -> (bool, bool) { - let usecase = DefaultAdminUsecase::from_global(); - let readiness = usecase.execute_collect_dependency_readiness().await; + let readiness = collect_runtime_dependency_readiness().await; (readiness.storage_ready, readiness.iam_ready) } diff --git a/rustfs/src/app/admin_usecase.rs b/rustfs/src/app/admin_usecase.rs index 778eff608..51a65276e 100644 --- a/rustfs/src/app/admin_usecase.rs +++ b/rustfs/src/app/admin_usecase.rs @@ -17,6 +17,7 @@ use crate::app::context::{AppContext, get_global_app_context}; use crate::capacity::resolve_admin_used_capacity; use crate::error::ApiError; +use crate::server::{DependencyReadiness, collect_dependency_readiness as collect_runtime_dependency_readiness}; use rustfs_data_usage::DataUsageInfo; use rustfs_ecstore::admin_server_info::get_server_info; use rustfs_ecstore::data_usage::{apply_bucket_usage_memory_overlay, load_data_usage_from_backend}; @@ -24,12 +25,9 @@ use rustfs_ecstore::endpoints::EndpointServerPools; use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::pools::{PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free}; use rustfs_ecstore::store_api::StorageAPI; -use rustfs_madmin::{Disk, InfoMessage, StorageInfo}; +use rustfs_madmin::{InfoMessage, StorageInfo}; use s3s::S3ErrorCode; -use std::collections::{HashMap, HashSet}; -use std::sync::{Arc, OnceLock}; -use std::time::{Duration, Instant}; -use tokio::sync::Mutex; +use std::sync::Arc; use tracing::{debug, error, info, warn}; pub type AdminUsecaseResult = Result; @@ -49,12 +47,6 @@ impl std::fmt::Debug for QueryServerInfoResponse { } } -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] -pub struct DependencyReadiness { - pub storage_ready: bool, - pub iam_ready: bool, -} - #[derive(Debug, Clone, Default, PartialEq, Eq)] pub struct QueryPoolStatusRequest { pub pool: String, @@ -88,16 +80,7 @@ pub struct DefaultAdminUsecase { context: Option>, } -#[derive(Debug, Clone, Copy)] -struct StorageReadinessCacheEntry { - captured_at: Instant, - storage_ready: bool, -} - impl DefaultAdminUsecase { - const DISK_STATE_OK: &'static str = "ok"; - const DISK_STATE_UNFORMATTED: &'static str = "unformatted"; - const RUNTIME_STATE_RETURNING: &'static str = "returning"; const POOL_STATUS_ACTIVE: &'static str = "active"; const POOL_STATUS_CANCELED: &'static str = "canceled"; const POOL_STATUS_COMPLETE: &'static str = "complete"; @@ -321,160 +304,8 @@ impl DefaultAdminUsecase { used_size as f64 / total_size as f64 } - fn disk_is_online_for_readiness(disk: &Disk) -> bool { - let state_is_acceptable = disk.state.eq_ignore_ascii_case(Self::DISK_STATE_OK) - || disk.state.eq_ignore_ascii_case(rustfs_madmin::ITEM_ONLINE) - || disk.state.eq_ignore_ascii_case(Self::DISK_STATE_UNFORMATTED); - - if let Some(runtime_state) = disk.runtime_state.as_deref() { - let runtime_state_is_acceptable = runtime_state.eq_ignore_ascii_case(rustfs_madmin::ITEM_ONLINE) - || runtime_state.eq_ignore_ascii_case(Self::RUNTIME_STATE_RETURNING); - return runtime_state_is_acceptable && state_is_acceptable; - } - - state_is_acceptable - } - - fn health_readiness_cache_ttl() -> Duration { - Duration::from_millis(rustfs_utils::get_env_u64( - rustfs_config::ENV_HEALTH_READINESS_CACHE_TTL_MS, - rustfs_config::DEFAULT_HEALTH_READINESS_CACHE_TTL_MS, - )) - } - - fn storage_readiness_cache() -> &'static Mutex> { - static CACHE: OnceLock>> = OnceLock::new(); - CACHE.get_or_init(|| Mutex::new(None)) - } - - async fn load_cached_storage_readiness() -> Option { - let ttl = Self::health_readiness_cache_ttl(); - if ttl.is_zero() { - return None; - } - - let cache = Self::storage_readiness_cache().lock().await; - let entry = cache.as_ref()?; - if entry.captured_at.elapsed() <= ttl { - return Some(entry.storage_ready); - } - - None - } - - async fn update_storage_readiness_cache(storage_ready: bool) { - if Self::health_readiness_cache_ttl().is_zero() { - return; - } - - let mut cache = Self::storage_readiness_cache().lock().await; - *cache = Some(StorageReadinessCacheEntry { - captured_at: Instant::now(), - storage_ready, - }); - } - - fn pool_write_quorum(info: &StorageInfo, pool_idx: usize, set_drive_count: usize) -> usize { - if set_drive_count == 0 { - return 1; - } - - let data_drives = info - .backend - .standard_sc_data - .get(pool_idx) - .copied() - .filter(|count| *count > 0) - .unwrap_or_else(|| (set_drive_count / 2).max(1)); - - let parity_drives = if let Some(drives_per_set) = info.backend.drives_per_set.get(pool_idx).copied() { - drives_per_set.saturating_sub(data_drives) - } else if let Some(parity) = info.backend.standard_sc_parities.get(pool_idx).copied() { - parity - } else if let Some(parity) = info.backend.standard_sc_parity { - parity - } else { - set_drive_count.saturating_sub(data_drives) - }; - - let mut write_quorum = data_drives; - if data_drives == parity_drives { - write_quorum += 1; - } - write_quorum.max(1) - } - - fn storage_ready_from_runtime_state(info: &StorageInfo) -> bool { - if info.disks.is_empty() { - return false; - } - - let mut total_online = 0usize; - let mut set_online_counts: HashMap<(usize, usize), usize> = HashMap::new(); - let mut set_drive_counts: HashMap<(usize, usize), usize> = HashMap::new(); - let mut seen_disks: HashSet<(String, String, i32, i32, i32)> = HashSet::new(); - - for disk in &info.disks { - if disk.pool_index < 0 || disk.set_index < 0 { - continue; - } - - let dedup_key = ( - disk.endpoint.clone(), - disk.drive_path.clone(), - disk.pool_index, - disk.set_index, - disk.disk_index, - ); - if !seen_disks.insert(dedup_key) { - continue; - } - - let pool_idx = disk.pool_index as usize; - let set_idx = disk.set_index as usize; - let key = (pool_idx, set_idx); - *set_drive_counts.entry(key).or_default() += 1; - - if Self::disk_is_online_for_readiness(disk) { - total_online += 1; - *set_online_counts.entry(key).or_default() += 1; - } - } - - if total_online == 0 { - return false; - } - - if set_drive_counts.is_empty() { - return false; - } - - set_drive_counts.into_iter().all(|((pool_idx, set_idx), set_drive_count)| { - let online = set_online_counts.get(&(pool_idx, set_idx)).copied().unwrap_or_default(); - let write_quorum = Self::pool_write_quorum(info, pool_idx, set_drive_count); - online >= write_quorum - }) - } - pub async fn execute_collect_dependency_readiness(&self) -> DependencyReadiness { - let iam_ready = self.context.as_ref().map(|context| context.iam().is_ready()).unwrap_or(false); - let storage_ready = if let Some(cached) = Self::load_cached_storage_readiness().await { - cached - } else { - let computed = if let Some(store) = new_object_layer_fn() { - let storage_info = store.storage_info().await; - Self::storage_ready_from_runtime_state(&storage_info) - } else { - false - }; - Self::update_storage_readiness_cache(computed).await; - computed - }; - - DependencyReadiness { - storage_ready, - iam_ready: iam_ready && storage_ready, - } + collect_runtime_dependency_readiness().await } } @@ -509,121 +340,6 @@ mod tests { let _ = readiness.iam_ready; } - #[test] - fn storage_ready_from_runtime_state_returns_false_when_all_disks_faulty() { - let info = StorageInfo { - backend: rustfs_madmin::BackendInfo { - standard_sc_data: vec![1], - drives_per_set: vec![1], - ..Default::default() - }, - disks: vec![Disk { - pool_index: 0, - set_index: 0, - state: "offline".to_string(), - runtime_state: Some("offline".to_string()), - ..Default::default() - }], - }; - - assert!(!DefaultAdminUsecase::storage_ready_from_runtime_state(&info)); - } - - #[test] - fn storage_ready_from_runtime_state_returns_true_when_set_meets_write_quorum() { - let info = StorageInfo { - backend: rustfs_madmin::BackendInfo { - standard_sc_data: vec![1], - drives_per_set: vec![1], - ..Default::default() - }, - disks: vec![Disk { - pool_index: 0, - set_index: 0, - state: "ok".to_string(), - runtime_state: Some("online".to_string()), - ..Default::default() - }], - }; - - assert!(DefaultAdminUsecase::storage_ready_from_runtime_state(&info)); - } - - #[test] - fn storage_ready_from_runtime_state_deduplicates_duplicate_disk_rows() { - let duplicate_disk = Disk { - endpoint: "127.0.0.1:9000".to_string(), - drive_path: "/data0".to_string(), - pool_index: 0, - set_index: 0, - disk_index: 0, - state: "ok".to_string(), - runtime_state: Some("online".to_string()), - ..Default::default() - }; - let info = StorageInfo { - backend: rustfs_madmin::BackendInfo { - standard_sc_data: vec![2], - drives_per_set: vec![4], - ..Default::default() - }, - disks: vec![duplicate_disk.clone(), duplicate_disk], - }; - - assert!( - !DefaultAdminUsecase::storage_ready_from_runtime_state(&info), - "duplicate rows must not satisfy write quorum" - ); - } - - #[test] - fn disk_online_for_readiness_requires_runtime_and_state_both_acceptable() { - let disk = Disk { - state: "disk io error".to_string(), - runtime_state: Some("online".to_string()), - ..Default::default() - }; - assert!(!DefaultAdminUsecase::disk_is_online_for_readiness(&disk)); - } - - #[test] - fn storage_ready_from_runtime_state_requires_all_sets_meet_quorum() { - let info = StorageInfo { - backend: rustfs_madmin::BackendInfo { - standard_sc_data: vec![1], - drives_per_set: vec![2], - ..Default::default() - }, - disks: vec![ - Disk { - endpoint: "127.0.0.1:9000".to_string(), - drive_path: "/set0d0".to_string(), - pool_index: 0, - set_index: 0, - disk_index: 0, - state: "ok".to_string(), - runtime_state: Some("online".to_string()), - ..Default::default() - }, - Disk { - endpoint: "127.0.0.1:9000".to_string(), - drive_path: "/set1d0".to_string(), - pool_index: 0, - set_index: 1, - disk_index: 0, - state: "offline".to_string(), - runtime_state: Some("offline".to_string()), - ..Default::default() - }, - ], - }; - - assert!( - !DefaultAdminUsecase::storage_ready_from_runtime_state(&info), - "if any set fails write quorum, readiness must be false" - ); - } - #[test] fn admin_pool_list_item_maps_capacity_and_active_status() { let now = OffsetDateTime::UNIX_EPOCH; diff --git a/rustfs/src/embedded.rs b/rustfs/src/embedded.rs index 4a55a1f09..6f047cdee 100644 --- a/rustfs/src/embedded.rs +++ b/rustfs/src/embedded.rs @@ -49,7 +49,10 @@ use crate::app::context::{AppContext, init_global_app_context}; use crate::config::Config; use crate::init::{add_bucket_notification_configuration, init_buffer_profile_system, init_kms_system}; -use crate::server::{init_event_notifier, shutdown_event_notifier, start_audit_system, start_http_server, stop_audit_system}; +use crate::server::{ + ShutdownHandle, init_event_notifier, publish_ready_when_runtime_ready, shutdown_event_notifier, start_audit_system, + start_http_server, stop_audit_system, +}; use crate::startup_fs_guard::enforce_unsupported_fs_policy; use rustfs_common::{GlobalReadiness, SystemStage, set_global_addr}; use rustfs_config::ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS; @@ -368,10 +371,10 @@ impl RustFSServerBuilder { // Start HTTP server. let mut s3_config = config.clone(); s3_config.console_enable = false; - let (shutdown_tx, bound_addr) = start_http_server(&s3_config, readiness.clone()).await?; + let (shutdown_handle, bound_addr) = start_http_server(&s3_config, readiness.clone()).await?; let ctx = CancellationToken::new(); let shutdown_embedded_server = || { - let _ = shutdown_tx.send(()); + shutdown_handle.signal(); ctx.cancel(); }; @@ -464,8 +467,12 @@ impl RustFSServerBuilder { warn!("notification system: {e}"); } - // Mark fully ready. - readiness.mark_stage(SystemStage::FullReady); + publish_ready_when_runtime_ready(readiness.as_ref(), None) + .await + .map_err(|e| { + shutdown_embedded_server(); + ServerError::Init(format!("runtime readiness: {e}")) + })?; rustfs_common::set_global_init_time_now().await; let server = RustFSServer { @@ -473,7 +480,7 @@ impl RustFSServerBuilder { access_key: self.access_key.clone(), secret_key: self.secret_key.clone(), region: self.region.clone(), - shutdown_tx: Some(shutdown_tx), + shutdown_handle: Some(shutdown_handle), cancel_token: ctx, temp_dir: temp_dir_guard.map(|g| g.keep()), }; @@ -500,7 +507,7 @@ pub struct RustFSServer { access_key: String, secret_key: String, region: String, - shutdown_tx: Option>, + shutdown_handle: Option, cancel_token: CancellationToken, temp_dir: Option, } @@ -564,13 +571,10 @@ impl RustFSServer { } // Signal HTTP server to stop. - if let Some(tx) = self.shutdown_tx.take() { - let _ = tx.send(()); + if let Some(shutdown_handle) = self.shutdown_handle.take() { + shutdown_handle.shutdown().await; } - // Brief grace period for connections to drain. - tokio::time::sleep(std::time::Duration::from_secs(2)).await; - // Clean up temp directory if we created it. if let Some(ref dir) = self.temp_dir && let Err(e) = tokio::fs::remove_dir_all(dir).await @@ -586,8 +590,8 @@ impl Drop for RustFSServer { fn drop(&mut self) { // Best-effort synchronous cleanup. self.cancel_token.cancel(); - if let Some(tx) = self.shutdown_tx.take() { - let _ = tx.send(()); + if let Some(shutdown_handle) = self.shutdown_handle.take() { + shutdown_handle.signal(); } if let Some(ref dir) = self.temp_dir { let _ = std::fs::remove_dir_all(dir); diff --git a/rustfs/src/init.rs b/rustfs/src/init.rs index f42a8da5c..8eedf02b2 100644 --- a/rustfs/src/init.rs +++ b/rustfs/src/init.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use crate::server::ShutdownHandle; use crate::storage::{process_lambda_configurations, process_queue_configurations, process_topic_configurations}; use crate::{admin, config, version}; use rustfs_config::{ @@ -496,7 +497,7 @@ pub async fn init_auto_tuner(ctx: tokio_util::sync::CancellationToken) { /// This function initializes the FTP server (non-encrypted) if enabled in the configuration. #[cfg(feature = "ftps")] #[instrument(skip_all)] -pub async fn init_ftp_system() -> Result>, Box> { +pub async fn init_ftp_system() -> Result, Box> { { use crate::protocols::ProtocolStorageClient; use rustfs_config::{DEFAULT_FTP_ADDRESS, ENV_FTP_ADDRESS, ENV_FTP_ENABLE, ENV_FTP_EXTERNAL_IP, ENV_FTP_PASSIVE_PORTS}; @@ -547,7 +548,7 @@ pub async fn init_ftp_system() -> Result Result Result Result>, Box> { +pub async fn init_ftps_system() -> Result, Box> { { use crate::protocols::ProtocolStorageClient; use rustfs_config::{ @@ -623,7 +624,7 @@ pub async fn init_ftps_system() -> Result Result Result Result>, Box> -{ +pub async fn init_webdav_system() -> Result, Box> { { use crate::protocols::ProtocolStorageClient; use rustfs_config::{ @@ -693,7 +693,7 @@ pub async fn init_webdav_system() -> Result Result Result Result>, Box> { +pub async fn init_sftp_system() -> Result, Box> { { use crate::protocols::ProtocolStorageClient; use rustfs_config::{ @@ -777,7 +777,7 @@ pub async fn init_sftp_system() -> Result Result Result<()> { }; #[cfg(not(feature = "ftps"))] - let ftp_shutdown_tx: Option> = None; + let ftp_shutdown_tx: Option = None; // Initialize FTPS system if enabled #[cfg(feature = "ftps")] @@ -443,7 +443,7 @@ async fn run(config: rustfs::config::Config) -> Result<()> { }; #[cfg(not(feature = "ftps"))] - let ftps_shutdown_tx: Option> = None; + let ftps_shutdown_tx: Option = None; // Initialize WebDAV system if enabled #[cfg(feature = "webdav")] @@ -463,7 +463,7 @@ async fn run(config: rustfs::config::Config) -> Result<()> { }; #[cfg(not(feature = "webdav"))] - let webdav_shutdown_tx: Option> = None; + let webdav_shutdown_tx: Option = None; // Initialize SFTP system if enabled #[cfg(feature = "sftp")] @@ -483,7 +483,7 @@ async fn run(config: rustfs::config::Config) -> Result<()> { }; #[cfg(not(feature = "sftp"))] - let sftp_shutdown_tx: Option> = None; + let sftp_shutdown_tx: Option = None; // Initialize buffer profiling system init_buffer_profile_system(&config); @@ -607,55 +607,31 @@ async fn run(config: rustfs::config::Config) -> Result<()> { &server_address, jiff::Zoned::now() ); - // 4. Mark as Full Ready now that critical components are warm - readiness.mark_stage(SystemStage::FullReady); + publish_ready_when_runtime_ready(readiness.as_ref(), Some(&state_manager)).await?; // Set the global RustFS initialization time to now rustfs_common::set_global_init_time_now().await; - // Publish ready only after all critical bootstrap metadata is in place - state_manager.update(ServiceState::Ready); if enable_scanner { init_data_scanner(ctx.clone(), store.clone()).await; } - // Perform hibernation for 1 second - tokio::time::sleep(SHUTDOWN_TIMEOUT).await; // listen to the shutdown signal - match wait_for_shutdown().await { - #[cfg(unix)] - ShutdownSignal::CtrlC | ShutdownSignal::Sigint | ShutdownSignal::Sigterm => { - handle_shutdown( - &state_manager, - s3_shutdown_tx, - console_shutdown_tx, - ProtocolShutdownSenders { - ftp: ftp_shutdown_tx, - ftps: ftps_shutdown_tx, - webdav: webdav_shutdown_tx, - sftp: sftp_shutdown_tx, - }, - ctx.clone(), - ) - .await; - } - #[cfg(not(unix))] - ShutdownSignal::CtrlC => { - handle_shutdown( - &state_manager, - s3_shutdown_tx, - console_shutdown_tx, - ProtocolShutdownSenders { - ftp: ftp_shutdown_tx, - ftps: ftps_shutdown_tx, - webdav: webdav_shutdown_tx, - sftp: sftp_shutdown_tx, - }, - ctx.clone(), - ) - .await; - } - } + let shutdown_signal = wait_for_shutdown().await; + handle_shutdown( + &state_manager, + shutdown_signal, + s3_shutdown_tx, + console_shutdown_tx, + ProtocolShutdownSenders { + ftp: ftp_shutdown_tx, + ftps: ftps_shutdown_tx, + webdav: webdav_shutdown_tx, + sftp: sftp_shutdown_tx, + }, + ctx.clone(), + ) + .await; info!(target: "rustfs::main::run","server is stopped state: {:?}", state_manager.current_state()); Ok(()) @@ -664,17 +640,18 @@ async fn run(config: rustfs::config::Config) -> Result<()> { /// Shutdown channels for every protocol server. None means the protocol was /// disabled at startup. struct ProtocolShutdownSenders { - ftp: Option>, - ftps: Option>, - webdav: Option>, - sftp: Option>, + ftp: Option, + ftps: Option, + webdav: Option, + sftp: Option, } /// Handles the shutdown process of the server async fn handle_shutdown( state_manager: &ServiceStateManager, - s3_shutdown_tx: Option>, - console_shutdown_tx: Option>, + shutdown_signal: ShutdownSignal, + s3_shutdown_handle: Option, + console_shutdown_handle: Option, protocols: ProtocolShutdownSenders, ctx: CancellationToken, ) { @@ -688,6 +665,7 @@ async fn handle_shutdown( info!( target: "rustfs::main::handle_shutdown", + signal = shutdown_signal.log_label(), "Shutdown signal received in main thread" ); // update the status to stopping first @@ -722,12 +700,13 @@ async fn handle_shutdown( } // Shutdown FTP and FTPS servers + let mut protocol_shutdowns = Vec::new(); if let Some(ftp_shutdown_tx) = ftp_shutdown_tx { info!( target: "rustfs::main::handle_shutdown", "Shutting down FTP server..." ); - let _ = ftp_shutdown_tx.send(()); + protocol_shutdowns.push(ftp_shutdown_tx.shutdown()); } if let Some(ftps_shutdown_tx) = ftps_shutdown_tx { @@ -735,7 +714,7 @@ async fn handle_shutdown( target: "rustfs::main::handle_shutdown", "Shutting down FTPS server..." ); - let _ = ftps_shutdown_tx.send(()); + protocol_shutdowns.push(ftps_shutdown_tx.shutdown()); } // Shutdown WebDAV server @@ -744,7 +723,7 @@ async fn handle_shutdown( target: "rustfs::main::handle_shutdown", "Shutting down WebDAV server..." ); - let _ = webdav_shutdown_tx.send(()); + protocol_shutdowns.push(webdav_shutdown_tx.shutdown()); } // Shutdown SFTP server @@ -753,7 +732,7 @@ async fn handle_shutdown( target: "rustfs::main::handle_shutdown", "Shutting down SFTP server..." ); - let _ = sftp_shutdown_tx.send(()); + protocol_shutdowns.push(sftp_shutdown_tx.shutdown()); } // Stop the notification system @@ -784,15 +763,13 @@ async fn handle_shutdown( target: "rustfs::main::handle_shutdown", "Server is stopping..." ); - if let Some(s3_shutdown_tx) = s3_shutdown_tx { - let _ = s3_shutdown_tx.send(()); + if let Some(s3_shutdown_handle) = s3_shutdown_handle { + s3_shutdown_handle.shutdown().await; } - if let Some(console_shutdown_tx) = console_shutdown_tx { - let _ = console_shutdown_tx.send(()); + if let Some(console_shutdown_handle) = console_shutdown_handle { + console_shutdown_handle.shutdown().await; } - - // Wait for the worker thread to complete the cleaning work - tokio::time::sleep(SHUTDOWN_TIMEOUT).await; + join_all(protocol_shutdowns).await; // the last updated status is stopped state_manager.update(ServiceState::Stopped); diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index accf6bbed..0e92451b4 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -18,7 +18,7 @@ use crate::auth::IAMAuth; use crate::auth_keystone; use crate::config; use crate::server::{ - ReadinessGateLayer, RemoteAddr, + ReadinessGateLayer, RemoteAddr, ShutdownHandle, compress::{CompressionConfig, PathAwareCompressionPredicate, PathCategoryInjectionLayer}, hybrid::hybrid, layer::{ @@ -37,7 +37,7 @@ use http::{HeaderMap, Method, Request as HttpRequest, Response}; use hyper_util::{ rt::{TokioExecutor, TokioIo, TokioTimer}, server::conn::auto::Builder as ConnBuilder, - server::graceful::GracefulShutdown, + server::graceful::{GracefulShutdown, Watcher}, service::TowerToHyperService, }; use metrics::{counter, gauge, histogram}; @@ -128,10 +128,7 @@ pub(crate) fn active_http_requests() -> u64 { ACTIVE_HTTP_REQUESTS.load(Ordering::Relaxed) } -pub async fn start_http_server( - config: &config::Config, - readiness: Arc, -) -> Result<(tokio::sync::broadcast::Sender<()>, SocketAddr)> { +pub async fn start_http_server(config: &config::Config, readiness: Arc) -> Result<(ShutdownHandle, SocketAddr)> { let server_addr = parse_and_resolve_address(config.address.as_str()).map_err(Error::other)?; // The listening address and port are obtained from the parameters @@ -331,8 +328,6 @@ pub async fn start_http_server( // Create shutdown channel let (shutdown_tx, mut shutdown_rx) = tokio::sync::broadcast::channel(1); - let shutdown_tx_clone = shutdown_tx.clone(); - // Create compression configuration from environment variables let compression_config = CompressionConfig::from_env(); if compression_config.enabled { @@ -345,21 +340,12 @@ pub async fn start_http_server( } let is_console = config.console_enable; - tokio::spawn(async move { + let task_handle = tokio::spawn(async move { // Note: CORS layer is removed from global middleware stack // - S3 API CORS is handled by bucket-level CORS configuration in apply_cors_headers() // - Console CORS is handled by its own cors_layer in setup_console_middleware_stack() // This ensures S3 API CORS behavior matches AWS S3 specification - #[cfg(unix)] - let (mut sigterm_inner, mut sigint_inner) = { - use tokio::signal::unix::{SignalKind, signal}; - // Unix platform specific code - let sigterm_inner = signal(SignalKind::terminate()).expect("Failed to create SIGTERM signal handler"); - let sigint_inner = signal(SignalKind::interrupt()).expect("Failed to create SIGINT signal handler"); - (sigterm_inner, sigint_inner) - }; - // ── HTTP Transport Tuning (configurable via env vars) ── // Read all transport parameters from environment, falling back to defaults. // H2 frame size is clamped to RFC 7540 range: 2^14 (16KB) to 2^24 (16MB). @@ -433,8 +419,7 @@ pub async fn start_http_server( .keep_alive_timeout(Duration::from_secs(h2_keep_alive_timeout)); let http_server = Arc::new(conn_builder); - let mut ctrl_c = std::pin::pin!(tokio::signal::ctrl_c()); - let graceful = Arc::new(GracefulShutdown::new()); + let graceful = GracefulShutdown::new(); debug!("graceful initiated"); loop { @@ -450,21 +435,6 @@ pub async fn start_http_server( continue; } }, - _ = ctrl_c.as_mut() => { - info!("Ctrl-C received in worker thread"); - let _ = shutdown_tx_clone.send(()); - break; - }, - Some(_) = sigint_inner.recv() => { - info!("SIGINT received in worker thread"); - let _ = shutdown_tx_clone.send(()); - break; - }, - Some(_) = sigterm_inner.recv() => { - info!("SIGTERM received in worker thread"); - let _ = shutdown_tx_clone.send(()); - break; - }, _ = shutdown_rx.recv() => { info!("Shutdown signal received in worker thread"); break; @@ -481,11 +451,6 @@ pub async fn start_http_server( continue; } }, - _ = ctrl_c.as_mut() => { - info!("Ctrl-C received in worker thread"); - let _ = shutdown_tx_clone.send(()); - break; - }, _ = shutdown_rx.recv() => { info!("Shutdown signal received in worker thread"); break; @@ -531,29 +496,24 @@ pub async fn start_http_server( trusted_proxy_layer: rustfs_trusted_proxies::is_enabled().then(|| rustfs_trusted_proxies::layer().clone()), }; - process_connection(socket, tls_acceptor.clone(), connection_ctx, graceful.clone()); + process_connection(socket, tls_acceptor.clone(), connection_ctx, graceful.watcher()); } - match Arc::try_unwrap(graceful) { - Ok(g) => { - tokio::select! { - () = g.shutdown() => { - debug!("Gracefully shutdown!"); - }, - () = tokio::time::sleep(Duration::from_secs(10)) => { - debug!("Waited 10 seconds for graceful shutdown, aborting..."); - } - } - } - Err(arc_graceful) => { - error!("Cannot perform graceful shutdown, other references exist err: {:?}", arc_graceful); - tokio::time::sleep(Duration::from_secs(10)).await; - debug!("Timeout reached, forcing shutdown"); + let active_connections = graceful.count(); + if active_connections > 0 { + info!(active_connections, "Draining active HTTP connections before shutdown"); + } + tokio::select! { + () = graceful.shutdown() => { + debug!("Gracefully shutdown!"); + }, + () = tokio::time::sleep(Duration::from_secs(10)) => { + warn!(active_connections, "Timed out waiting for HTTP connections to drain during shutdown"); } } }); - Ok((shutdown_tx, local_addr)) + Ok((ShutdownHandle::new(shutdown_tx, task_handle), local_addr)) } #[derive(Clone)] @@ -655,7 +615,7 @@ fn process_connection( socket: TcpStream, tls_acceptor: Option>, context: ConnectionContext, - graceful: Arc, + graceful: Watcher, ) { tokio::spawn(async move { let ConnectionContext { diff --git a/rustfs/src/server/mod.rs b/rustfs/src/server/mod.rs index dec4d90bc..0ce5fa5d4 100644 --- a/rustfs/src/server/mod.rs +++ b/rustfs/src/server/mod.rs @@ -26,6 +26,8 @@ mod runtime; mod service_state; pub mod tls_material; +use tracing::warn; + // Items used by main.rs (binary crate) and/or embedded.rs — must be fully pub. pub use audit::{is_audit_module_enabled, refresh_audit_module_enabled, start_audit_system, stop_audit_system}; pub use event::{init_event_notifier, is_notify_module_enabled, refresh_notify_module_enabled, shutdown_event_notifier}; @@ -48,7 +50,51 @@ pub(crate) use prefix::{ ADMIN_PREFIX, CONSOLE_PREFIX, FAVICON_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, LICENSE, MINIO_ADMIN_PREFIX, MINIO_ADMIN_V3_PREFIX, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH, RPC_PREFIX, RUSTFS_ADMIN_PREFIX, TONIC_PREFIX, VERSION, }; +pub(crate) use readiness::DependencyReadiness; pub(crate) use readiness::ReadinessGateLayer; +pub(crate) use readiness::collect_dependency_readiness; +pub use readiness::publish_ready_when_runtime_ready; #[derive(Clone, Copy, Debug)] pub struct RemoteAddr(pub std::net::SocketAddr); + +pub struct ShutdownHandle { + shutdown_tx: Option>, + task_handle: Option>, +} + +impl ShutdownHandle { + pub fn new(shutdown_tx: tokio::sync::broadcast::Sender<()>, task_handle: tokio::task::JoinHandle<()>) -> Self { + Self { + shutdown_tx: Some(shutdown_tx), + task_handle: Some(task_handle), + } + } + + pub fn signal(&self) { + if let Some(tx) = &self.shutdown_tx { + let _ = tx.send(()); + } + } + + pub async fn shutdown(self) { + self.signal(); + self.wait().await; + } + + pub async fn wait(mut self) { + if let Some(task_handle) = self.task_handle.take() + && let Err(err) = task_handle.await + { + warn!(?err, "Server task join failed during shutdown"); + } + } +} + +impl Drop for ShutdownHandle { + fn drop(&mut self) { + if let Some(tx) = &self.shutdown_tx { + let _ = tx.send(()); + } + } +} diff --git a/rustfs/src/server/readiness.rs b/rustfs/src/server/readiness.rs index bbc9a7be7..384c6513d 100644 --- a/rustfs/src/server/readiness.rs +++ b/rustfs/src/server/readiness.rs @@ -12,18 +12,39 @@ // See the License for the specific language governing permissions and // limitations under the License. +use crate::server::{ServiceState, ServiceStateManager}; use bytes::Bytes; use http::{Request as HttpRequest, Response, StatusCode}; use http_body::Body; use http_body_util::{BodyExt, Full}; use hyper::body::Incoming; use rustfs_common::GlobalReadiness; +use rustfs_ecstore::new_object_layer_fn; +use rustfs_ecstore::store_api::StorageAPI; +use rustfs_iam::get_global_iam_sys; +use rustfs_madmin::{Disk, StorageInfo}; use std::future::Future; use std::pin::Pin; use std::sync::Arc; use std::task::{Context, Poll}; +use std::time::Duration; +use std::{ + collections::{HashMap, HashSet}, + sync::OnceLock, + time::Instant, +}; +use tokio::sync::Mutex; use tower::{Layer, Service}; -use tracing::debug; +use tracing::{debug, info}; + +pub const STARTUP_RUNTIME_READINESS_MAX_WAIT: Duration = Duration::from_secs(30); +pub const STARTUP_RUNTIME_READINESS_POLL_INTERVAL: Duration = Duration::from_secs(1); + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct DependencyReadiness { + pub storage_ready: bool, + pub iam_ready: bool, +} /// ReadinessGateLayer ensures that the system components (IAM, Storage) /// are fully initialized before allowing any request to proceed. @@ -133,3 +154,396 @@ where }) } } + +pub async fn publish_ready_when_runtime_ready( + readiness: &GlobalReadiness, + state_manager: Option<&ServiceStateManager>, +) -> Result<(), std::io::Error> { + wait_for_runtime_readiness_with( + STARTUP_RUNTIME_READINESS_MAX_WAIT, + STARTUP_RUNTIME_READINESS_POLL_INTERVAL, + collect_dependency_readiness_uncached, + |dependency_readiness| { + readiness.mark_stage(rustfs_common::SystemStage::FullReady); + if let Some(state_manager) = state_manager { + state_manager.update(ServiceState::Ready); + } + info!( + target: "rustfs::server::readiness", + storage_ready = dependency_readiness.storage_ready, + iam_ready = dependency_readiness.iam_ready, + "Runtime readiness reached write quorum; publishing ready state" + ); + }, + ) + .await +} + +#[derive(Debug, Clone, Copy)] +struct StorageReadinessCacheEntry { + captured_at: Instant, + storage_ready: bool, +} + +const DISK_STATE_OK: &str = "ok"; +const DISK_STATE_UNFORMATTED: &str = "unformatted"; +const RUNTIME_STATE_RETURNING: &str = "returning"; + +fn health_readiness_cache_ttl() -> Duration { + Duration::from_millis(rustfs_utils::get_env_u64( + rustfs_config::ENV_HEALTH_READINESS_CACHE_TTL_MS, + rustfs_config::DEFAULT_HEALTH_READINESS_CACHE_TTL_MS, + )) +} + +fn storage_readiness_cache() -> &'static Mutex> { + static CACHE: OnceLock>> = OnceLock::new(); + CACHE.get_or_init(|| Mutex::new(None)) +} + +async fn load_cached_storage_readiness() -> Option { + let ttl = health_readiness_cache_ttl(); + if ttl.is_zero() { + return None; + } + + let cache = storage_readiness_cache().lock().await; + let entry = cache.as_ref()?; + if entry.captured_at.elapsed() <= ttl { + return Some(entry.storage_ready); + } + + None +} + +async fn update_storage_readiness_cache(storage_ready: bool) { + if health_readiness_cache_ttl().is_zero() { + return; + } + + let mut cache = storage_readiness_cache().lock().await; + *cache = Some(StorageReadinessCacheEntry { + captured_at: Instant::now(), + storage_ready, + }); +} + +fn disk_is_online_for_readiness(disk: &Disk) -> bool { + let state_is_acceptable = disk.state.eq_ignore_ascii_case(DISK_STATE_OK) + || disk.state.eq_ignore_ascii_case(rustfs_madmin::ITEM_ONLINE) + || disk.state.eq_ignore_ascii_case(DISK_STATE_UNFORMATTED); + + if let Some(runtime_state) = disk.runtime_state.as_deref() { + let runtime_state_is_acceptable = runtime_state.eq_ignore_ascii_case(rustfs_madmin::ITEM_ONLINE) + || runtime_state.eq_ignore_ascii_case(RUNTIME_STATE_RETURNING); + return runtime_state_is_acceptable && state_is_acceptable; + } + + state_is_acceptable +} + +fn pool_write_quorum(info: &StorageInfo, pool_idx: usize, set_drive_count: usize) -> usize { + if set_drive_count == 0 { + return 1; + } + + let data_drives = info + .backend + .standard_sc_data + .get(pool_idx) + .copied() + .filter(|count| *count > 0) + .unwrap_or_else(|| (set_drive_count / 2).max(1)); + + let parity_drives = if let Some(drives_per_set) = info.backend.drives_per_set.get(pool_idx).copied() { + drives_per_set.saturating_sub(data_drives) + } else if let Some(parity) = info.backend.standard_sc_parities.get(pool_idx).copied() { + parity + } else if let Some(parity) = info.backend.standard_sc_parity { + parity + } else { + set_drive_count.saturating_sub(data_drives) + }; + + let mut write_quorum = data_drives; + if data_drives == parity_drives { + write_quorum += 1; + } + write_quorum.max(1) +} + +fn storage_ready_from_runtime_state(info: &StorageInfo) -> bool { + if info.disks.is_empty() { + return false; + } + + let mut total_online = 0usize; + let mut set_online_counts: HashMap<(usize, usize), usize> = HashMap::new(); + let mut set_drive_counts: HashMap<(usize, usize), usize> = HashMap::new(); + let mut seen_disks: HashSet<(String, String, i32, i32, i32)> = HashSet::new(); + + for disk in &info.disks { + if disk.pool_index < 0 || disk.set_index < 0 { + continue; + } + + let dedup_key = ( + disk.endpoint.clone(), + disk.drive_path.clone(), + disk.pool_index, + disk.set_index, + disk.disk_index, + ); + if !seen_disks.insert(dedup_key) { + continue; + } + + let pool_idx = disk.pool_index as usize; + let set_idx = disk.set_index as usize; + let key = (pool_idx, set_idx); + *set_drive_counts.entry(key).or_default() += 1; + + if disk_is_online_for_readiness(disk) { + total_online += 1; + *set_online_counts.entry(key).or_default() += 1; + } + } + + if total_online == 0 || set_drive_counts.is_empty() { + return false; + } + + set_drive_counts.into_iter().all(|((pool_idx, set_idx), set_drive_count)| { + let online = set_online_counts.get(&(pool_idx, set_idx)).copied().unwrap_or_default(); + let write_quorum = pool_write_quorum(info, pool_idx, set_drive_count); + online >= write_quorum + }) +} + +pub async fn collect_dependency_readiness() -> DependencyReadiness { + let iam_ready = get_global_iam_sys().is_some_and(|sys| sys.is_ready()); + let storage_ready = if let Some(cached) = load_cached_storage_readiness().await { + cached + } else { + let computed = collect_storage_readiness_uncached().await; + update_storage_readiness_cache(computed).await; + computed + }; + + DependencyReadiness { + storage_ready, + iam_ready: iam_ready && storage_ready, + } +} + +async fn collect_dependency_readiness_uncached() -> DependencyReadiness { + let iam_ready = get_global_iam_sys().is_some_and(|sys| sys.is_ready()); + let storage_ready = collect_storage_readiness_uncached().await; + + DependencyReadiness { + storage_ready, + iam_ready: iam_ready && storage_ready, + } +} + +async fn collect_storage_readiness_uncached() -> bool { + if let Some(store) = new_object_layer_fn() { + let storage_info = store.storage_info().await; + storage_ready_from_runtime_state(&storage_info) + } else { + false + } +} + +pub async fn wait_for_runtime_readiness_with( + max_wait: Duration, + poll_interval: Duration, + mut load_readiness: F, + mut on_ready: ReadyFn, +) -> Result<(), std::io::Error> +where + F: FnMut() -> Fut, + Fut: Future, + ReadyFn: FnMut(DependencyReadiness), +{ + let startup_deadline = tokio::time::Instant::now() + max_wait; + + loop { + let readiness = load_readiness().await; + if readiness.storage_ready && readiness.iam_ready { + on_ready(readiness); + return Ok(()); + } + + if tokio::time::Instant::now() >= startup_deadline { + let reason = format!( + "startup readiness timed out after {}s: storage_ready={}, iam_ready={}", + max_wait.as_secs(), + readiness.storage_ready, + readiness.iam_ready + ); + return Err(std::io::Error::other(reason)); + } + + info!( + target: "rustfs::server::readiness", + storage_ready = readiness.storage_ready, + iam_ready = readiness.iam_ready, + "Runtime readiness has not reached write quorum yet; delaying ready state publication" + ); + tokio::time::sleep(poll_interval).await; + } +} + +#[cfg(test)] +mod tests { + use super::*; + use rustfs_madmin::{BackendInfo, Disk}; + use std::future; + + #[test] + fn startup_runtime_readiness_wait_constants_are_ordered() { + assert!(STARTUP_RUNTIME_READINESS_MAX_WAIT > STARTUP_RUNTIME_READINESS_POLL_INTERVAL); + assert_eq!(STARTUP_RUNTIME_READINESS_MAX_WAIT.as_secs(), 30); + assert_eq!(STARTUP_RUNTIME_READINESS_POLL_INTERVAL.as_secs(), 1); + } + + #[tokio::test] + async fn wait_for_runtime_readiness_with_does_not_publish_ready_when_runtime_readiness_is_not_reached() { + let readiness = GlobalReadiness::new(); + let state_manager = ServiceStateManager::new(); + + let err = wait_for_runtime_readiness_with( + Duration::ZERO, + Duration::from_millis(1), + || { + future::ready(DependencyReadiness { + storage_ready: false, + iam_ready: false, + }) + }, + |_| { + readiness.mark_stage(rustfs_common::SystemStage::FullReady); + state_manager.update(ServiceState::Ready); + }, + ) + .await + .expect_err("unready startup should time out"); + + assert!(err.to_string().contains("startup readiness timed out")); + assert!(!readiness.is_ready()); + assert_eq!(state_manager.current_state(), ServiceState::Starting); + } + + #[test] + fn storage_ready_from_runtime_state_returns_false_when_all_disks_faulty() { + let info = StorageInfo { + backend: BackendInfo { + standard_sc_data: vec![1], + drives_per_set: vec![1], + ..Default::default() + }, + disks: vec![Disk { + pool_index: 0, + set_index: 0, + state: "offline".to_string(), + runtime_state: Some("offline".to_string()), + ..Default::default() + }], + }; + + assert!(!storage_ready_from_runtime_state(&info)); + } + + #[test] + fn storage_ready_from_runtime_state_returns_true_when_set_meets_write_quorum() { + let info = StorageInfo { + backend: BackendInfo { + standard_sc_data: vec![1], + drives_per_set: vec![1], + ..Default::default() + }, + disks: vec![Disk { + pool_index: 0, + set_index: 0, + state: "ok".to_string(), + runtime_state: Some("online".to_string()), + ..Default::default() + }], + }; + + assert!(storage_ready_from_runtime_state(&info)); + } + + #[test] + fn storage_ready_from_runtime_state_deduplicates_duplicate_disk_rows() { + let duplicate_disk = Disk { + endpoint: "127.0.0.1:9000".to_string(), + drive_path: "/data0".to_string(), + pool_index: 0, + set_index: 0, + disk_index: 0, + state: "ok".to_string(), + runtime_state: Some("online".to_string()), + ..Default::default() + }; + let info = StorageInfo { + backend: BackendInfo { + standard_sc_data: vec![2], + drives_per_set: vec![4], + ..Default::default() + }, + disks: vec![duplicate_disk.clone(), duplicate_disk], + }; + + assert!(!storage_ready_from_runtime_state(&info), "duplicate rows must not satisfy write quorum"); + } + + #[test] + fn disk_online_for_readiness_requires_runtime_and_state_both_acceptable() { + let disk = Disk { + state: "disk io error".to_string(), + runtime_state: Some("online".to_string()), + ..Default::default() + }; + assert!(!disk_is_online_for_readiness(&disk)); + } + + #[test] + fn storage_ready_from_runtime_state_requires_all_sets_meet_quorum() { + let info = StorageInfo { + backend: BackendInfo { + standard_sc_data: vec![1], + drives_per_set: vec![2], + ..Default::default() + }, + disks: vec![ + Disk { + endpoint: "127.0.0.1:9000".to_string(), + drive_path: "/set0d0".to_string(), + pool_index: 0, + set_index: 0, + disk_index: 0, + state: "ok".to_string(), + runtime_state: Some("online".to_string()), + ..Default::default() + }, + Disk { + endpoint: "127.0.0.1:9000".to_string(), + drive_path: "/set1d0".to_string(), + pool_index: 0, + set_index: 1, + disk_index: 0, + state: "offline".to_string(), + runtime_state: Some("offline".to_string()), + ..Default::default() + }, + ], + }; + + assert!( + !storage_ready_from_runtime_state(&info), + "if any set fails write quorum, readiness must be false" + ); + } +} diff --git a/rustfs/src/server/service_state.rs b/rustfs/src/server/service_state.rs index be400c2ba..3022e7346 100644 --- a/rustfs/src/server/service_state.rs +++ b/rustfs/src/server/service_state.rs @@ -35,6 +35,18 @@ pub enum ShutdownSignal { Sigint, } +impl ShutdownSignal { + pub fn log_label(&self) -> &'static str { + match self { + ShutdownSignal::CtrlC => "Ctrl-C", + #[cfg(unix)] + ShutdownSignal::Sigterm => "SIGTERM", + #[cfg(unix)] + ShutdownSignal::Sigint => "SIGINT", + } + } +} + #[atomic_enum] #[derive(PartialEq)] pub enum ServiceState { @@ -234,4 +246,14 @@ mod tests { fn test_ready_maps_to_running_status() { assert_eq!(systemd_status_text(ServiceState::Ready), SERVICE_STATUS_RUNNING); } + + #[test] + fn shutdown_signal_log_label_matches_variants() { + assert_eq!(ShutdownSignal::CtrlC.log_label(), "Ctrl-C"); + #[cfg(unix)] + { + assert_eq!(ShutdownSignal::Sigterm.log_label(), "SIGTERM"); + assert_eq!(ShutdownSignal::Sigint.log_label(), "SIGINT"); + } + } }