From 151103a6097203d7e1fc0a3d5857425a4860aee2 Mon Sep 17 00:00:00 2001 From: Chris Date: Mon, 28 Sep 2026 18:40:30 +0800 Subject: [PATCH] fix(ecstore): recover remote disks after transient stalls Merge the approved fix from pull request #8149. --- crates/ecstore/src/cluster/rpc/remote_disk.rs | 105 +++++++++++++++++- crates/ecstore/src/disk/disk_store.rs | 11 ++ docs/architecture/readiness-matrix.md | 2 + rustfs/src/server/health.rs | 35 +++++- rustfs/src/server/readiness.rs | 79 ++++++++++--- rustfs/src/shared_types.rs | 15 ++- 6 files changed, 221 insertions(+), 26 deletions(-) diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index 1f8cb9042..0ba276567 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -1244,13 +1244,14 @@ impl RemoteDisk { { return; } + // Own the flag before spawning: shutdown can drop the task without polling it. + let lease = RecoveryMonitorLease { + active: Arc::clone(&active), + }; let span = Self::recovery_monitor_span(&addr, &endpoint, handle_id); super::spawn_background_monitor(span, async move { #[cfg(test)] test_state.start_count.fetch_add(1, Ordering::AcqRel); - let lease = RecoveryMonitorLease { - active: Arc::clone(&active), - }; Self::monitor_remote_disk_recovery(addr.clone(), endpoint.clone(), Arc::clone(&health), cancel_token.clone()).await; #[cfg(test)] if let Some(hook) = test_state.teardown_hook.lock().await.take() { @@ -1556,7 +1557,8 @@ impl RemoteDisk { }; if evict_cached_connection { - evict_failed_connection(addr).await; + // Cache contention must not hold the recovery lease indefinitely. + let _ = timeout(get_drive_active_check_timeout(), evict_failed_connection(addr)).await; } result @@ -1715,7 +1717,7 @@ impl RemoteDisk { if timeout_duration == Duration::ZERO { let operation_result = operation().await; if operation_result.is_ok() { - self.health.log_success(); + self.health.record_operation_success(&self.endpoint, "operation_success"); } self.handle_network_like_error(op, timeout_duration, &operation_result, failure_health_action) .await; @@ -1729,7 +1731,7 @@ impl RemoteDisk { Ok(operation_result) => { // Log success; the waiting guard balances every exit path. if operation_result.is_ok() { - self.health.log_success(); + self.health.record_operation_success(&self.endpoint, "operation_success"); } self.handle_network_like_error(op, timeout_duration, &operation_result, failure_health_action) .await; @@ -8204,6 +8206,97 @@ mod tests { assert!(result.is_ok()); } + #[tokio::test] + async fn successful_operation_recovers_suspect_remote_disk_without_monitor() { + let endpoint = Endpoint { + url: url::Url::parse("http://remote-recovery:9000/data").expect("valid endpoint"), + is_local: false, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, + }; + let disk = RemoteDisk::new( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + Arc::new(TcpHttpInternodeDataTransport), + ) + .await + .expect("create remote disk"); + + for duration in [Duration::ZERO, Duration::from_secs(1)] { + disk.health.mark_failure(&endpoint, "read_operation_deadline"); + assert_eq!(disk.runtime_state(), RuntimeDriveHealthState::Suspect); + assert!(!disk.recovery_monitor_is_active()); + let result = disk + .execute_with_timeout(|| async { Err::<(), _>(DiskError::FileNotFound) }, duration) + .await; + assert!(matches!(result, Err(DiskError::FileNotFound))); + assert_eq!( + disk.runtime_state(), + RuntimeDriveHealthState::Suspect, + "a failed operation must not restore readiness" + ); + + disk.execute_with_timeout(|| async { Ok(()) }, duration) + .await + .expect("successful disk RPC"); + assert_eq!( + disk.runtime_state(), + RuntimeDriveHealthState::Online, + "successful traffic must restore readiness without a recovery monitor" + ); + assert!(disk.offline_duration_secs().is_none()); + assert_eq!(disk.health.waiting.load(Ordering::Acquire), 0); + } + + disk.health.mark_offline(&endpoint, "test_offline"); + let result = disk.execute_with_timeout(|| async { Ok(()) }, Duration::from_secs(1)).await; + assert!( + matches!(result, Err(DiskError::FaultyDisk)), + "offline handles still require recovery probes before data I/O" + ); + assert_eq!(disk.runtime_state(), RuntimeDriveHealthState::Offline); + } + + #[test] + fn recovery_monitor_releases_lease_when_dropped_before_first_poll() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("create test runtime"); + let endpoint = Endpoint { + url: url::Url::parse("http://remote-unpolled:9000/data").expect("valid endpoint"), + is_local: false, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, + }; + let active = Arc::new(AtomicBool::new(false)); + let start_count = Arc::new(AtomicU32::new(0)); + { + let _guard = runtime.enter(); + RemoteDisk::schedule_recovery_monitor( + "http://remote-unpolled:9000".to_string(), + endpoint, + Uuid::new_v4(), + Arc::new(DiskHealthTracker::new()), + CancellationToken::new(), + Arc::clone(&active), + RecoveryMonitorTestState { + start_count: Arc::clone(&start_count), + teardown_hook: Arc::new(tokio::sync::Mutex::new(None)), + }, + ); + assert!(active.load(Ordering::Acquire)); + } + drop(runtime); + assert_eq!(start_count.load(Ordering::Acquire), 0, "task must never be polled"); + assert!(!active.load(Ordering::Acquire), "dropping an unpolled monitor must release its lease"); + } + #[tokio::test] async fn test_execute_with_timeout_marks_remote_disk_faulty() { let url = url::Url::parse("http://remote-timeout:9000").expect("operation should succeed"); diff --git a/crates/ecstore/src/disk/disk_store.rs b/crates/ecstore/src/disk/disk_store.rs index 8c4fb2a72..612a5276b 100644 --- a/crates/ecstore/src/disk/disk_store.rs +++ b/crates/ecstore/src/disk/disk_store.rs @@ -1236,6 +1236,17 @@ impl DiskHealthTracker { record_drive_recovery_class(classify_drive_recovery(duration)); } self.offline_since_unix_secs.store(0, Ordering::Release); + info!( + event = EVENT_DISK_RECOVERY_PROBE_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_DISK, + endpoint = %endpoint, + state = "recovered", + previous_state = current.as_str(), + runtime_state = next.as_str(), + reason, + "Disk recovered" + ); } else if let Some(duration) = self.offline_duration() { record_drive_offline_duration(endpoint, duration); } diff --git a/docs/architecture/readiness-matrix.md b/docs/architecture/readiness-matrix.md index c034c3024..58d98c620 100644 --- a/docs/architecture/readiness-matrix.md +++ b/docs/architecture/readiness-matrix.md @@ -69,6 +69,8 @@ The existing `details.storage.ready` boolean and `connected` / `disconnected` st Node readiness additionally reports `details.storage.readQuorum`, `details.storage.writeQuorum`, and `details.poolMetadata.ready`. `details.storage.ready` follows read quorum. `details.lock.ready` on this probe follows shared-lock quorum; on `/minio/health/cluster` it follows exclusive-lock quorum. The metadata component's status is `writable` or `unavailable`. A blocked metadata writer is visible there and on the cluster-write probe; it does not by itself make `/health/ready` return 503. +`details.storage.unavailableDrives` lists inventoried handles excluded from the node's quorum snapshot. Each entry has zero-based `poolIndex`, `setIndex`, and `diskIndex`, its `runtimeState`, and `hostOnline` from the lock reachability observation (always true for local drives). A `suspect` drive with `hostOnline: true` is reachable but not counted; an `online` drive with `hostOnline: false` is excluded because its host did not answer. Internal addresses and filesystem paths are omitted. This list does not enumerate missing inventory slots, and an empty list does not override a failed inventory or quorum check. It is omitted with the other details in minimal responses; HEAD responses have no body. Admin storage info still reports each owner's disk view, which can differ from these node-local handle states. + Node storage quorum uses configured drives per set, all configured pools/sets, and their Standard storage-class data/parity layout. Missing, duplicate, unreachable, or unhealthy disk observations cannot supply extra quorum votes. The read quorum is the data-drive count; the write quorum is that count plus one when data and parity counts are equal. These are observations of available storage slots, not guarantees that a particular object's metadata, shards, or required locks are available. For a healthy IAM and metadata writer in a four-node, one-drive-per-node EC 2+2 set, after startup has published `FullReady`: diff --git a/rustfs/src/server/health.rs b/rustfs/src/server/health.rs index e4103d5ff..ff687f73c 100644 --- a/rustfs/src/server/health.rs +++ b/rustfs/src/server/health.rs @@ -355,11 +355,12 @@ pub(crate) fn build_health_response_parts( kms_ready, include_dependency_details, }); - if let Some(details) = readiness_report.and_then(|report| report.storage_details) + if let Some(details) = readiness_report.and_then(|report| report.storage_details.as_ref()) && payload.get("details").is_some() { payload["details"]["storage"]["readQuorum"] = json!(details.read_quorum_ready); payload["details"]["storage"]["writeQuorum"] = json!(details.write_quorum_ready); + payload["details"]["storage"]["unavailableDrives"] = json!(details.unavailable_drives); payload["details"]["poolMetadata"] = json!({ "ready": details.pool_metadata_write_ready, "status": if details.pool_metadata_write_ready { "writable" } else { "unavailable" }, @@ -478,6 +479,7 @@ mod tests { read_quorum_ready: read_quorum, write_quorum_ready: write_quorum, pool_metadata_write_ready: metadata_ready, + ..Default::default() }); let parts = build_health_response_parts(Method::GET, HealthProbe::Readiness, Some(&report), "rustfs", None, None); let expected_ready = read_quorum && lock_ready; @@ -510,6 +512,36 @@ mod tests { }); } + #[test] + #[serial] + fn node_storage_details_explain_excluded_drives_without_exposing_endpoints() { + with_var(rustfs_config::ENV_HEALTH_MINIMAL_RESPONSE_ENABLE, Some("false"), || { + let mut report = ready_report(); + report.readiness.storage_ready = false; + report.storage_details = Some(crate::shared_types::StorageReadinessDetails { + pool_metadata_write_ready: true, + unavailable_drives: vec![crate::shared_types::UnavailableReadinessDrive { + pool_index: 0, + set_index: 1, + disk_index: 2, + runtime_state: "suspect".to_string(), + host_online: true, + }], + ..Default::default() + }); + let parts = build_health_response_parts(Method::GET, HealthProbe::Readiness, Some(&report), "rustfs", None, None); + assert_eq!(parts.status_code, StatusCode::SERVICE_UNAVAILABLE); + let payload = parts.payload.expect("readiness GET body"); + assert_eq!( + payload["details"]["storage"]["unavailableDrives"], + json!([{ + "poolIndex": 0, "setIndex": 1, "diskIndex": 2, + "runtimeState": "suspect", "hostOnline": true, + }]) + ); + }); + } + #[test] #[serial] fn node_storage_details_do_not_expand_minimal_or_liveness_payloads() { @@ -518,6 +550,7 @@ mod tests { read_quorum_ready: true, write_quorum_ready: true, pool_metadata_write_ready: true, + ..Default::default() }); with_var(rustfs_config::ENV_HEALTH_MINIMAL_RESPONSE_ENABLE, Some("true"), || { let parts = build_health_response_parts(Method::GET, HealthProbe::Readiness, Some(&report), "rustfs", None, None); diff --git a/rustfs/src/server/readiness.rs b/rustfs/src/server/readiness.rs index bdb455aff..127d7b7e4 100644 --- a/rustfs/src/server/readiness.rs +++ b/rustfs/src/server/readiness.rs @@ -77,7 +77,9 @@ const METRIC_RUNTIME_READINESS_READY: &str = "rustfs_runtime_readiness_ready"; const METRIC_RUNTIME_READINESS_DEGRADED_TOTAL: &str = "rustfs_runtime_readiness_degraded_total"; const METRIC_POOL_METADATA_CHECK_TIMEOUT_TOTAL: &str = "rustfs_pool_metadata_check_timeouts_total"; -pub use crate::shared_types::{DependencyReadiness, DependencyReadinessReport, ReadinessDegradedReason, StorageReadinessDetails}; +pub use crate::shared_types::{ + DependencyReadiness, DependencyReadinessReport, ReadinessDegradedReason, StorageReadinessDetails, UnavailableReadinessDrive, +}; /// ReadinessGateLayer ensures that the system components (IAM, Storage) /// are fully initialized before allowing any request to proceed. @@ -920,7 +922,8 @@ pub async fn collect_node_readiness_report() -> DependencyReadinessReport { let pool_metadata_status = store.pool_meta_write_status().await; storage = apply_node_pool_metadata_timeout_policy(pool_metadata_write_readiness(pool_metadata_status), Instant::now()); match node_storage_snapshot(store.as_ref(), &lock_observation.online_hosts).await { - Ok(info) => { + Ok((info, unavailable_drives)) => { + details.unavailable_drives = unavailable_drives; details.read_quorum_ready = storage_read_ready_from_runtime_state(&info); details.write_quorum_ready = storage_ready_from_runtime_state(&info); } @@ -947,7 +950,10 @@ pub async fn collect_node_readiness_report() -> DependencyReadinessReport { report } -async fn node_storage_snapshot(store: &S, online_hosts: &HashSet) -> Result +async fn node_storage_snapshot( + store: &S, + online_hosts: &HashSet, +) -> Result<(StorageInfo, Vec), StorageError> where S: StorageAdminApi, { @@ -956,8 +962,9 @@ where backend: store.backend_info().await, ..Default::default() }; + let mut unavailable_drives = Vec::new(); if configured_readiness_topology(&info).is_none() { - return Ok(info); + return Ok((info, unavailable_drives)); } // Inventory and runtime health are local snapshots. Reuse the lock probe's @@ -971,26 +978,38 @@ where .into_iter() .flatten() { - info.disks.push(node_disk_snapshot( - disk_endpoint_snapshot(&disk), - disk.runtime_state().as_str(), - online_hosts, - )); + let (snapshot, unavailable) = + node_disk_snapshot(disk_endpoint_snapshot(&disk), disk.runtime_state().as_str(), online_hosts); + if let Some(unavailable) = unavailable { + unavailable_drives.push(unavailable); + } + info.disks.push(snapshot); } } } - Ok(info) + Ok((info, unavailable_drives)) }) .await .map_err(|_| StorageError::Timeout)? } -fn node_disk_snapshot(endpoint: Endpoint, runtime_state: &str, online_hosts: &HashSet) -> Disk { +fn node_disk_snapshot( + endpoint: Endpoint, + runtime_state: &str, + online_hosts: &HashSet, +) -> (Disk, Option) { let reachable = endpoint.is_local || online_hosts.contains(&endpoint.host_port()); // Returning drives can still reject data I/O as faulty. Without a fresh // disk-info probe, only an Online runtime observation can supply quorum. let online = reachable && runtime_state == rustfs_madmin::ITEM_ONLINE; - Disk { + let unavailable = (!online).then(|| UnavailableReadinessDrive { + pool_index: endpoint.pool_idx, + set_index: endpoint.set_idx, + disk_index: endpoint.disk_idx, + runtime_state: runtime_state.to_string(), + host_online: reachable, + }); + let disk = Disk { endpoint: endpoint.to_string(), drive_path: endpoint.get_file_path(), pool_index: endpoint.pool_idx, @@ -999,7 +1018,8 @@ fn node_disk_snapshot(endpoint: Endpoint, runtime_state: &str, online_hosts: &Ha state: if online { DISK_STATE_OK } else { "offline" }.to_string(), runtime_state: Some(runtime_state.to_string()), ..Default::default() - } + }; + (disk, unavailable) } async fn collect_cluster_health_report_with( @@ -1392,10 +1412,16 @@ mod tests { let write_quorum = data + usize::from(data == parity); for survivors in (0..=drive_count).rev().chain(std::iter::once(drive_count)) { let online_hosts = (0..survivors).map(|idx| format!("node-{idx}:9000")).collect(); - let info = node_storage_snapshot(&store, &online_hosts) + let (info, unavailable_drives) = node_storage_snapshot(&store, &online_hosts) .await .expect("read local runtime inventory"); assert_eq!(info.disks.len(), drive_count, "offline members retain their topology slots"); + assert_eq!(unavailable_drives.len(), drive_count - survivors); + assert!( + unavailable_drives + .iter() + .all(|disk| !disk.host_online && disk.runtime_state == "online") + ); assert_eq!( storage_read_ready_from_runtime_state(&info), survivors >= data, @@ -1428,7 +1454,22 @@ mod tests { disk_idx: 2, }; let mut disks = online_readiness_disks(0, 2); - disks.push(node_disk_snapshot(endpoint, runtime_state, &online_hosts)); + let (disk, unavailable) = node_disk_snapshot(endpoint, runtime_state, &online_hosts); + if (is_local || reachable) && runtime_state == "online" { + assert!(unavailable.is_none()); + } else { + assert_eq!( + unavailable, + Some(UnavailableReadinessDrive { + pool_index: 0, + set_index: 0, + disk_index: 2, + runtime_state: runtime_state.to_string(), + host_online: is_local || reachable, + }) + ); + } + disks.push(disk); let info = StorageInfo { backend: BackendInfo { total_sets: vec![1], @@ -1455,21 +1496,21 @@ mod tests { async fn node_storage_snapshot_requires_each_configured_set_and_distinct_drives() { let mut store = runtime_inventory(&[(2, 4, 2), (1, 8, 2)]).await; let online_hosts = (0..8).map(|idx| format!("node-{idx}:9000")).collect(); - let info = node_storage_snapshot(&store, &online_hosts) + let (info, _) = node_storage_snapshot(&store, &online_hosts) .await .expect("healthy mixed layout"); assert!(storage_ready_from_runtime_state(&info)); let selector = DiskSetSelector::new(0, 1); let original = store.disks.remove(&selector).expect("second configured set"); - let info = node_storage_snapshot(&store, &online_hosts) + let (info, _) = node_storage_snapshot(&store, &online_hosts) .await .expect("missing set snapshot"); assert!(!storage_read_ready_from_runtime_state(&info)); assert!(!storage_ready_from_runtime_state(&info)); store.disks.insert(selector, vec![original[0].clone(); 4]); - let info = node_storage_snapshot(&store, &online_hosts) + let (info, _) = node_storage_snapshot(&store, &online_hosts) .await .expect("duplicate drive snapshot"); assert!(!storage_read_ready_from_runtime_state(&info)); @@ -1480,6 +1521,7 @@ mod tests { &node_storage_snapshot(&store, &online_hosts) .await .expect("restored inventory") + .0 )); } @@ -1502,6 +1544,7 @@ mod tests { &node_storage_snapshot(&store, &online_hosts) .await .expect("inspection recovered") + .0 )); } diff --git a/rustfs/src/shared_types.rs b/rustfs/src/shared_types.rs index 3a872eb8c..b29b328a4 100644 --- a/rustfs/src/shared_types.rs +++ b/rustfs/src/shared_types.rs @@ -85,11 +85,24 @@ pub struct DependencyReadinessReport { pub storage_details: Option, } -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Default, PartialEq, Eq)] pub struct StorageReadinessDetails { pub read_quorum_ready: bool, pub write_quorum_ready: bool, pub pool_metadata_write_ready: bool, + pub unavailable_drives: Vec, +} + +/// Node-local reasons for excluding an inventoried drive from readiness quorum. +/// Use topology indices instead of exposing internal addresses or filesystem paths. +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] +#[serde(rename_all = "camelCase")] +pub struct UnavailableReadinessDrive { + pub pool_index: i32, + pub set_index: i32, + pub disk_index: i32, + pub runtime_state: String, + pub host_online: bool, } pub(crate) fn convert_ecstore_object_info(object: StorageObjectInfo) -> NotifyObjectInfo {