diff --git a/crates/common/src/metrics.rs b/crates/common/src/metrics.rs index cb524b371..4c409279f 100644 --- a/crates/common/src/metrics.rs +++ b/crates/common/src/metrics.rs @@ -915,11 +915,13 @@ const SCAN_CYCLE_RESULT_SUCCESS: u8 = 1; const SCAN_CYCLE_RESULT_ERROR: u8 = 2; const SCAN_CYCLE_RESULT_PARTIAL: u8 = 3; const SCAN_CYCLE_RESULT_SUPERSEDED: u8 = 4; +const SCAN_CYCLE_RESULT_DEFERRED: u8 = 5; const SCAN_CYCLE_RESULT_UNKNOWN_LABEL: &str = "unknown"; const SCAN_CYCLE_RESULT_SUCCESS_LABEL: &str = "success"; const SCAN_CYCLE_RESULT_ERROR_LABEL: &str = "error"; const SCAN_CYCLE_RESULT_PARTIAL_LABEL: &str = "partial"; const SCAN_CYCLE_RESULT_SUPERSEDED_LABEL: &str = "superseded"; +const SCAN_CYCLE_RESULT_DEFERRED_LABEL: &str = "deferred"; #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] pub enum ScanCyclePartialReason { @@ -1424,6 +1426,7 @@ fn scan_cycle_result_label(result: u8) -> &'static str { SCAN_CYCLE_RESULT_ERROR => SCAN_CYCLE_RESULT_ERROR_LABEL, SCAN_CYCLE_RESULT_PARTIAL => SCAN_CYCLE_RESULT_PARTIAL_LABEL, SCAN_CYCLE_RESULT_SUPERSEDED => SCAN_CYCLE_RESULT_SUPERSEDED_LABEL, + SCAN_CYCLE_RESULT_DEFERRED => SCAN_CYCLE_RESULT_DEFERRED_LABEL, _ => SCAN_CYCLE_RESULT_UNKNOWN_LABEL, } } @@ -1752,6 +1755,11 @@ pub fn emit_scan_cycle_superseded(duration: Duration) { metrics::counter!(OTEL_SCANNER_CYCLES, "result" => SCAN_CYCLE_RESULT_SUPERSEDED_LABEL).increment(1); } +pub fn emit_scan_cycle_deferred(duration: Duration) { + global_metrics().record_scan_cycle_deferred(duration); + metrics::counter!(OTEL_SCANNER_CYCLES, "result" => SCAN_CYCLE_RESULT_DEFERRED_LABEL).increment(1); +} + pub fn emit_scan_bucket_drive_complete(success: bool, bucket: &str, disk: &str, duration: Duration) { let result = if success { "success" } else { "error" }; global_metrics().record_scanner_bucket_drive_result(bucket, disk, result); @@ -2549,6 +2557,17 @@ impl Metrics { .store(duration_millis_saturated(duration), Ordering::Relaxed); } + pub fn record_scan_cycle_deferred(&self, duration: Duration) { + self.record_scanner_cycle_end_time(); + self.last_scan_cycle_result + .store(SCAN_CYCLE_RESULT_DEFERRED, Ordering::Relaxed); + self.last_scan_cycle_partial_reason + .store(ScanCyclePartialReason::Unknown as u8, Ordering::Relaxed); + self.last_scan_cycle_partial_source.store(0, Ordering::Relaxed); + self.last_scan_cycle_duration_millis + .store(duration_millis_saturated(duration), Ordering::Relaxed); + } + pub fn record_scan_cycle_partial(&self, duration: Duration, reason: ScanCyclePartialReason) { self.record_scan_cycle_partial_with_source(duration, reason, None); } @@ -4264,6 +4283,21 @@ mod tests { assert_eq!(report.partial_cycles, 0); } + #[tokio::test] + async fn report_tracks_deferred_cycle_without_failed_increment() { + let metrics = Metrics::new(); + metrics.record_scan_cycle_deferred(Duration::from_millis(250)); + + let report = metrics.report().await; + + assert_eq!(report.last_cycle_result, SCAN_CYCLE_RESULT_DEFERRED_LABEL); + assert_eq!(report.last_cycle_result_code, u64::from(SCAN_CYCLE_RESULT_DEFERRED)); + assert_eq!(report.last_cycle_duration_seconds, 0.25); + assert_eq!(report.failed_cycles, 0); + assert_eq!(report.superseded_cycles, 0); + assert_eq!(report.partial_cycles, 0); + } + #[tokio::test] async fn report_tracks_successful_scan_cycle_without_failed_increment() { let metrics = Metrics::new(); diff --git a/crates/obs/src/metrics/collectors/scanner.rs b/crates/obs/src/metrics/collectors/scanner.rs index 1689b722e..7ce05c23d 100644 --- a/crates/obs/src/metrics/collectors/scanner.rs +++ b/crates/obs/src/metrics/collectors/scanner.rs @@ -113,7 +113,7 @@ pub struct ScannerStats { pub current_cycle_usage_saves: u64, /// Current scanner mode: 0 unknown or idle, 1 normal, 2 deep bitrot scan pub current_scan_mode: u64, - /// Last scanner cycle result: 0 unknown, 1 success, 2 error, 3 partial, 4 superseded + /// Last scanner cycle result: 0 unknown, 1 success, 2 error, 3 partial, 4 superseded, 5 deferred pub last_cycle_result: u64, /// Last scanner partial cycle reason: 0 unknown, 1 runtime, 2 objects, 3 directories pub last_cycle_partial_reason: u64, diff --git a/crates/obs/src/metrics/schema/scanner.rs b/crates/obs/src/metrics/schema/scanner.rs index 965810ff7..a9027bf07 100644 --- a/crates/obs/src/metrics/schema/scanner.rs +++ b/crates/obs/src/metrics/schema/scanner.rs @@ -460,7 +460,7 @@ pub static SCANNER_CURRENT_SCAN_MODE_MD: LazyLock = LazyLock:: pub static SCANNER_LAST_CYCLE_RESULT_MD: LazyLock = LazyLock::new(|| { new_gauge_md( MetricName::ScannerLastCycleResult, - "Last scanner cycle result: 0 unknown, 1 success, 2 error, 3 partial, 4 superseded.", + "Last scanner cycle result: 0 unknown, 1 success, 2 error, 3 partial, 4 superseded, 5 deferred.", &[], subsystems::SCANNER, ) diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index 59cd21e7e..11023eaac 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -30,8 +30,8 @@ use crate::runtime_config::{ use crate::scanner_budget::{ScannerCycleBudget, ScannerCycleBudgetConfig, ScannerCycleBudgetReason}; use crate::scanner_folder::{data_usage_update_dir_cycles, heal_object_select_prob}; use crate::scanner_io::{ - ScannerCycleStatus, ScannerIOCycle, dirty_usage_bucket_notified, dirty_usage_buckets_pending, dirty_usage_generation, - scanner_dirty_usage_state, scanner_maintenance_changed, scanner_maintenance_generation, + ScannerCycleDeferReason, ScannerCycleStatus, ScannerIOCycle, dirty_usage_bucket_notified, dirty_usage_buckets_pending, + dirty_usage_generation, scanner_dirty_usage_state, scanner_maintenance_changed, scanner_maintenance_generation, }; use crate::sleeper::{SCANNER_SLEEPER, set_scanner_default_speed}; use crate::{DataUsageInfo, ScannerActivityGuard, ScannerError, ScannerRuntimeGuard}; @@ -41,7 +41,8 @@ use chrono::{DateTime, Utc}; use rustfs_common::heal_channel::HealScanMode; use rustfs_common::metrics::{ CurrentCycle, Metric, Metrics, ScanCyclePartialReason, ScanCycleWorkSnapshot, ScannerUsageSaveResult, ScannerWorkSource, - emit_scan_cycle_complete, emit_scan_cycle_partial_with_source, emit_scan_cycle_superseded, global_metrics, + emit_scan_cycle_complete, emit_scan_cycle_deferred, emit_scan_cycle_partial_with_source, emit_scan_cycle_superseded, + global_metrics, }; use rustfs_config::ScannerSpeed; #[cfg(test)] @@ -84,7 +85,7 @@ const METRIC_SCANNER_LEADER_LOCK_TOTAL: &str = "rustfs_scanner_leader_lock_total const CLEAN_IDLE_MAX_INTERVAL: Duration = Duration::from_secs(24 * 60 * 60); const MAX_SCANNER_SCHEDULE_DELAY: Duration = Duration::from_secs(365 * 24 * 60 * 60); const CLEAN_IDLE_BACKOFF_FACTOR: u32 = 2; -/// First-retry delay after a usage snapshot is superseded by concurrent writes. +/// First-retry delay after a scanner cycle cannot publish authoritative usage. /// /// A superseded cycle is the *expected* outcome of the dirty-usage fast path: /// a write burst marks buckets dirty, the scanner wakes within milliseconds, @@ -94,14 +95,15 @@ const CLEAN_IDLE_BACKOFF_FACTOR: u32 = 2; /// otherwise idle instance whose clean-idle backoff had doubled a 60 s /// interval), which defeats the fast path it is meant to protect. /// -/// The exponential growth in [`ScannerSupersededBackoff::retry_interval`] is +/// The exponential growth in [`ScannerRetryBackoff::retry_interval`] is /// what protects against a persistently hot bucket driving an unbroken /// full-scan loop, so it can start small: 5 s, 10 s, 20 s … capped by -/// [`SUPERSEDED_RETRY_MAX_INTERVAL`]. A one-off race recovers in seconds; a +/// [`SCANNER_RETRY_MAX_INTERVAL`]. A one-off race recovers in seconds; a /// genuinely hot bucket still reaches minute-scale backoff within a handful of -/// cycles. -const SUPERSEDED_RETRY_BASE_INTERVAL: Duration = Duration::from_secs(5); -const SUPERSEDED_RETRY_MAX_INTERVAL: Duration = Duration::from_secs(30 * 60); +/// cycles. Preflight deferrals use the same bounded schedule so a temporarily +/// unavailable peer cannot drive a tight retry loop. +const SCANNER_RETRY_BASE_INTERVAL: Duration = Duration::from_secs(5); +const SCANNER_RETRY_MAX_INTERVAL: Duration = Duration::from_secs(30 * 60); const SCANNER_LEADER_LOCK_POLL_INTERVAL: Duration = Duration::from_secs(1); #[cfg(not(test))] const SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(30); @@ -338,6 +340,7 @@ pub(crate) enum ScannerCycleOutcome { CompletedWithPendingMaintenance, Partial, Superseded, + Deferred(ScannerCycleDeferReason), Failed, } @@ -382,22 +385,16 @@ struct ScannerCleanIdleBackoff { } #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] -struct ScannerSupersededBackoff { +struct ScannerRetryBackoff { consecutive_cycles: u32, } -impl ScannerSupersededBackoff { - fn record_cycle(&mut self, outcome: ScannerCycleOutcome) { - match outcome { - ScannerCycleOutcome::Superseded => { - self.consecutive_cycles = self.consecutive_cycles.saturating_add(1); - } - ScannerCycleOutcome::Completed - | ScannerCycleOutcome::CompletedWithPendingMaintenance - | ScannerCycleOutcome::Partial - | ScannerCycleOutcome::Failed => { - self.consecutive_cycles = 0; - } +impl ScannerRetryBackoff { + fn record_retryable_cycle(&mut self, retryable: bool) { + if retryable { + self.consecutive_cycles = self.consecutive_cycles.saturating_add(1); + } else { + self.consecutive_cycles = 0; } } @@ -406,8 +403,8 @@ impl ScannerSupersededBackoff { let multiplier = 1u32.checked_shl(exponent).unwrap_or(u32::MAX); let base_interval = configured_interval .max(Duration::from_secs(1)) - .min(SUPERSEDED_RETRY_BASE_INTERVAL); - let cap = SUPERSEDED_RETRY_MAX_INTERVAL.max(configured_interval.max(Duration::from_secs(1))); + .min(SCANNER_RETRY_BASE_INTERVAL); + let cap = SCANNER_RETRY_MAX_INTERVAL.max(configured_interval.max(Duration::from_secs(1))); Some(base_interval.saturating_mul(multiplier).min(cap)) } } @@ -3008,6 +3005,21 @@ async fn run_data_scanner_cycle( ScannerCycleOutcome::Failed }; } + ScannerCycleOutcome::Deferred(reason) => { + info!( + target: "rustfs::scanner", + event = EVENT_SCANNER_CYCLE_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + cycle = cycle_info.current, + reason = reason.as_str(), + state = "deferred", + "Scanner cycle deferred before usage scanning began" + ); + emit_scan_cycle_deferred(cycle_start.elapsed()); + mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; + return ScannerCycleOutcome::Deferred(reason); + } ScannerCycleOutcome::Superseded => { info!( target: "rustfs::scanner", @@ -3201,7 +3213,8 @@ async fn run_data_scanner_with_maintenance_state( let mut dirty_usage_generation_seen = dirty_usage_generation(); let mut runtime_config_generation_seen = scanner_runtime_config_generation(); let mut clean_idle_backoff = ScannerCleanIdleBackoff::default(); - let mut superseded_backoff = ScannerSupersededBackoff::default(); + let mut superseded_backoff = ScannerRetryBackoff::default(); + let mut deferred_backoff = ScannerRetryBackoff::default(); let initial_runtime_config = resolve_scanner_runtime_config(); if clean_idle_topology_supported && scanner_clean_idle_backoff_configured(&initial_runtime_config) @@ -3331,7 +3344,8 @@ async fn run_data_scanner_with_maintenance_state( ) .await .unwrap_or(ScannerCycleOutcome::Failed); - superseded_backoff.record_cycle(initial_outcome); + superseded_backoff.record_retryable_cycle(initial_outcome == ScannerCycleOutcome::Superseded); + deferred_backoff.record_retryable_cycle(matches!(initial_outcome, ScannerCycleOutcome::Deferred(_))); dirty_usage_generation_seen = dirty_generation_before_cycle; if guard.is_lock_lost() { record_scanner_leader_lock_lost("Scanner leader lock lost during the initial cycle").await; @@ -3416,7 +3430,9 @@ async fn run_data_scanner_with_maintenance_state( let mut wait_plan = scanner_cycle_wait_plan(&runtime_config, clean_idle_backoff, backoff_enabled, randomized_cycle_delay_for); let superseded_retry_interval = superseded_backoff.retry_interval(runtime_config.cycle_interval); - if let Some(retry_interval) = superseded_retry_interval { + let deferred_retry_interval = deferred_backoff.retry_interval(runtime_config.cycle_interval); + let convergence_retry_interval = superseded_retry_interval.or(deferred_retry_interval); + if let Some(retry_interval) = convergence_retry_interval { wait_plan.effective_interval = retry_interval; wait_plan.delay = randomized_cycle_delay_for(retry_interval).min(retry_interval); } @@ -3443,6 +3459,8 @@ async fn run_data_scanner_with_maintenance_state( clean_idle_backoff_enabled = backoff_enabled, superseded_retry_backoff_enabled = superseded_retry_interval.is_some(), superseded_cycles = superseded_backoff.consecutive_cycles, + deferred_retry_backoff_enabled = deferred_retry_interval.is_some(), + deferred_cycles = deferred_backoff.consecutive_cycles, lifecycle_active = maintenance_features.lifecycle, replication_active = maintenance_features.replication, feature_inspection_failed = maintenance_features.inspection_failed, @@ -3457,13 +3475,12 @@ async fn run_data_scanner_with_maintenance_state( activity_poll_interval, &mut scanner_activity_seen, ScannerCycleObservedGenerations { - // A superseded cycle already observed concurrent writes. Hold - // further dirty notifications until the bounded retry timer so - // a hot bucket cannot drive an unbroken full-scan loop. - dirty_usage: superseded_retry_interval.is_none().then_some(dirty_usage_generation_seen), + // A non-converged cycle holds further activity notifications + // until its bounded retry timer to avoid an unbroken scan loop. + dirty_usage: convergence_retry_interval.is_none().then_some(dirty_usage_generation_seen), runtime_config: runtime_config_generation_seen, maintenance: maintenance_generation_before_wait, - defer_cluster_activity: superseded_retry_interval.is_some(), + defer_cluster_activity: convergence_retry_interval.is_some(), }, || guard.is_lock_lost(), || probe_scanner_activity(storeapi.as_ref(), distributed), @@ -3540,7 +3557,8 @@ async fn run_data_scanner_with_maintenance_state( ) .await .unwrap_or(ScannerCycleOutcome::Failed); - superseded_backoff.record_cycle(outcome); + superseded_backoff.record_retryable_cycle(outcome == ScannerCycleOutcome::Superseded); + deferred_backoff.record_retryable_cycle(matches!(outcome, ScannerCycleOutcome::Deferred(_))); dirty_usage_generation_seen = dirty_generation_before_cycle; if guard.is_lock_lost() { record_scanner_leader_lock_lost("Scanner leader lock lost during a scanner cycle").await; @@ -3710,6 +3728,12 @@ fn scanner_cycle_completion_outcome( ) -> ScannerCycleOutcome { match (scan_status, usage_persist_outcome) { (_, DataUsagePersistOutcome::Failed) => ScannerCycleOutcome::Failed, + (ScannerCycleStatus::Deferred(reason), DataUsagePersistOutcome::NoUpdate) + if !has_dirty_usage && !has_failed_dirty_usage => + { + ScannerCycleOutcome::Deferred(reason) + } + (ScannerCycleStatus::Deferred(_), _) => ScannerCycleOutcome::Failed, (ScannerCycleStatus::Superseded, _) if !has_failed_dirty_usage => ScannerCycleOutcome::Superseded, (ScannerCycleStatus::Superseded, _) => ScannerCycleOutcome::Failed, ( @@ -6553,6 +6577,46 @@ mod tests { #[test] fn test_scanner_cycle_completion_prioritizes_persist_failure() { + assert_eq!( + scanner_cycle_completion_outcome( + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable), + DataUsagePersistOutcome::NoUpdate, + false, + false, + ), + ScannerCycleOutcome::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable) + ); + assert_eq!( + scanner_cycle_completion_outcome( + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::DataMovement), + DataUsagePersistOutcome::Saved, + false, + false, + ), + ScannerCycleOutcome::Failed + ); + assert_eq!( + scanner_cycle_completion_outcome( + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::DataMovement), + DataUsagePersistOutcome::NoUpdate, + true, + false, + ), + ScannerCycleOutcome::Failed + ); + assert_eq!( + scanner_cycle_completion_outcome( + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::DataMovement), + DataUsagePersistOutcome::Failed, + false, + false, + ), + ScannerCycleOutcome::Failed + ); + assert_eq!( + scanner_cycle_completion_outcome(ScannerCycleStatus::Incomplete, DataUsagePersistOutcome::NoUpdate, false, false), + ScannerCycleOutcome::Failed + ); assert_eq!( scanner_cycle_completion_outcome(ScannerCycleStatus::Incomplete, DataUsagePersistOutcome::Failed, true, true), ScannerCycleOutcome::Failed @@ -6940,47 +7004,47 @@ mod tests { #[test] fn superseded_retry_backoff_grows_caps_and_resets_after_convergence() { - let mut backoff = ScannerSupersededBackoff::default(); + let mut backoff = ScannerRetryBackoff::default(); assert_eq!(backoff.retry_interval(Duration::from_secs(24 * 60 * 60)), None); for expected in [5, 10, 20, 40, 80, 160, 320] { - backoff.record_cycle(ScannerCycleOutcome::Superseded); + backoff.record_retryable_cycle(true); assert_eq!( backoff.retry_interval(Duration::from_secs(24 * 60 * 60)), Some(Duration::from_secs(expected)) ); } for _ in 0..20 { - backoff.record_cycle(ScannerCycleOutcome::Superseded); + backoff.record_retryable_cycle(true); } assert_eq!( backoff.retry_interval(Duration::from_secs(24 * 60 * 60)), Some(Duration::from_secs(24 * 60 * 60)) ); - backoff.record_cycle(ScannerCycleOutcome::Completed); + backoff.record_retryable_cycle(false); assert_eq!(backoff.retry_interval(Duration::from_secs(24 * 60 * 60)), None); } #[test] fn superseded_retry_backoff_respects_a_faster_configured_cycle() { - let mut backoff = ScannerSupersededBackoff::default(); - backoff.record_cycle(ScannerCycleOutcome::Superseded); + let mut backoff = ScannerRetryBackoff::default(); + backoff.record_retryable_cycle(true); // A configured cycle shorter than the base still wins: retrying sooner // than the operator's own cadence buys nothing. assert_eq!(backoff.retry_interval(Duration::from_secs(3)), Some(Duration::from_secs(3))); - backoff.record_cycle(ScannerCycleOutcome::Superseded); + backoff.record_retryable_cycle(true); assert_eq!(backoff.retry_interval(Duration::from_secs(3)), Some(Duration::from_secs(6))); } #[test] fn superseded_retry_backoff_grows_from_the_default_cycle() { - let mut backoff = ScannerSupersededBackoff::default(); + let mut backoff = ScannerRetryBackoff::default(); // The first race after a write burst retries in seconds, not a whole // cycle, while repeated supersedes still climb toward the cap. for expected in [5, 10, 20, 40] { - backoff.record_cycle(ScannerCycleOutcome::Superseded); + backoff.record_retryable_cycle(true); assert_eq!(backoff.retry_interval(Duration::from_secs(60)), Some(Duration::from_secs(expected))); } } diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 78987fa13..019a1fe64 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -2194,11 +2194,45 @@ pub(crate) async fn scanner_set_disk_inventory(set: &SetDisks) -> Vec> disks } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum ScannerCycleDeferReason { + ActivityBaselineUnavailable, + DataMovement, +} + +impl ScannerCycleDeferReason { + pub(crate) fn as_str(self) -> &'static str { + match self { + Self::ActivityBaselineUnavailable => "activity_baseline_unavailable", + Self::DataMovement => "data_movement", + } + } +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub(crate) enum ScannerCycleStatus { Complete, Incomplete, Superseded, + Deferred(ScannerCycleDeferReason), +} + +enum ScannerActivityPreflight { + Ready(crate::scanner::ScannerActivitySnapshot), + ActivityBaselineUnavailable(String), + DataMovement, +} + +fn scanner_activity_preflight( + activity: std::result::Result, +) -> ScannerActivityPreflight { + match activity { + Err(error) => ScannerActivityPreflight::ActivityBaselineUnavailable(error), + Ok(snapshot) if !crate::scanner::scanner_activity_allows_usage_publication(&snapshot) => { + ScannerActivityPreflight::DataMovement + } + Ok(snapshot) => ScannerActivityPreflight::Ready(snapshot), + } } #[derive(Debug)] @@ -2307,9 +2341,9 @@ impl ScannerIOCycle for ECStore { let child_token = ctx.child_token(); let distributed = self.setup_is_dist_erasure().await; - let activity_before = match crate::scanner::probe_scanner_activity(self, distributed).await { - Ok(snapshot) => snapshot, - Err(err) => { + let activity_before = match scanner_activity_preflight(crate::scanner::probe_scanner_activity(self, distributed).await) { + ScannerActivityPreflight::Ready(snapshot) => snapshot, + ScannerActivityPreflight::ActivityBaselineUnavailable(err) => { warn!( target: "rustfs::scanner::io", event = EVENT_SCANNER_SET_STATE, @@ -2319,20 +2353,26 @@ impl ScannerIOCycle for ECStore { error = %err, "Scanner cycle skipped because cluster activity could not be baselined" ); - return Ok(ScannerCycleResult::new(ScannerCycleStatus::Incomplete, None)); + return Ok(ScannerCycleResult::new( + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable), + None, + )); + } + ScannerActivityPreflight::DataMovement => { + debug!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_SET_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + state = "cycle_data_movement_active", + "Scanner cycle deferred while rebalance or decommission data movement is active" + ); + return Ok(ScannerCycleResult::new( + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::DataMovement), + None, + )); } }; - if !crate::scanner::scanner_activity_allows_usage_publication(&activity_before) { - debug!( - target: "rustfs::scanner::io", - event = EVENT_SCANNER_SET_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_IO, - state = "cycle_data_movement_active", - "Scanner cycle deferred while rebalance or decommission data movement is active" - ); - return Ok(ScannerCycleResult::new(ScannerCycleStatus::Incomplete, None)); - } let dirty_generation_before_bucket_list = dirty_usage_generation(); let bucket_listing = self.list_bucket_for_scanner(&BucketOptions::default()).await?; let mut bucket_plan_complete = bucket_listing.topology_complete; @@ -3982,6 +4022,7 @@ mod tests { use super::*; use crate::scanner_budget::ScannerCycleBudgetConfig; use crate::scanner_folder::ScannerItem; + use crate::storage_api::owner::{EcstoreRebalStatus, EcstoreRebalanceInfo, EcstoreRebalanceMeta, EcstoreRebalanceStats}; use crate::storage_api::scan::{BucketOperations as _, MakeBucketOptions, ObjectIO as _}; use crate::{ DiskOption, ECStore, Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints, ScannerObjectOptions, @@ -4004,6 +4045,20 @@ mod tests { } } + #[test] + fn scanner_activity_preflight_defers_a_temporarily_offline_peer() { + let preflight = scanner_activity_preflight(Err("peer rustfs-node3:9000 is temporarily offline".to_string())); + + match preflight { + ScannerActivityPreflight::ActivityBaselineUnavailable(error) => { + assert_eq!(error, "peer rustfs-node3:9000 is temporarily offline"); + } + ScannerActivityPreflight::Ready(_) | ScannerActivityPreflight::DataMovement => { + panic!("an unavailable activity baseline must defer the scanner cycle"); + } + } + } + async fn setup_two_pool_scanner_store() -> (tempfile::TempDir, Arc) { init_ecstore_config_for_scanner_tests(); let temp_dir = tempfile::tempdir().expect("multi-pool scanner test directory should be created"); @@ -4096,6 +4151,42 @@ mod tests { assert!(!second.is_lock_lost()); } + #[tokio::test] + #[serial] + async fn scanner_cycle_is_deferred_while_rebalance_is_active() { + let (_temp_dir, store) = setup_two_pool_scanner_store().await; + let mut pool_stats = vec![EcstoreRebalanceStats::default(); store.pools.len()]; + pool_stats[0] = EcstoreRebalanceStats { + participating: true, + info: EcstoreRebalanceInfo { + start_time: Some(OffsetDateTime::now_utc()), + status: EcstoreRebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }; + *store.rebalance_meta.write().await = Some(EcstoreRebalanceMeta { + id: Uuid::new_v4().to_string(), + pool_stats, + ..Default::default() + }); + assert!(store.scanner_data_movement_active().await); + + let ctx = CancellationToken::new(); + let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default()); + let (updates, mut receiver) = mpsc::channel(1); + let result = tokio::time::timeout( + Duration::from_secs(30), + ScannerIOCycle::nsscanner_with_status(store.as_ref(), ctx, budget, updates, 1, 1, HealScanMode::Normal), + ) + .await + .expect("rebalance-deferred scanner cycle should finish") + .expect("rebalance-deferred scanner cycle should succeed"); + + assert_eq!(result.status, ScannerCycleStatus::Deferred(ScannerCycleDeferReason::DataMovement)); + assert!(receiver.recv().await.is_none(), "rebalance-deferred cycle must not publish usage"); + } + #[tokio::test] async fn data_usage_publish_fails_when_receiver_is_closed() { let (updates, receiver) = mpsc::channel(1); diff --git a/crates/scanner/src/storage_api.rs b/crates/scanner/src/storage_api.rs index 455b17085..634f20e45 100644 --- a/crates/scanner/src/storage_api.rs +++ b/crates/scanner/src/storage_api.rs @@ -83,6 +83,11 @@ pub(crate) use rustfs_ecstore::api::layout::{ EndpointServerPools as EcstoreEndpointServerPools, Endpoints as EcstoreEndpoints, PoolEndpoints as EcstorePoolEndpoints, }; #[cfg(test)] +pub(crate) use rustfs_ecstore::api::rebalance::{ + RebalStatus as EcstoreRebalStatus, RebalanceInfo as EcstoreRebalanceInfo, RebalanceMeta as EcstoreRebalanceMeta, + RebalanceStats as EcstoreRebalanceStats, +}; +#[cfg(test)] pub(crate) use rustfs_ecstore::api::runtime::InstanceContext as EcstoreInstanceContext; pub(crate) use rustfs_ecstore::api::runtime::{ expiry_state_handle as ecstore_expiry_state_handle, global_tier_config_mgr as ecstore_get_global_tier_config_mgr, @@ -122,8 +127,9 @@ pub(crate) mod owner { #[cfg(test)] pub(crate) use super::{ EcstoreDiskOption, EcstoreDiskStore, EcstoreEndpoint, EcstoreEndpointServerPools, EcstoreEndpoints, - EcstoreInstanceContext, EcstorePoolEndpoints, ecstore_config_init, ecstore_init_bucket_metadata_sys, - ecstore_init_local_disks_with_instance_ctx, ecstore_new_disk, + EcstoreInstanceContext, EcstorePoolEndpoints, EcstoreRebalStatus, EcstoreRebalanceInfo, EcstoreRebalanceMeta, + EcstoreRebalanceStats, ecstore_config_init, ecstore_init_bucket_metadata_sys, ecstore_init_local_disks_with_instance_ctx, + ecstore_new_disk, }; }