From 902cf1d3dccbf3e773e25c40ca41afe53d3e10cf Mon Sep 17 00:00:00 2001 From: cxymds Date: Thu, 2 Jul 2026 16:34:56 +0800 Subject: [PATCH] fix: handle scanner task failures (#4185) --- crates/scanner/src/scanner_io.rs | 74 ++++++++++++++++++++++++++++++-- 1 file changed, 70 insertions(+), 4 deletions(-) diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 07e332730..41f1fb5c8 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -445,6 +445,10 @@ fn record_set_scan_failure(first_err: &mut Option, err: Error) { } } +fn scanner_task_join_error(stage: &str, err: tokio::task::JoinError) -> Error { + Error::other(format!("{stage} task join failed: {err}")) +} + fn finalize_nsscanner_result(results: &[DataUsageCache], first_err: Option) -> Result<()> { if results.iter().any(|result| result.info.last_update.is_some()) { return Ok(()); @@ -906,7 +910,21 @@ impl ScannerIO for ECStore { } }); - let _ = join_all(wait_futs).await; + for join_result in join_all(wait_futs).await { + if let Err(err) = join_result { + error!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_SET_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + state = "set_task_join_failed", + error = %err, + "Scanner set task join failed" + ); + let mut first_err = first_err_mutex.lock().await; + record_set_scan_failure(&mut first_err, scanner_task_join_error("scanner set", err)); + } + } record_set_scan_concurrency_limit(0); record_set_scans_queued(0); record_set_scans_active(0); @@ -993,7 +1011,20 @@ impl ScannerIOCache for SetDisks { let _reset_disk_bucket_scan_gauges = DiskBucketScanGaugeReset::new(pool_label.clone(), set_label.clone()); let mut old_cache = DataUsageCache::default(); - old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await?; + if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await { + warn!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_CACHE_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + pool = self.pool_index, + set = self.set_index, + cache_name = DATA_USAGE_CACHE_NAME, + state = "old_cache_load_failed", + error = %e, + "Scanner old data usage cache load failed; rebuilding from bucket caches" + ); + } let mut cache = DataUsageCache { info: DataUsageCacheInfo { @@ -1334,14 +1365,36 @@ impl ScannerIOCache for SetDisks { })); } - let _ = join_all(futs).await; + let mut first_join_err = None; + for join_result in join_all(futs).await { + if let Err(err) = join_result { + error!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_DISK_BUCKET_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + pool = self.pool_index, + set = self.set_index, + state = "disk_bucket_task_join_failed", + error = %err, + "Scanner disk bucket task join failed" + ); + record_set_scan_failure(&mut first_join_err, scanner_task_join_error("scanner disk bucket", err)); + } + } record_disk_scan_concurrency_limit(&pool_label, &set_label, 0); record_disk_bucket_scans_queued(0, &pool_label, &set_label); record_disk_bucket_scans_active(0, &pool_label, &set_label); drop(bucket_result_tx_clone); - collect_bucket_results_fut.await?; + if let Err(err) = collect_bucket_results_fut.await { + return Err(scanner_task_join_error("scanner bucket result collector", err)); + } + + if let Some(err) = first_join_err { + return Err(err); + } let completed_count = completed_bucket_count.load(Ordering::Relaxed); if should_publish_completed_snapshot(completed_count, buckets.len(), budget.budget_elapsed(), ctx.is_cancelled()) { @@ -1746,6 +1799,19 @@ mod tests { assert!(first.to_string().contains("first")); } + #[tokio::test] + async fn scanner_task_join_error_includes_stage() { + let handle = tokio::spawn(async { + tokio::time::sleep(Duration::from_secs(60)).await; + }); + handle.abort(); + + let join_err = handle.await.expect_err("aborted task should return a join error"); + let err = scanner_task_join_error("scanner set", join_err); + + assert!(err.to_string().contains("scanner set task join failed")); + } + #[test] fn finalize_nsscanner_result_returns_ok_when_any_set_succeeds() { let mut results = vec![DataUsageCache::default(), DataUsageCache::default()];