diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index abab9a503..3c7f354df 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -231,6 +231,8 @@ pub struct RemoteDisk { recovery_monitor_active: Arc, #[cfg(test)] recovery_monitor_start_count: Arc, + #[cfg(test)] + recovery_monitor_teardown_hook: Arc>>>, data_transport: Arc, } @@ -244,6 +246,13 @@ impl Drop for RecoveryMonitorLease { } } +#[cfg(test)] +#[derive(Debug, Default)] +struct RecoveryMonitorTeardownHook { + arrived: tokio::sync::Notify, + release: tokio::sync::Notify, +} + // ── Connection lifecycle (grpc-optimization P3) ── /// Whether to prewarm the internode control channel in the background at construction (default off). @@ -438,6 +447,8 @@ impl RemoteDisk { recovery_monitor_active: Arc::new(AtomicBool::new(false)), #[cfg(test)] recovery_monitor_start_count: Arc::new(AtomicU32::new(0)), + #[cfg(test)] + recovery_monitor_teardown_hook: Arc::new(tokio::sync::Mutex::new(None)), data_transport, }; record_drive_runtime_state(ep, RuntimeDriveHealthState::Online); @@ -612,6 +623,8 @@ impl RemoteDisk { Arc::clone(&self.recovery_monitor_active), #[cfg(test)] Arc::clone(&self.recovery_monitor_start_count), + #[cfg(test)] + Arc::clone(&self.recovery_monitor_teardown_hook), ); } @@ -623,6 +636,7 @@ impl RemoteDisk { cancel_token: CancellationToken, active: Arc, #[cfg(test)] start_count: Arc, + #[cfg(test)] teardown_hook: Arc>>>, ) { if active .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire) @@ -638,6 +652,11 @@ impl RemoteDisk { 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) = teardown_hook.lock().await.take() { + hook.arrived.notify_one(); + hook.release.notified().await; + } drop(lease); if !cancel_token.is_cancelled() && health.runtime_state() != RuntimeDriveHealthState::Online { Self::schedule_recovery_monitor( @@ -649,6 +668,8 @@ impl RemoteDisk { active, #[cfg(test)] start_count, + #[cfg(test)] + teardown_hook, ); } }); @@ -692,9 +713,21 @@ impl RemoteDisk { let endpoint = self.endpoint.clone(); let handle_id = self.handle_id; let recovery_monitor_active = Arc::clone(&self.recovery_monitor_active); + #[cfg(test)] + let recovery_monitor_teardown_hook = Arc::clone(&self.recovery_monitor_teardown_hook); tokio::spawn(async move { - Self::monitor_remote_disk_health(addr, endpoint, handle_id, health, cancel_token, recovery_monitor_active).await; + Self::monitor_remote_disk_health( + addr, + endpoint, + handle_id, + health, + cancel_token, + recovery_monitor_active, + #[cfg(test)] + recovery_monitor_teardown_hook, + ) + .await; }); } @@ -706,6 +739,7 @@ impl RemoteDisk { health: Arc, cancel_token: CancellationToken, recovery_monitor_active: Arc, + #[cfg(test)] recovery_monitor_teardown_hook: Arc>>>, ) { let mut interval = time::interval(get_drive_active_check_interval()); @@ -739,6 +773,8 @@ impl RemoteDisk { Arc::clone(&recovery_monitor_active), #[cfg(test)] Arc::new(AtomicU32::new(0)), + #[cfg(test)] + Arc::clone(&recovery_monitor_teardown_hook), ); } @@ -807,6 +843,8 @@ impl RemoteDisk { Arc::clone(&recovery_monitor_active), #[cfg(test)] Arc::new(AtomicU32::new(0)), + #[cfg(test)] + Arc::clone(&recovery_monitor_teardown_hook), ); } } @@ -4755,6 +4793,78 @@ mod tests { assert!(!disk.recovery_monitor_is_active()); } + #[tokio::test] + #[serial(remote_disk_recovery_probe)] + async fn recovery_monitor_rearms_if_disk_fails_during_teardown() { + runtime_sources::ensure_test_rpc_secret(); + let Some(peer) = TestGrpcPeer::spawn(Bytes::new(), Bytes::new()).await else { + return; + }; + let endpoint = Endpoint { + url: url::Url::parse(&format!("{}/data/rustfs0", peer.addr)).expect("endpoint should parse"), + is_local: false, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, + }; + let disk = RemoteDisk::new( + &endpoint, + &DiskOption { + cleanup: false, + health_check: true, + }, + Arc::new(TcpHttpInternodeDataTransport), + ) + .await + .expect("remote disk should construct"); + if !disk.health_check { + peer.stop().await; + return; + } + + disk.force_runtime_state_for_test(RuntimeDriveHealthState::Offline); + let hook = Arc::new(RecoveryMonitorTeardownHook::default()); + *disk.recovery_monitor_teardown_hook.lock().await = Some(Arc::clone(&hook)); + + temp_env::async_with_vars( + [ + (rustfs_config::ENV_DRIVE_RETURNING_PROBE_INTERVAL_SECS, Some("1")), + (rustfs_config::ENV_DRIVE_RETURNING_SUCCESS_THRESHOLD, Some("1")), + (rustfs_config::ENV_DRIVE_ACTIVE_CHECK_TIMEOUT_SECS, Some("1")), + ], + async { + disk.spawn_recovery_monitor_if_needed(); + tokio::time::timeout(Duration::from_secs(5), hook.arrived.notified()) + .await + .expect("first recovery monitor should reach teardown"); + assert_eq!(disk.runtime_state(), RuntimeDriveHealthState::Online); + + disk.force_runtime_state_for_test(RuntimeDriveHealthState::Offline); + hook.release.notify_one(); + tokio::time::timeout(Duration::from_secs(2), async { + while disk.recovery_monitor_start_count() < 2 { + tokio::task::yield_now().await; + } + }) + .await + .expect("teardown failure should re-arm recovery monitoring"); + assert!(disk.recovery_monitor_is_active(), "re-armed monitor should retain single-flight ownership"); + + disk.cancel_token.cancel(); + tokio::time::timeout(Duration::from_secs(2), async { + while disk.recovery_monitor_is_active() { + tokio::task::yield_now().await; + } + }) + .await + .expect("cancelled re-armed monitor should release single-flight state"); + }, + ) + .await; + + peer.stop().await; + } + #[tokio::test] #[serial(remote_disk_recovery_probe)] async fn recovery_monitor_restores_online_then_real_reads_use_replacement_handle() { @@ -6433,6 +6543,17 @@ mod tests { ) .await .expect("remote disk should construct"); + let replacement = RemoteDisk::new( + &endpoint, + &DiskOption { + cleanup: false, + health_check: true, + }, + Arc::new(TcpHttpInternodeDataTransport), + ) + .await + .expect("replacement remote disk should construct"); + assert_ne!(remote_disk.handle_id, replacement.handle_id, "replacement handles need distinct log identities"); let span = tracing::info_span!("request-span", request_id = "req-remote-disk"); let _entered = span.enter(); @@ -6448,11 +6569,24 @@ mod tests { assert_eq!(log["span"]["name"], Value::String("recovery-monitor".to_string())); assert_eq!(log["span"]["kind"], Value::String("remote_disk".to_string())); + assert_eq!(log["span"]["handle_id"], Value::String(remote_disk.handle_id.to_string())); let spans = log["spans"].as_array().expect("spans should be present"); assert!(spans.iter().any(|span| { span.get("name").and_then(Value::as_str) == Some("request-span") && span.get("request_id").and_then(Value::as_str) == Some("req-remote-disk") })); + + remote_disk.force_runtime_state_for_test(RuntimeDriveHealthState::Offline); + remote_disk + .execute_with_timeout(|| async { Ok::<(), Error>(()) }, Duration::from_secs(1)) + .await + .expect_err("faulty handle should short-circuit"); + let faulty_log = logs + .lines() + .into_iter() + .find(|value| value.get("state").and_then(Value::as_str) == Some("faulty_short_circuit")) + .expect("expected faulty short-circuit log"); + assert_eq!(faulty_log["handle_id"], Value::String(remote_disk.handle_id.to_string())); } #[tokio::test(flavor = "current_thread")] diff --git a/crates/ecstore/src/disk/disk_store.rs b/crates/ecstore/src/disk/disk_store.rs index 111b8fabc..35fdf2ced 100644 --- a/crates/ecstore/src/disk/disk_store.rs +++ b/crates/ecstore/src/disk/disk_store.rs @@ -2962,32 +2962,54 @@ mod tests { } #[test] + #[serial_test::serial] fn concurrent_failure_and_recovery_publish_one_health_snapshot() { - let endpoint = Endpoint::try_from("/tmp/concurrent-health-snapshot").expect("endpoint should parse"); - let health = Arc::new(DiskHealthTracker::new()); - let workers = (0..8) - .map(|_| { - let health = Arc::clone(&health); - let endpoint = endpoint.clone(); - std::thread::spawn(move || { - for _ in 0..32 { + temp_env::with_var(rustfs_config::ENV_DRIVE_SUSPECT_FAILURE_THRESHOLD, Some("2"), || { + let endpoint = Endpoint::try_from("/tmp/concurrent-health-snapshot").expect("endpoint should parse"); + let health = Arc::new(DiskHealthTracker::new()); + let transition_guard = health + .transition_lock + .lock() + .expect("health transition lock should not be poisoned"); + let start = Arc::new(std::sync::Barrier::new(3)); + let (completed_tx, completed_rx) = std::sync::mpsc::channel(); + let workers = (0..2) + .map(|_| { + let health = Arc::clone(&health); + let endpoint = endpoint.clone(); + let start = Arc::clone(&start); + let completed_tx = completed_tx.clone(); + std::thread::spawn(move || { + start.wait(); health.mark_failure(&endpoint, "concurrent_test"); - health.mark_recovery_success(&endpoint, "concurrent_test"); - let (runtime, faulty) = health.health_state_snapshot(); - assert!(matches!( - (runtime, faulty), - (RuntimeDriveHealthState::Online, false) - | (RuntimeDriveHealthState::Suspect, false) - | (RuntimeDriveHealthState::Offline, true) - | (RuntimeDriveHealthState::Returning, true) - )); - } + completed_tx.send(()).expect("completion receiver should remain available"); + }) }) - }) - .collect::>(); - for worker in workers { - worker.join().expect("health transition worker should not panic"); - } + .collect::>(); + + start.wait(); + assert!( + matches!( + completed_rx.recv_timeout(Duration::from_millis(250)), + Err(std::sync::mpsc::RecvTimeoutError::Timeout) + ), + "concurrent transitions must wait for the serialization lock" + ); + drop(transition_guard); + completed_rx + .recv_timeout(Duration::from_secs(1)) + .expect("first failure transition should complete after lock release"); + completed_rx + .recv_timeout(Duration::from_secs(1)) + .expect("second failure transition should complete after lock release"); + for worker in workers { + worker.join().expect("health transition worker should not panic"); + } + + assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Offline); + assert!(health.is_faulty()); + assert_eq!(health.consecutive_failures.load(Ordering::Acquire), 2); + }); } #[test]