From b8b5511b68d3472d9bde0e7dd7a20944dd479675 Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Wed, 23 Jul 2025 10:35:34 +0800 Subject: [PATCH] fix: heal data part lose Signed-off-by: junxiang Mu <1948535941@qq.com> --- crates/ahm/src/heal/manager.rs | 4 ++++ crates/ahm/src/scanner/data_scanner.rs | 14 +++++++++++--- crates/ahm/src/scanner/metrics.rs | 22 ++++++++++++++++++++++ crates/ecstore/src/set_disk.rs | 8 ++++++++ crates/ecstore/src/sets.rs | 7 +++++++ crates/ecstore/src/store.rs | 7 +++++++ 6 files changed, 59 insertions(+), 3 deletions(-) diff --git a/crates/ahm/src/heal/manager.rs b/crates/ahm/src/heal/manager.rs index 358ad3668..cb2418838 100644 --- a/crates/ahm/src/heal/manager.rs +++ b/crates/ahm/src/heal/manager.rs @@ -188,6 +188,10 @@ impl HealManager { } /// Get task progress + pub async fn get_active_tasks_count(&self) -> usize { + self.active_heals.lock().await.len() + } + pub async fn get_task_progress(&self, task_id: &str) -> Result { let active_heals = self.active_heals.lock().await; if let Some(task) = active_heals.get(task_id) { diff --git a/crates/ahm/src/scanner/data_scanner.rs b/crates/ahm/src/scanner/data_scanner.rs index ec20a0b95..98f810b90 100644 --- a/crates/ahm/src/scanner/data_scanner.rs +++ b/crates/ahm/src/scanner/data_scanner.rs @@ -2114,13 +2114,17 @@ mod tests { if final_heal_stats.successful_tasks > 0 { println!("Healing completed successfully, checking file recovery..."); if recovered_files > 0 { - println!("Successfully recovered {}/{} deleted xl.meta files", recovered_files, deleted_meta_paths.len()); + println!( + "Successfully recovered {}/{} deleted xl.meta files", + recovered_files, + deleted_meta_paths.len() + ); } else { println!("No xl.meta files recovered yet - healing may have recreated metadata elsewhere"); } } else { println!("No successful heal tasks completed yet - healing may still be in progress or failed"); - + // If healing failed, this is acceptable for this test scenario // The important thing is that the scanner detected the issue and submitted heal tasks if final_heal_stats.failed_tasks > 0 { @@ -2140,7 +2144,11 @@ mod tests { println!(" - Scanner submitted {} heal tasks", final_heal_stats.total_tasks); println!(" - Scanner handled the situation gracefully"); if recovered_files > 0 { - println!(" - Successfully recovered {}/{} xl.meta files", recovered_files, deleted_meta_paths.len()); + println!( + " - Successfully recovered {}/{} xl.meta files", + recovered_files, + deleted_meta_paths.len() + ); } else { println!(" - Note: xl.meta file recovery may require additional time or manual intervention"); } diff --git a/crates/ahm/src/scanner/metrics.rs b/crates/ahm/src/scanner/metrics.rs index 10d010058..f5d7b73bd 100644 --- a/crates/ahm/src/scanner/metrics.rs +++ b/crates/ahm/src/scanner/metrics.rs @@ -42,6 +42,10 @@ pub struct ScannerMetrics { pub heal_tasks_completed: u64, /// Total heal tasks failed pub heal_tasks_failed: u64, + /// Total healthy objects found + pub healthy_objects: u64, + /// Total corrupted objects found + pub corrupted_objects: u64, /// Last scan activity time pub last_activity: Option, /// Current scan cycle @@ -122,6 +126,8 @@ pub struct MetricsCollector { heal_tasks_failed: AtomicU64, current_cycle: AtomicU64, total_cycles: AtomicU64, + healthy_objects: AtomicU64, + corrupted_objects: AtomicU64, } impl MetricsCollector { @@ -139,6 +145,8 @@ impl MetricsCollector { heal_tasks_failed: AtomicU64::new(0), current_cycle: AtomicU64::new(0), total_cycles: AtomicU64::new(0), + healthy_objects: AtomicU64::new(0), + corrupted_objects: AtomicU64::new(0), } } @@ -197,6 +205,16 @@ impl MetricsCollector { self.total_cycles.fetch_add(1, Ordering::Relaxed); } + /// Increment healthy objects count + pub fn increment_healthy_objects(&self) { + self.healthy_objects.fetch_add(1, Ordering::Relaxed); + } + + /// Increment corrupted objects count + pub fn increment_corrupted_objects(&self) { + self.corrupted_objects.fetch_add(1, Ordering::Relaxed); + } + /// Get current metrics snapshot pub fn get_metrics(&self) -> ScannerMetrics { ScannerMetrics { @@ -209,6 +227,8 @@ impl MetricsCollector { heal_tasks_queued: self.heal_tasks_queued.load(Ordering::Relaxed), heal_tasks_completed: self.heal_tasks_completed.load(Ordering::Relaxed), heal_tasks_failed: self.heal_tasks_failed.load(Ordering::Relaxed), + healthy_objects: self.healthy_objects.load(Ordering::Relaxed), + corrupted_objects: self.corrupted_objects.load(Ordering::Relaxed), last_activity: Some(SystemTime::now()), current_cycle: self.current_cycle.load(Ordering::Relaxed), total_cycles: self.total_cycles.load(Ordering::Relaxed), @@ -234,6 +254,8 @@ impl MetricsCollector { self.heal_tasks_failed.store(0, Ordering::Relaxed); self.current_cycle.store(0, Ordering::Relaxed); self.total_cycles.store(0, Ordering::Relaxed); + self.healthy_objects.store(0, Ordering::Relaxed); + self.corrupted_objects.store(0, Ordering::Relaxed); info!("Scanner metrics reset"); } diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 678480d47..d00ead5c2 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -5950,6 +5950,14 @@ impl StorageAPI for SetDisks { async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> { unimplemented!() } + + #[tracing::instrument(skip(self))] + async fn verify_object_integrity(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> { + let mut get_object_reader = + ::get_object_reader(self, bucket, object, None, HeaderMap::new(), opts).await?; + let _ = get_object_reader.read_all().await?; + Ok(()) + } } #[derive(Debug, PartialEq, Eq)] diff --git a/crates/ecstore/src/sets.rs b/crates/ecstore/src/sets.rs index bea7d9232..40e692ae9 100644 --- a/crates/ecstore/src/sets.rs +++ b/crates/ecstore/src/sets.rs @@ -875,6 +875,13 @@ impl StorageAPI for Sets { async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> { unimplemented!() } + + #[tracing::instrument(skip(self))] + async fn verify_object_integrity(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> { + self.get_disks_by_key(object) + .verify_object_integrity(bucket, object, opts) + .await + } } async fn _close_storage_disks(disks: &[Option]) { diff --git a/crates/ecstore/src/store.rs b/crates/ecstore/src/store.rs index 765415fb8..88cffaeca 100644 --- a/crates/ecstore/src/store.rs +++ b/crates/ecstore/src/store.rs @@ -2235,6 +2235,13 @@ impl StorageAPI for ECStore { Ok(()) } + + async fn verify_object_integrity(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> { + let mut get_object_reader = + ::get_object_reader(self, bucket, object, None, HeaderMap::new(), opts).await?; + let _ = get_object_reader.read_all().await?; + Ok(()) + } } async fn init_local_peer(endpoint_pools: &EndpointServerPools, host: &String, port: &String) {