From 5d03414edb245aae114e402070322198e50349ca Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Sat, 22 Aug 2026 15:33:36 +0800 Subject: [PATCH] fix(scanner): cancel scan workers with cycle scope --- crates/scanner/src/scanner_budget.rs | 40 +++++++++++++++++++++-- crates/scanner/src/scanner_io.rs | 1 + crates/scanner/src/scanner_io/io_cache.rs | 8 ++--- crates/scanner/src/scanner_io/io_cycle.rs | 4 +-- 4 files changed, 45 insertions(+), 8 deletions(-) diff --git a/crates/scanner/src/scanner_budget.rs b/crates/scanner/src/scanner_budget.rs index b0f7df7c8..65743439f 100644 --- a/crates/scanner/src/scanner_budget.rs +++ b/crates/scanner/src/scanner_budget.rs @@ -255,7 +255,7 @@ impl ScannerCycleBudget { } if self .max_directories - .is_some_and(|max_directories| directories > max_directories) + .is_some_and(|max_directories| directory_budget_exhausted(directories, max_directories)) { self.cancel_for(ScannerCycleBudgetReason::Directories); } @@ -281,7 +281,7 @@ impl ScannerCycleBudget { } if self .max_directories - .is_some_and(|max_directories| directories > max_directories) + .is_some_and(|max_directories| directory_budget_exhausted(directories, max_directories)) { self.cancel_for(ScannerCycleBudgetReason::Directories); return false; @@ -329,6 +329,13 @@ fn saturating_fetch_add(value: &AtomicU64, delta: u64) -> u64 { } } +fn directory_budget_exhausted(directories: u64, max_directories: u64) -> bool { + // Saturation hides a remote max+1 update when the configured limit is the + // largest representable counter. Treat that boundary as exhausted rather + // than allowing work to continue indefinitely. + directories > max_directories || (directories == u64::MAX && max_directories == u64::MAX) +} + impl Drop for ScannerCycleBudget { fn drop(&mut self) { self.token.cancel(); @@ -471,6 +478,35 @@ mod tests { assert_eq!(directory_budget.reason(), Some(ScannerCycleBudgetReason::Directories)); } + #[test] + fn directory_budget_fails_closed_when_progress_saturates() { + let parent = CancellationToken::new(); + let budget = ScannerCycleBudget::new( + &parent, + ScannerCycleBudgetConfig { + max_directories: Some(u64::MAX), + ..Default::default() + }, + ); + + budget.record_remote_progress(0, u64::MAX); + + assert_eq!(budget.reason(), Some(ScannerCycleBudgetReason::Directories)); + assert!(budget.token().is_cancelled()); + + let local_budget = ScannerCycleBudget::new( + &parent, + ScannerCycleBudgetConfig { + max_directories: Some(u64::MAX), + ..Default::default() + }, + ); + local_budget.record_remote_progress(0, u64::MAX - 1); + assert!(!local_budget.budget_elapsed()); + assert!(!local_budget.try_start_directory()); + assert_eq!(local_budget.reason(), Some(ScannerCycleBudgetReason::Directories)); + } + #[test] fn explicit_progress_tracking_counts_unbounded_remote_work_without_cancelling() { let parent = CancellationToken::new(); diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 677719a69..c104fb263 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -48,6 +48,7 @@ use time::OffsetDateTime; use tokio::sync::{Mutex, Notify, Semaphore, mpsc}; use tokio::time::Duration; use tokio_util::sync::CancellationToken; +use tokio_util::task::AbortOnDropHandle; use tracing::{debug, error, warn}; use crate::ScannerObjectInfo as ObjectInfo; diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index 3151b4c0c..aa0339e4b 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -314,7 +314,7 @@ impl ScannerIOCache for SetDisks { let ctx_clone = ctx.clone(); let completed_bucket_count = Arc::new(AtomicUsize::new(0)); let completed_bucket_count_clone = completed_bucket_count.clone(); - let collect_bucket_results_fut = tokio::spawn(async move { + let collect_bucket_results_fut = AbortOnDropHandle::new(tokio::spawn(async move { let mut cancelled = false; loop { @@ -333,7 +333,7 @@ impl ScannerIOCache for SetDisks { } } } - }); + })); let mut futs = Vec::new(); @@ -365,7 +365,7 @@ impl ScannerIOCache for SetDisks { NamespaceScannerWorkerMode::RemoteV4(server_epoch) => Some(server_epoch), NamespaceScannerWorkerMode::Coordinator => None, }; - futs.push(tokio::spawn(async move { + futs.push(AbortOnDropHandle::new(tokio::spawn(async move { let remote_session_id = uuid::Uuid::new_v4(); let mut remote_session_sequence = 0_u64; loop { @@ -1038,7 +1038,7 @@ impl ScannerIOCache for SetDisks { ); } } - })); + }))); } drop(bucket_tx); drop(bucket_result_tx); diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index c2a8eb253..64763655f 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -242,7 +242,7 @@ impl ScannerIOCycle for ECStore { results[results_index_clone] = result; } }); - wait_futs.push(receiver_fut); + wait_futs.push(AbortOnDropHandle::new(receiver_fut)); let scan_plan = ScannerBucketScanPlan { buckets: set_buckets, @@ -318,7 +318,7 @@ impl ScannerIOCycle for ECStore { record_set_scan_failure(&mut first_err, e); } }); - wait_futs.push(scanner_fut); + wait_futs.push(AbortOnDropHandle::new(scanner_fut)); } }