diff --git a/crates/ahm/src/heal/manager.rs b/crates/ahm/src/heal/manager.rs index 0fd8c7ab5..180eaf665 100644 --- a/crates/ahm/src/heal/manager.rs +++ b/crates/ahm/src/heal/manager.rs @@ -16,7 +16,7 @@ use crate::error::{Error, Result}; use crate::heal::{ progress::{HealProgress, HealStatistics}, storage::HealStorageAPI, - task::{HealRequest, HealTask, HealTaskStatus}, + task::{HealOptions, HealPriority, HealRequest, HealTask, HealTaskStatus, HealType}, }; use rustfs_ecstore::disk::error::DiskError; use rustfs_ecstore::disk::DiskAPI; @@ -332,14 +332,14 @@ impl HealManager { } // enqueue erasure set heal request for this disk - let set_disk_id = format!("{}_{}", ep.pool_idx, ep.set_idx); - let req = crate::heal::task::HealRequest::new( - crate::heal::task::HealType::ErasureSet { - buckets: buckets.clone(), - set_disk_id: set_disk_id.clone() + let set_disk_id = format!("pool_{}_set_{}", ep.pool_idx, ep.set_idx); + let req = HealRequest::new( + HealType::ErasureSet { + buckets: buckets.clone(), + set_disk_id: set_disk_id.clone() }, - crate::heal::task::HealOptions::default(), - crate::heal::task::HealPriority::Normal, + HealOptions::default(), + HealPriority::Normal, ); let mut queue = heal_queue.lock().await; queue.push_back(req); diff --git a/crates/ahm/src/heal/task.rs b/crates/ahm/src/heal/task.rs index 44f12daca..8e9d856d7 100644 --- a/crates/ahm/src/heal/task.rs +++ b/crates/ahm/src/heal/task.rs @@ -795,6 +795,13 @@ impl HealTask { progress.update_progress(2, 4, 0, 0); } + // Step 3: Heal bucket structure + for bucket in buckets.iter() { + if let Err(err) = self.heal_bucket(bucket).await { + info!("{}", err.to_string()); + } + } + // Step 3: Create erasure set healer with resume support info!("Step 3: Creating erasure set healer with resume support"); let erasure_healer = ErasureSetHealer::new(self.storage.clone(), self.progress.clone(), self.cancel_token.clone(), disk); diff --git a/crates/ahm/src/scanner/data_scanner.rs b/crates/ahm/src/scanner/data_scanner.rs index 9c5e1aa1c..0496f70e1 100644 --- a/crates/ahm/src/scanner/data_scanner.rs +++ b/crates/ahm/src/scanner/data_scanner.rs @@ -507,29 +507,24 @@ impl Scanner { let enable_healing = self.config.read().await.enable_healing; if enable_healing { if let Some(heal_manager) = &self.heal_manager { - // Get bucket list for erasure set healing - let buckets = match rustfs_ecstore::new_object_layer_fn() { - Some(ecstore) => { - match ecstore.list_bucket(&ecstore::store_api::BucketOptions::default()).await { + // Get bucket list for erasure set healing + let buckets = match rustfs_ecstore::new_object_layer_fn() { + Some(ecstore) => match ecstore.list_bucket(&ecstore::store_api::BucketOptions::default()).await { Ok(buckets) => buckets.iter().map(|b| b.name.clone()).collect::>(), Err(e) => { error!("Failed to get bucket list for disk healing: {}", e); - return Err(Error::Storage(e.into())); + return Err(Error::Storage(e)); } - } - } - None => { - error!("No ECStore available for getting bucket list"); - return Err(Error::Storage(ecstore::error::StorageError::other("No ECStore available"))); - } - }; - - let set_disk_id = format!("{}_{}", disk.endpoint().pool_idx, disk.endpoint().set_idx); - let req = HealRequest::new( - crate::heal::task::HealType::ErasureSet { - buckets, - set_disk_id, }, + None => { + error!("No ECStore available for getting bucket list"); + return Err(Error::Storage(ecstore::error::StorageError::other("No ECStore available"))); + } + }; + + let set_disk_id = format!("pool_{}_set_{}", disk.endpoint().pool_idx, disk.endpoint().set_idx); + let req = HealRequest::new( + crate::heal::task::HealType::ErasureSet { buckets, set_disk_id }, crate::heal::task::HealOptions::default(), crate::heal::task::HealPriority::High, ); @@ -571,27 +566,22 @@ impl Scanner { if let Some(heal_manager) = &self.heal_manager { // Get bucket list for erasure set healing let buckets = match rustfs_ecstore::new_object_layer_fn() { - Some(ecstore) => { - match ecstore.list_bucket(&ecstore::store_api::BucketOptions::default()).await { - Ok(buckets) => buckets.iter().map(|b| b.name.clone()).collect::>(), - Err(e) => { - error!("Failed to get bucket list for disk healing: {}", e); - return Err(Error::Storage(e.into())); - } + Some(ecstore) => match ecstore.list_bucket(&ecstore::store_api::BucketOptions::default()).await { + Ok(buckets) => buckets.iter().map(|b| b.name.clone()).collect::>(), + Err(e) => { + error!("Failed to get bucket list for disk healing: {}", e); + return Err(Error::Storage(e)); } - } + }, None => { error!("No ECStore available for getting bucket list"); return Err(Error::Storage(ecstore::error::StorageError::other("No ECStore available"))); } }; - let set_disk_id = format!("{}_{}", disk.endpoint().pool_idx, disk.endpoint().set_idx); + let set_disk_id = format!("pool_{}_set_{}", disk.endpoint().pool_idx, disk.endpoint().set_idx); let req = HealRequest::new( - crate::heal::task::HealType::ErasureSet { - buckets, - set_disk_id, - }, + crate::heal::task::HealType::ErasureSet { buckets, set_disk_id }, crate::heal::task::HealOptions::default(), crate::heal::task::HealPriority::Urgent, ); diff --git a/crates/ahm/tests/heal_integration_test.rs b/crates/ahm/tests/heal_integration_test.rs index c5a6f3a5b..76603e53b 100644 --- a/crates/ahm/tests/heal_integration_test.rs +++ b/crates/ahm/tests/heal_integration_test.rs @@ -303,6 +303,54 @@ async fn test_heal_format_basic() { info!("Heal format basic test passed"); } +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +#[serial] +async fn test_heal_format_with_data() { + let (disk_paths, ecstore, heal_storage) = setup_test_env().await; + + // Create test bucket and object + let bucket_name = "test-bucket"; + let object_name = "test-object.txt"; + let test_data = b"Hello, this is test data for healing!"; + + create_test_bucket(&ecstore, bucket_name).await; + upload_test_object(&ecstore, bucket_name, object_name, test_data).await; + + let obj_dir = disk_paths[0].join(bucket_name).join(object_name); + let target_part = WalkDir::new(&obj_dir) + .min_depth(2) + .max_depth(2) + .into_iter() + .filter_map(Result::ok) + .find(|e| e.file_type().is_file() && e.file_name().to_str().map(|n| n.starts_with("part.")).unwrap_or(false)) + .map(|e| e.into_path()) + .expect("Failed to locate part file to delete"); + + // ─── 1️⃣ delete format.json on one disk ────────────── + let format_path = disk_paths[0].join(".rustfs.sys").join("format.json"); + std::fs::remove_dir_all(&disk_paths[0]).expect("failed to delete all contents under disk_paths[0]"); + std::fs::create_dir_all(&disk_paths[0]).expect("failed to recreate disk_paths[0] directory"); + println!("✅ Deleted format.json on disk: {:?}", disk_paths[0]); + + // Create heal manager with faster interval + let cfg = HealConfig { + heal_interval: Duration::from_secs(2), + ..Default::default() + }; + let heal_manager = HealManager::new(heal_storage.clone(), Some(cfg)); + heal_manager.start().await.unwrap(); + + // Wait for task completion + tokio::time::sleep(tokio::time::Duration::from_secs(5)).await; + + // ─── 2️⃣ verify format.json is restored ─────── + assert!(format_path.exists(), "format.json does not exist on disk after heal"); + // ─── 3 verify each part file is restored ─────── + assert!(target_part.exists()); + + info!("Heal format basic test passed"); +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] async fn test_heal_storage_api_direct() {