From 56a4099f99150f4f3b63a30746f28e7bbb2aba28 Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 23 Sep 2026 14:51:24 +0800 Subject: [PATCH] fix(heal): recreate missing bucket volumes (#8082) * fix(heal): recreate missing bucket volumes * fix(heal): use typed error for empty targets --- crates/ecstore/src/set_disk/mod.rs | 11 +- crates/ecstore/src/set_disk/ops/heal.rs | 158 ++++++++++++++++++++++++ crates/ecstore/src/store/heal.rs | 9 +- 3 files changed, 172 insertions(+), 6 deletions(-) diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 34460dd64..0a0225654 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -6809,7 +6809,10 @@ pub fn should_heal_object_on_disk( latest_meta: &FileInfo, ) -> (bool, bool, Option) { if let Some(err) = err - && (err == &DiskError::FileNotFound || err == &DiskError::FileVersionNotFound || err == &DiskError::FileCorrupt) + && (err == &DiskError::FileNotFound + || err == &DiskError::FileVersionNotFound + || err == &DiskError::FileCorrupt + || err == &DiskError::VolumeNotFound) { return (true, true, Some(err.clone())); } @@ -11263,6 +11266,12 @@ mod tests { let (should_heal, _, _) = should_heal_object_on_disk(&err, &[], &meta, &latest_meta); assert!(should_heal); + let err = Some(DiskError::VolumeNotFound); + let (should_heal, is_meta, reason) = should_heal_object_on_disk(&err, &[], &meta, &latest_meta); + assert!(should_heal); + assert!(is_meta); + assert_eq!(reason, Some(DiskError::VolumeNotFound)); + let err = Some(DiskError::FileCorrupt); let (should_heal, is_meta, reason) = should_heal_object_on_disk(&err, &[], &meta, &latest_meta); assert!(should_heal); diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index e469826d5..6c0547b29 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -1398,6 +1398,26 @@ impl SetDisks { )); } + // `errs` is reported in physical-disk order while the + // reconstruction loop below uses the shard order from + // the selected metadata. Preserve which targets are + // missing their bucket volume across that permutation. + let mut missing_bucket_volumes = vec![false; out_dated_disks.len()]; + for (physical_index, error) in errs.iter().enumerate() { + if error.as_ref() != Some(&DiskError::VolumeNotFound) { + continue; + } + let shard_index = latest_meta + .erasure + .distribution + .get(physical_index) + .and_then(|index| index.checked_sub(1)) + .unwrap_or(physical_index); + if let Some(target) = missing_bucket_volumes.get_mut(shard_index) { + *target = true; + } + } + out_dated_disks = Self::shuffle_disks(&out_dated_disks, &latest_meta.erasure.distribution); let mut parts_metadata = Self::shuffle_parts_metadata(&parts_metadata, &latest_meta.erasure.distribution); let mut copy_parts_metadata = vec![None; parts_metadata.len()]; @@ -1424,6 +1444,58 @@ impl SetDisks { } } + // A missing bucket volume is a repairable object target, + // but the final rename cannot create its parent volume. + // Recreate only the affected volumes after quorum and + // lifecycle checks have passed. Keep failures visible so + // a partial repair cannot be reported as successful. + let mut volume_creation_error = None; + for (index, outdated_disk) in out_dated_disks.iter_mut().enumerate() { + if !missing_bucket_volumes.get(index).copied().unwrap_or(false) || outdated_disk.is_none() { + continue; + } + if let Some(scope) = &bucket_heal_scope { + scope.check()?; + } + let disk = outdated_disk.as_ref().expect("checked above").clone(); + let endpoint = disk.endpoint().to_string(); + match disk.make_volume(bucket).await { + Ok(()) | Err(DiskError::VolumeExists) => {} + Err(error) => { + warn!( + event = EVENT_SET_DISK_HEAL, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_HEAL, + bucket, + object, + version_id, + disk_index = index, + endpoint = %endpoint, + error = %error, + state = "bucket_volume_recreate_failed", + "Heal object could not recreate the target bucket volume" + ); + if volume_creation_error.is_none() { + volume_creation_error = Some(error.clone()); + } + *outdated_disk = None; + disks_to_heal_count = disks_to_heal_count.saturating_sub(1); + for drive in &mut result.after.drives { + if drive.endpoint == endpoint { + drive.state = heal_drive_state_for_error(&error).to_string(); + } + } + } + } + } + + if let Some(scope) = &bucket_heal_scope { + scope.check()?; + } + if disks_to_heal_count == 0 { + return Ok((result, volume_creation_error.or(Some(DiskError::ErasureWriteQuorum)))); + } + // We write at temporary location and then rename to final location. let tmp_id = Uuid::new_v4().to_string(); // Delete markers and remote (transitioned) objects carry no data_dir and @@ -1865,6 +1937,10 @@ impl SetDisks { )); } + if let Some(error) = volume_creation_error { + return Ok((result, Some(error))); + } + if result.integrity_verified && let Err(error) = crate::io_support::shard_integrity::restore_proof_replicas(&latest_meta, &disks, bucket, object) @@ -4223,6 +4299,88 @@ mod heal_result_report_tests { ); } + #[tokio::test] + async fn deep_heal_recreates_missing_bucket_volume_before_rebuilding_shard() { + let (temp_dirs, disks, set) = hermetic_set_disks_isolated(4).await; + let bucket = "deep-heal-missing-bucket-volume"; + let object = "object.bin"; + for disk in &disks { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + + let mut reader = PutObjReader::from_vec(vec![0x5c; 1024 * 1024]); + set.put_object( + bucket, + object, + &mut reader, + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + .expect("source object should be written before volume loss"); + let source = disks[2] + .read_version("", bucket, object, "", &ReadOptions::default()) + .await + .expect("source metadata should be readable"); + let data_dir = source.data_dir.expect("non-inline source should have a data directory"); + + disks[1] + .delete_volume(bucket, true) + .await + .expect("target bucket volume should be removable for the regression setup"); + assert!(matches!(disks[1].stat_volume(bucket).await, Err(DiskError::VolumeNotFound))); + + let (result, error) = set + .heal_object( + bucket, + object, + "", + &HealOpts { + no_lock: true, + scan_mode: HealScanMode::Deep, + ..Default::default() + }, + ) + .await + .expect("deep heal should finish after bucket volume loss"); + + assert!(error.is_none(), "deep heal should recreate the missing bucket volume: {error:?}"); + assert!(disks[1].stat_volume(bucket).await.is_ok(), "target bucket volume should be restored"); + let healed = disks[1] + .read_version("", bucket, object, "", &ReadOptions::default()) + .await + .expect("healed target metadata should be readable"); + assert_eq!(healed.data_dir, Some(data_dir)); + assert_eq!(result.after.drives[1].state, DriveState::Ok.to_string()); + assert!( + temp_dirs[1] + .path() + .join(bucket) + .join(object) + .join(data_dir.to_string()) + .join("part.1") + .exists(), + "deep heal should rebuild the object shard after recreating its volume" + ); + + let (_, retry_error) = set + .heal_object( + bucket, + object, + "", + &HealOpts { + no_lock: true, + scan_mode: HealScanMode::Deep, + ..Default::default() + }, + ) + .await + .expect("repeat heal should complete"); + assert!(retry_error.is_none(), "repeat heal should be idempotent: {retry_error:?}"); + } + #[tokio::test] async fn deep_heal_rebuilds_a_stale_current_delete_marker_metadata_replica() { let (_temp_dirs, disks, set) = hermetic_set_disks_isolated(4).await; diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index 216fc7778..8e20fa184 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -857,11 +857,10 @@ impl ECStore { opts: &HealOpts, allow_unversioned_absence: bool, ) -> Result { - if opts.dry_run - || opts.no_lock - || (!allow_unversioned_absence && version_id.is_empty()) - || super::utils::is_reserved_or_invalid_bucket(bucket, false) - { + // Object heal may recreate a missing bucket volume before committing a + // shard. Keep the current bucket generation fenced for the whole + // operation, including the ordinary unversioned path. + if opts.dry_run || opts.no_lock || super::utils::is_reserved_or_invalid_bucket(bucket, false) { let (item, error) = self.handle_heal_object(bucket, object, version_id, opts).await?; return Ok(HealObjectStorageResult { item,