From bb04355a3ded1ec3cc1793ea59daa673ac6073d0 Mon Sep 17 00:00:00 2001 From: Chris Date: Wed, 30 Sep 2026 15:45:59 +0800 Subject: [PATCH] fix(heal): retry contended local metadata CAS (#8267) --- crates/heal/src/error.rs | 3 ++ crates/heal/src/heal/manager/tests.rs | 64 +++++++++++++++++++++++++++ 2 files changed, 67 insertions(+) diff --git a/crates/heal/src/error.rs b/crates/heal/src/error.rs index cf665187c..8d2f32aae 100644 --- a/crates/heal/src/error.rs +++ b/crates/heal/src/error.rs @@ -151,6 +151,9 @@ impl Error { ) || is_recoverable_heal_error_message(&err.to_string()) } + // Nonblocking local CAS and replacement leases report lock + // contention as WouldBlock; retain the existing task retry budget. + Error::Disk(DiskError::Io(error)) if error.kind() == std::io::ErrorKind::WouldBlock => true, Error::Disk(err) => { if err.is_dangling_delete_grace() { return true; diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index 3a32bf84a..12a690399 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -2569,6 +2569,70 @@ fn test_retry_request_for_recoverable_lock_timeout() { assert!(retry_error.contains("Lock acquisition timeout")); } +#[cfg(unix)] +#[tokio::test] +async fn contended_healing_marker_cas_retains_bounded_task_retries() { + let temp = TempDir::new().expect("marker contention directory"); + let disk = make_manager_resume_disk(&temp, "marker-contention").await; + let metadata = temp.path().join("marker-contention").join(super::super::RUSTFS_META_BUCKET); + let lock = std::fs::OpenOptions::new() + .create(true) + .truncate(false) + .read(true) + .write(true) + .open(metadata.join(".rustfs-cas.lock")) + .expect("marker CAS lock should open"); + lock.lock().expect("fixture should own the marker CAS lock"); + let result = super::super::apply_healing_markers_to_targets(vec![disk.clone()], Some("owner"), None, false).await; + assert!( + matches!(&result, Err(Error::Disk(DiskError::Io(error))) if error.kind() == std::io::ErrorKind::WouldBlock), + "the real contended marker CAS must return WouldBlock: {result:?}" + ); + assert!( + !metadata.join(super::super::HEALING_MARKER_PATH).exists(), + "lock contention must not publish a healing marker" + ); + + let mut request = HealRequest::new( + HealType::ErasureSet { + buckets: vec!["bucket".to_string()], + set_disk_id: "pool_0_set_0".to_string(), + }, + HealOptions { + timeout: Some(Duration::from_secs(60)), + ..HealOptions::default() + }, + HealPriority::Low, + ); + request.source = HealRequestSource::AutoHeal; + let original = request.clone(); + let storage: Arc = Arc::new(MockStorage); + for attempt in 1..=MAX_RECOVERABLE_HEAL_RETRIES { + let task = HealTask::from_request(request, storage.clone()); + let (retry, delay, _) = retry_request_for_result_with_budget(&task, &result) + .await + .expect("marker CAS contention should retain the existing task retry budget"); + assert_eq!(retry.id, original.id); + assert_eq!(retry.heal_type, original.heal_type); + assert_eq!(retry.source, original.source); + assert_eq!(retry.options.timeout, original.options.timeout); + assert_eq!(retry.retry_attempts, attempt); + assert!(delay > Duration::ZERO); + request = retry; + } + let exhausted = HealTask::from_request(request, storage); + assert!(retry_request_for_result_with_budget(&exhausted, &result).await.is_none()); + + drop(lock); + super::super::apply_healing_markers_to_targets(vec![disk], Some("owner"), None, false) + .await + .expect("the marker CAS should succeed once contention ends"); + assert_eq!( + std::fs::read(metadata.join(super::super::HEALING_MARKER_PATH)).expect("published marker should be readable"), + b"owner" + ); +} + #[tokio::test] async fn retry_request_for_result_preserves_remaining_timeout_budget() { let storage: Arc = Arc::new(MockStorage);