mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-02 05:14:35 +08:00
fix(heal): retry contended local metadata CAS (#8267)
This commit is contained in:
@@ -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;
|
||||
|
||||
@@ -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<dyn HealStorageAPI> = 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<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
|
||||
Reference in New Issue
Block a user