mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-30 08:49:26 +00:00
@@ -16,7 +16,7 @@ use crate::error::{Error, Result};
|
|||||||
use crate::heal::{
|
use crate::heal::{
|
||||||
progress::{HealProgress, HealStatistics},
|
progress::{HealProgress, HealStatistics},
|
||||||
storage::HealStorageAPI,
|
storage::HealStorageAPI,
|
||||||
task::{HealRequest, HealTask, HealTaskStatus},
|
task::{HealOptions, HealPriority, HealRequest, HealTask, HealTaskStatus, HealType},
|
||||||
};
|
};
|
||||||
use rustfs_ecstore::disk::error::DiskError;
|
use rustfs_ecstore::disk::error::DiskError;
|
||||||
use rustfs_ecstore::disk::DiskAPI;
|
use rustfs_ecstore::disk::DiskAPI;
|
||||||
@@ -332,14 +332,14 @@ impl HealManager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// enqueue erasure set heal request for this disk
|
// enqueue erasure set heal request for this disk
|
||||||
let set_disk_id = format!("{}_{}", ep.pool_idx, ep.set_idx);
|
let set_disk_id = format!("pool_{}_set_{}", ep.pool_idx, ep.set_idx);
|
||||||
let req = crate::heal::task::HealRequest::new(
|
let req = HealRequest::new(
|
||||||
crate::heal::task::HealType::ErasureSet {
|
HealType::ErasureSet {
|
||||||
buckets: buckets.clone(),
|
buckets: buckets.clone(),
|
||||||
set_disk_id: set_disk_id.clone()
|
set_disk_id: set_disk_id.clone()
|
||||||
},
|
},
|
||||||
crate::heal::task::HealOptions::default(),
|
HealOptions::default(),
|
||||||
crate::heal::task::HealPriority::Normal,
|
HealPriority::Normal,
|
||||||
);
|
);
|
||||||
let mut queue = heal_queue.lock().await;
|
let mut queue = heal_queue.lock().await;
|
||||||
queue.push_back(req);
|
queue.push_back(req);
|
||||||
|
|||||||
@@ -795,6 +795,13 @@ impl HealTask {
|
|||||||
progress.update_progress(2, 4, 0, 0);
|
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
|
// Step 3: Create erasure set healer with resume support
|
||||||
info!("Step 3: Creating 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);
|
let erasure_healer = ErasureSetHealer::new(self.storage.clone(), self.progress.clone(), self.cancel_token.clone(), disk);
|
||||||
|
|||||||
@@ -507,29 +507,24 @@ impl Scanner {
|
|||||||
let enable_healing = self.config.read().await.enable_healing;
|
let enable_healing = self.config.read().await.enable_healing;
|
||||||
if enable_healing {
|
if enable_healing {
|
||||||
if let Some(heal_manager) = &self.heal_manager {
|
if let Some(heal_manager) = &self.heal_manager {
|
||||||
// Get bucket list for erasure set healing
|
// Get bucket list for erasure set healing
|
||||||
let buckets = match rustfs_ecstore::new_object_layer_fn() {
|
let buckets = match rustfs_ecstore::new_object_layer_fn() {
|
||||||
Some(ecstore) => {
|
Some(ecstore) => match ecstore.list_bucket(&ecstore::store_api::BucketOptions::default()).await {
|
||||||
match ecstore.list_bucket(&ecstore::store_api::BucketOptions::default()).await {
|
|
||||||
Ok(buckets) => buckets.iter().map(|b| b.name.clone()).collect::<Vec<String>>(),
|
Ok(buckets) => buckets.iter().map(|b| b.name.clone()).collect::<Vec<String>>(),
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
error!("Failed to get bucket list for disk healing: {}", 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::HealOptions::default(),
|
||||||
crate::heal::task::HealPriority::High,
|
crate::heal::task::HealPriority::High,
|
||||||
);
|
);
|
||||||
@@ -571,27 +566,22 @@ impl Scanner {
|
|||||||
if let Some(heal_manager) = &self.heal_manager {
|
if let Some(heal_manager) = &self.heal_manager {
|
||||||
// Get bucket list for erasure set healing
|
// Get bucket list for erasure set healing
|
||||||
let buckets = match rustfs_ecstore::new_object_layer_fn() {
|
let buckets = match rustfs_ecstore::new_object_layer_fn() {
|
||||||
Some(ecstore) => {
|
Some(ecstore) => match ecstore.list_bucket(&ecstore::store_api::BucketOptions::default()).await {
|
||||||
match ecstore.list_bucket(&ecstore::store_api::BucketOptions::default()).await {
|
Ok(buckets) => buckets.iter().map(|b| b.name.clone()).collect::<Vec<String>>(),
|
||||||
Ok(buckets) => buckets.iter().map(|b| b.name.clone()).collect::<Vec<String>>(),
|
Err(e) => {
|
||||||
Err(e) => {
|
error!("Failed to get bucket list for disk healing: {}", e);
|
||||||
error!("Failed to get bucket list for disk healing: {}", e);
|
return Err(Error::Storage(e));
|
||||||
return Err(Error::Storage(e.into()));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
},
|
||||||
None => {
|
None => {
|
||||||
error!("No ECStore available for getting bucket list");
|
error!("No ECStore available for getting bucket list");
|
||||||
return Err(Error::Storage(ecstore::error::StorageError::other("No ECStore available")));
|
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(
|
let req = HealRequest::new(
|
||||||
crate::heal::task::HealType::ErasureSet {
|
crate::heal::task::HealType::ErasureSet { buckets, set_disk_id },
|
||||||
buckets,
|
|
||||||
set_disk_id,
|
|
||||||
},
|
|
||||||
crate::heal::task::HealOptions::default(),
|
crate::heal::task::HealOptions::default(),
|
||||||
crate::heal::task::HealPriority::Urgent,
|
crate::heal::task::HealPriority::Urgent,
|
||||||
);
|
);
|
||||||
|
|||||||
@@ -303,6 +303,54 @@ async fn test_heal_format_basic() {
|
|||||||
info!("Heal format basic test passed");
|
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)]
|
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn test_heal_storage_api_direct() {
|
async fn test_heal_storage_api_direct() {
|
||||||
|
|||||||
Reference in New Issue
Block a user