From e0b8c4fd4226863df45375043db6640642ba05bf Mon Sep 17 00:00:00 2001 From: weisd Date: Tue, 28 Apr 2026 14:38:31 +0800 Subject: [PATCH] fix(storage): avoid faulting local drives on transient timeouts (#2714) --- crates/ecstore/src/disk/disk_store.rs | 94 ++++++++++++++++++++------- crates/ecstore/src/rpc/remote_disk.rs | 8 +-- 2 files changed, 74 insertions(+), 28 deletions(-) diff --git a/crates/ecstore/src/disk/disk_store.rs b/crates/ecstore/src/disk/disk_store.rs index 9d28c05f4..1045163e5 100644 --- a/crates/ecstore/src/disk/disk_store.rs +++ b/crates/ecstore/src/disk/disk_store.rs @@ -231,9 +231,26 @@ impl DiskHealthTracker { RuntimeDriveHealthState::Offline => RuntimeDriveHealthState::Offline, }; - self.status.store(DISK_HEALTH_FAULTY, Ordering::Release); + let became_offline = next == RuntimeDriveHealthState::Offline && current != RuntimeDriveHealthState::Offline; + if next == RuntimeDriveHealthState::Offline { + self.status.store(DISK_HEALTH_FAULTY, Ordering::Release); + } else { + self.status.store(DISK_HEALTH_OK, Ordering::Release); + } self.transition_state(endpoint, current, next, reason); - current == RuntimeDriveHealthState::Online + became_offline + } + + pub fn mark_offline(&self, endpoint: &Endpoint, reason: &'static str) -> bool { + let current = self.runtime_state(); + if current == RuntimeDriveHealthState::Offline { + return false; + } + + self.consecutive_successes.store(0, Ordering::Release); + self.status.store(DISK_HEALTH_FAULTY, Ordering::Release); + self.transition_state(endpoint, current, RuntimeDriveHealthState::Offline, reason); + true } pub fn mark_recovery_success(&self, endpoint: &Endpoint, reason: &'static str) -> bool { @@ -268,6 +285,14 @@ impl DiskHealthTracker { became_online } + pub fn record_operation_success(&self, endpoint: &Endpoint, reason: &'static str) { + if self.runtime_state() == RuntimeDriveHealthState::Online { + self.log_success(); + } else { + self.mark_recovery_success(endpoint, reason); + } + } + fn transition_state( &self, endpoint: &Endpoint, @@ -707,7 +732,7 @@ impl LocalDiskWrapper { let result = operation().await; self.health.decrement_waiting(); if result.is_ok() { - self.health.log_success(); + self.health.record_operation_success(&self.endpoint(), "operation_success"); } return result; } @@ -718,7 +743,7 @@ impl LocalDiskWrapper { Ok(operation_result) => { // Log success and decrement waiting counter if operation_result.is_ok() { - self.health.log_success(); + self.health.record_operation_success(&self.endpoint(), "operation_success"); } self.health.decrement_waiting(); operation_result @@ -919,7 +944,7 @@ impl DiskAPI for LocalDiskWrapper { let has_err = result.iter().any(|e| e.is_some()); if !has_err { // Log success and decrement waiting counter - self.health.log_success(); + self.health.record_operation_success(&self.endpoint(), "operation_success"); } result @@ -1114,35 +1139,56 @@ mod tests { #[test] fn runtime_state_transitions_from_online_to_suspect_then_offline() { - let endpoint = Endpoint::try_from("/tmp/runtime-state-disk").expect("endpoint should parse"); - let health = DiskHealthTracker::new(); + temp_env::with_var(rustfs_config::ENV_DRIVE_SUSPECT_FAILURE_THRESHOLD, Some("2"), || { + let endpoint = Endpoint::try_from("/tmp/runtime-state-disk").expect("endpoint should parse"); + let health = DiskHealthTracker::new(); - assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Online); - assert!(health.mark_failure(&endpoint, "timeout")); - assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Suspect); + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Online); + assert!(!health.mark_failure(&endpoint, "timeout")); + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Suspect); + assert!(!health.is_faulty()); - assert!(!health.mark_failure(&endpoint, "timeout")); - assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Offline); - assert!(health.offline_duration().is_some()); + assert!(health.mark_failure(&endpoint, "timeout")); + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Offline); + assert!(health.is_faulty()); + assert!(health.offline_duration().is_some()); + }); } #[test] fn runtime_state_transitions_back_online_after_recovery_threshold() { - let endpoint = Endpoint::try_from("/tmp/runtime-state-recovery").expect("endpoint should parse"); + temp_env::with_var(rustfs_config::ENV_DRIVE_SUSPECT_FAILURE_THRESHOLD, Some("2"), || { + let endpoint = Endpoint::try_from("/tmp/runtime-state-recovery").expect("endpoint should parse"); + let health = DiskHealthTracker::new(); + + health.mark_failure(&endpoint, "timeout"); + health.mark_failure(&endpoint, "timeout"); + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Offline); + + assert!(!health.mark_recovery_success(&endpoint, "probe")); + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Returning); + + assert!(!health.mark_recovery_success(&endpoint, "probe")); + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Returning); + + assert!(health.mark_recovery_success(&endpoint, "probe")); + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Online); + assert!(health.offline_duration().is_none()); + }); + } + + #[test] + fn operation_success_recovers_suspect_drive_without_faulting() { + let endpoint = Endpoint::try_from("/tmp/runtime-state-suspect-success").expect("endpoint should parse"); let health = DiskHealthTracker::new(); - health.mark_failure(&endpoint, "timeout"); - health.mark_failure(&endpoint, "timeout"); - assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Offline); + assert!(!health.mark_failure(&endpoint, "timeout")); + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Suspect); + assert!(!health.is_faulty()); - assert!(!health.mark_recovery_success(&endpoint, "probe")); - assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Returning); - - assert!(!health.mark_recovery_success(&endpoint, "probe")); - assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Returning); - - assert!(health.mark_recovery_success(&endpoint, "probe")); + health.record_operation_success(&endpoint, "operation_success"); assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Online); + assert!(!health.is_faulty()); assert!(health.offline_duration().is_none()); } } diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index 4809b58c6..faf2e6a66 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -179,7 +179,7 @@ impl RemoteDisk { let mut interval = time::interval(CHECK_EVERY); // Perform basic connectivity check - if Self::perform_connectivity_check(&addr).await.is_err() && health.mark_failure(&endpoint, "connectivity_probe_failed") { + if Self::perform_connectivity_check(&addr).await.is_err() && health.mark_offline(&endpoint, "connectivity_probe_failed") { warn!("Remote disk health check failed for {}: marking as faulty", addr); // Start recovery monitoring @@ -222,7 +222,7 @@ impl RemoteDisk { } // Perform basic connectivity check - if Self::perform_connectivity_check(&addr).await.is_err() && health.mark_failure(&endpoint, "connectivity_probe_failed") { + if Self::perform_connectivity_check(&addr).await.is_err() && health.mark_offline(&endpoint, "connectivity_probe_failed") { warn!("Remote disk health check failed for {}: marking as faulty", addr); // Start recovery monitoring @@ -263,7 +263,7 @@ impl RemoteDisk { return; } } else { - health.mark_failure(&endpoint, "connectivity_probe_failed"); + health.mark_offline(&endpoint, "connectivity_probe_failed"); } } } @@ -380,7 +380,7 @@ impl RemoteDisk { } async fn mark_faulty_and_evict(&self, reason: &'static str) { - if self.health.mark_failure(&self.endpoint, reason) { + if self.health.mark_offline(&self.endpoint, reason) { self.spawn_recovery_monitor_if_needed(); counter!( "rustfs_drive_faulty_mark_total",