From d8bf26888511b1aa3121efd1cdf61aa2b64cc4a6 Mon Sep 17 00:00:00 2001 From: Hauser Date: Sat, 26 Sep 2026 14:06:50 +0800 Subject: [PATCH] fix(scanner): bound SNSD deep scans and checkpoint cloning (#8126) * fix(scanner): bound SNSD deep scans and checkpoint cloning * ci: run mount-dependent jobs on hosted VMs --------- Co-authored-by: hector <42570491+majinghe@users.noreply.github.com> --- .github/workflows/ci.yml | 2 +- .github/workflows/e2e-s3tests.yml | 2 +- crates/scanner/src/scanner.rs | 49 +++++++++++++---- crates/scanner/src/scanner/tests.rs | 61 ++++++++++++++++++---- crates/scanner/src/scanner_folder.rs | 18 +++++-- crates/scanner/src/scanner_folder/tests.rs | 10 ++++ 6 files changed, 115 insertions(+), 27 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2eae1ef5c..08e386509 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -690,7 +690,7 @@ jobs: # the probe succeeds and the tests exercise the real UringBackend/FdCache/ # latch paths instead of the StdBackend fallback (rustfs/backlog#1179). The # self-hosted sm-standard runners cannot guarantee this. - runs-on: sm-standard-2 + runs-on: ubuntu-latest timeout-minutes: 30 steps: - name: Checkout repository diff --git a/.github/workflows/e2e-s3tests.yml b/.github/workflows/e2e-s3tests.yml index 448b77150..859e0e0ed 100644 --- a/.github/workflows/e2e-s3tests.yml +++ b/.github/workflows/e2e-s3tests.yml @@ -146,7 +146,7 @@ jobs: # GitHub-hosted: reliably provides Docker + docker compose + python3/pip. # See the header note (ci-1) for why the self-hosted sm-standard-4 label # was abandoned. Scheduled failures are handled by alert-on-failure below. - runs-on: sm-standard-2 + runs-on: ubuntu-latest timeout-minutes: 180 strategy: fail-fast: false diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index fcd313ea9..099811fbe 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -1454,13 +1454,32 @@ fn bitrot_scan_cycle() -> Option { resolve_scanner_runtime_config().bitrot_cycle } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct ScannerBitrotPolicy { + cycle: Option, + deep_window_cycles: Option, +} + +impl ScannerBitrotPolicy { + fn new(cycle: Option, deep_scan_supported: bool, deep_window_cycles: u64) -> Self { + Self { + cycle, + deep_window_cycles: deep_scan_supported.then_some(deep_window_cycles), + } + } +} + fn get_cycle_scan_mode( current_cycle: u64, bitrot_start_cycle: u64, bitrot_start_time: Option>, - bitrot_cycle: Option, + policy: ScannerBitrotPolicy, ) -> HealScanMode { - let Some(bitrot_cycle) = bitrot_cycle else { + let Some(bitrot_cycle) = policy.cycle else { + return HealScanMode::Normal; + }; + + let Some(deep_window_cycles) = policy.deep_window_cycles else { return HealScanMode::Normal; }; @@ -1468,7 +1487,7 @@ fn get_cycle_scan_mode( return HealScanMode::Deep; } - if current_cycle.saturating_sub(bitrot_start_cycle) < heal_object_select_prob() as u64 { + if current_cycle.saturating_sub(bitrot_start_cycle) < deep_window_cycles { return HealScanMode::Deep; } @@ -1492,10 +1511,9 @@ fn background_heal_info_for_scan_start( current_cycle: u64, scan_mode: HealScanMode, now: DateTime, - bitrot_cycle: Option, + policy: ScannerBitrotPolicy, ) -> Option { - let reset_bitrot_start = - scan_mode == HealScanMode::Deep && should_reset_bitrot_start(&info, current_cycle, now, bitrot_cycle); + let reset_bitrot_start = scan_mode == HealScanMode::Deep && should_reset_bitrot_start(&info, current_cycle, now, policy); if info.current_scan_mode == scan_mode && !reset_bitrot_start { return None; } @@ -1513,13 +1531,17 @@ fn should_reset_bitrot_start( info: &BackgroundHealInfo, current_cycle: u64, now: DateTime, - bitrot_cycle: Option, + policy: ScannerBitrotPolicy, ) -> bool { let Some(bitrot_start_time) = info.bitrot_start_time else { return true; }; - let Some(bitrot_cycle) = bitrot_cycle else { + let Some(bitrot_cycle) = policy.cycle else { + return false; + }; + + let Some(deep_window_cycles) = policy.deep_window_cycles else { return false; }; @@ -1527,7 +1549,7 @@ fn should_reset_bitrot_start( return true; } - if current_cycle.saturating_sub(info.bitrot_start_cycle) < heal_object_select_prob() as u64 { + if current_cycle.saturating_sub(info.bitrot_start_cycle) < deep_window_cycles { return false; } @@ -1782,12 +1804,17 @@ where } let mut background_heal_info = background_heal_read.info; let background_heal_epoch = background_heal_read.expected_epoch; + let bitrot_policy = ScannerBitrotPolicy::new( + configured_bitrot_cycle, + !storeapi.setup_is_erasure_sd().await, + heal_object_select_prob() as u64, + ); let scan_mode = get_cycle_scan_mode( cycle_info.current, background_heal_info.bitrot_start_cycle, background_heal_info.bitrot_start_time, - configured_bitrot_cycle, + bitrot_policy, ); info!( target: "rustfs::scanner", @@ -1805,7 +1832,7 @@ where cycle_info.current, scan_mode, Utc::now(), - configured_bitrot_cycle, + bitrot_policy, ) { background_heal_info = new_heal_info.clone(); save_background_heal_info_for_epoch(storeapi.clone(), new_heal_info, background_heal_epoch).await; diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 778ff096f..18eadcf81 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -10015,7 +10015,7 @@ async fn scanner_activity_probe_wait_stops_after_leader_lock_loss() { #[serial] fn test_get_cycle_scan_mode_runs_deep_until_selection_window_completes() { with_var(ENV_SCANNER_BITROT_CYCLE_SECS, Some("3600"), || { - let mode = get_cycle_scan_mode(10, 0, Some(Utc::now()), bitrot_scan_cycle()); + let mode = get_cycle_scan_mode(10, 0, Some(Utc::now()), ScannerBitrotPolicy::new(bitrot_scan_cycle(), true, 1024)); assert_eq!(mode, HealScanMode::Deep); }); } @@ -10027,8 +10027,14 @@ fn test_get_cycle_scan_mode_respects_elapsed_bitrot_cycle() { let recent = Utc::now() - chrono::Duration::minutes(30); let old = Utc::now() - chrono::Duration::hours(2); - assert_eq!(get_cycle_scan_mode(2048, 0, Some(recent), bitrot_scan_cycle()), HealScanMode::Normal); - assert_eq!(get_cycle_scan_mode(2048, 0, Some(old), bitrot_scan_cycle()), HealScanMode::Deep); + assert_eq!( + get_cycle_scan_mode(2048, 0, Some(recent), ScannerBitrotPolicy::new(bitrot_scan_cycle(), true, 1024)), + HealScanMode::Normal + ); + assert_eq!( + get_cycle_scan_mode(2048, 0, Some(old), ScannerBitrotPolicy::new(bitrot_scan_cycle(), true, 1024)), + HealScanMode::Deep + ); }); } @@ -10036,17 +10042,48 @@ fn test_get_cycle_scan_mode_respects_elapsed_bitrot_cycle() { #[serial] fn test_get_cycle_scan_mode_can_disable_periodic_deep_scan() { with_var(ENV_SCANNER_BITROT_CYCLE_SECS, Some("off"), || { - assert_eq!(get_cycle_scan_mode(1, 0, None, bitrot_scan_cycle()), HealScanMode::Normal); + assert_eq!( + get_cycle_scan_mode(1, 0, None, ScannerBitrotPolicy::new(bitrot_scan_cycle(), true, 1)), + HealScanMode::Normal + ); }); } +#[test] +fn test_erasure_sd_does_not_enter_deep_scan_mode() { + let started = Utc::now(); + let cycle = Some(Duration::from_secs(3600)); + + let unsupported_policy = ScannerBitrotPolicy::new(cycle, false, 1024); + assert_eq!(get_cycle_scan_mode(10, 10, Some(started), unsupported_policy), HealScanMode::Normal); + assert_eq!(get_cycle_scan_mode(11, 10, None, unsupported_policy), HealScanMode::Normal); + assert_eq!( + get_cycle_scan_mode(10, 10, None, ScannerBitrotPolicy::new(Some(Duration::ZERO), false, 1024)), + HealScanMode::Normal + ); + + let info = BackgroundHealInfo { + bitrot_start_time: Some(started), + bitrot_start_cycle: 10, + current_scan_mode: HealScanMode::Deep, + }; + let normalized = background_heal_info_for_scan_start(info, 11, HealScanMode::Normal, started, unsupported_policy) + .expect("ErasureSD should persist a legacy Deep state as Normal"); + assert_eq!(normalized.current_scan_mode, HealScanMode::Normal); +} + #[test] #[serial] fn test_background_heal_info_for_scan_start_marks_deep_active() { let now = Utc::now(); - let info = - background_heal_info_for_scan_start(BackgroundHealInfo::default(), 7, HealScanMode::Deep, now, bitrot_scan_cycle()) - .expect("deep scan should update background heal info"); + let info = background_heal_info_for_scan_start( + BackgroundHealInfo::default(), + 7, + HealScanMode::Deep, + now, + ScannerBitrotPolicy::new(bitrot_scan_cycle(), true, 1024), + ) + .expect("deep scan should update background heal info"); assert_eq!(info.current_scan_mode, HealScanMode::Deep); assert_eq!(info.bitrot_start_cycle, 7); @@ -10077,8 +10114,14 @@ fn test_background_heal_info_for_scan_start_keeps_deep_window_start() { current_scan_mode: HealScanMode::Normal, }; - let info = background_heal_info_for_scan_start(info, 8, HealScanMode::Deep, Utc::now(), bitrot_scan_cycle()) - .expect("deep scan should mark active status"); + let info = background_heal_info_for_scan_start( + info, + 8, + HealScanMode::Deep, + Utc::now(), + ScannerBitrotPolicy::new(bitrot_scan_cycle(), true, 1024), + ) + .expect("deep scan should mark active status"); assert_eq!(info.current_scan_mode, HealScanMode::Deep); assert_eq!(info.bitrot_start_cycle, 7); diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index ac7c1abdb..ccb748f11 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -1219,7 +1219,7 @@ impl FolderScanner { } fn maybe_send_checkpoint(&mut self) { - let Some(checkpoint_tx) = self.checkpoint_tx.as_ref() else { return }; + let Some(checkpoint_tx) = self.checkpoint_tx.clone() else { return }; let elapsed = self.last_checkpoint_at.elapsed(); if self.new_cache.info.scan_progress.is_none() || elapsed < SCANNER_CHECKPOINT_MIN_INTERVAL @@ -1232,6 +1232,15 @@ impl FolderScanner { return; } + // Reserve the bounded queue slot before cloning the growing cache. + // When persistence is slower than traversal, constructing snapshots + // that can only be rejected by a full channel creates repeated O(N) + // allocations without improving restart progress. + if checkpoint_tx.is_closed() { + return; + } + let Ok(permit) = checkpoint_tx.try_reserve() else { return }; + let mut snapshot = self.new_cache.clone(); snapshot.info.last_update = Some(SystemTime::now()); snapshot.info.snapshot_complete = false; @@ -1257,10 +1266,9 @@ impl FolderScanner { { return; } - if checkpoint_tx.try_send(snapshot).is_ok() { - self.last_checkpoint_objects = self.checkpoint_objects; - self.last_checkpoint_at = Instant::now(); - } + permit.send(snapshot); + self.last_checkpoint_objects = self.checkpoint_objects; + self.last_checkpoint_at = Instant::now(); } fn carry_forward_old_children(&mut self, parent_hash: &DataUsageHash, entry: &mut DataUsageEntry) { diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index c1f488dfc..a940a2eb2 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -409,6 +409,16 @@ async fn periodic_checkpoint_emits_at_object_threshold_without_wall_clock_wait() scanner.maybe_send_checkpoint(); + scanner.checkpoint_objects = SCANNER_CHECKPOINT_OBJECT_INTERVAL * 2; + scanner.last_checkpoint_at = Instant::now() + .checked_sub(SCANNER_CHECKPOINT_MIN_INTERVAL) + .expect("test instant subtraction"); + scanner.maybe_send_checkpoint(); + assert_eq!( + scanner.last_checkpoint_objects, SCANNER_CHECKPOINT_OBJECT_INTERVAL, + "a full checkpoint queue must reject before cloning or advancing progress" + ); + let checkpoint = checkpoint_rx.try_recv().expect("object threshold emits a bounded checkpoint"); assert_eq!(checkpoint.info.name, "bucket"); assert!(!checkpoint.info.snapshot_complete);