mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 04:21:35 +00:00
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>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -1454,13 +1454,32 @@ fn bitrot_scan_cycle() -> Option<Duration> {
|
||||
resolve_scanner_runtime_config().bitrot_cycle
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
struct ScannerBitrotPolicy {
|
||||
cycle: Option<Duration>,
|
||||
deep_window_cycles: Option<u64>,
|
||||
}
|
||||
|
||||
impl ScannerBitrotPolicy {
|
||||
fn new(cycle: Option<Duration>, 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<DateTime<Utc>>,
|
||||
bitrot_cycle: Option<Duration>,
|
||||
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<Utc>,
|
||||
bitrot_cycle: Option<Duration>,
|
||||
policy: ScannerBitrotPolicy,
|
||||
) -> Option<BackgroundHealInfo> {
|
||||
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<Utc>,
|
||||
bitrot_cycle: Option<Duration>,
|
||||
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;
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user