diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index e3d461805..dd6d00bf6 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -14,7 +14,6 @@ use std::collections::BTreeMap; use std::future::Future; -#[cfg(test)] use std::sync::Mutex as StdMutex; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, LazyLock, RwLock}; @@ -1585,6 +1584,11 @@ async fn mark_scan_cycle_idle(cycle_info: &mut CurrentCycle, cycle_metrics_guard cycle_metrics_guard.finish(cycle_info.clone()).await; } +struct ScannerCycleScheduling { + requires_full_scan: bool, + service_cohort: Option>>, +} + #[cfg(test)] async fn run_data_scanner_cycle( ctx: &CancellationToken, @@ -1597,7 +1601,19 @@ where S: ScannerStorage, { let cycle_budget = ScannerCycleBudget::new(ctx, scanner_cycle_budget_config()); - run_data_scanner_cycle_with_budget(ctx, storeapi, cycle_info, cycle_revision, leader_epoch, cycle_budget, true).await + run_data_scanner_cycle_with_budget( + ctx, + storeapi, + cycle_info, + cycle_revision, + leader_epoch, + cycle_budget, + ScannerCycleScheduling { + requires_full_scan: true, + service_cohort: None, + }, + ) + .await } #[instrument(skip_all)] @@ -1609,7 +1625,7 @@ async fn run_data_scanner_cycle_with_budget( cycle_revision: &mut DataUsageCacheRevision, leader_epoch: u64, cycle_budget: Arc, - requires_full_scan: bool, + scheduling: ScannerCycleScheduling, ) -> ScannerCycleOutcome where S: ScannerStorage, @@ -1748,7 +1764,8 @@ where scan_mode, scan_scope: crate::scanner_io::ScannerBucketScanScope::default(), persisted_usage_baseline: usage_persist_baseline.data.clone(), - requires_full_scan, + requires_full_scan: scheduling.requires_full_scan, + service_cohort: scheduling.service_cohort, #[cfg(test)] resolved_scope_observer: None, }, @@ -2610,6 +2627,7 @@ where let mut clean_idle_backoff = ScannerCleanIdleBackoff::default(); let mut superseded_backoff = ScannerRetryBackoff::default(); let mut deferred_backoff = ScannerRetryBackoff::default(); + let service_cohort = Arc::new(StdMutex::new(crate::scanner_io::ScannerServiceCohort::default())); let initial_runtime_config = resolve_scanner_runtime_config(); if clean_idle_topology_supported && maintenance_generation_seen.is_none() { let Some((features, generation)) = detect_stable_scanner_maintenance_features(&ctx, &storeapi).await else { @@ -2819,7 +2837,10 @@ where &mut cycle_revision, leader_epoch, cycle_budget.clone(), - true, + ScannerCycleScheduling { + requires_full_scan: true, + service_cohort: Some(service_cohort.clone()), + }, ), guard.lock_lost_notified(), ) @@ -3110,11 +3131,14 @@ where &mut cycle_revision, leader_epoch, cycle_budget.clone(), - maintenance_features.requires_full_scan( - maintenance_generation_seen, - scanner_maintenance_generation(), - wake_reason, - ), + ScannerCycleScheduling { + requires_full_scan: maintenance_features.requires_full_scan( + maintenance_generation_seen, + scanner_maintenance_generation(), + wake_reason, + ), + service_cohort: Some(service_cohort.clone()), + }, ), guard.lock_lost_notified(), ) diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index e89253d05..04d548404 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -1221,7 +1221,18 @@ async fn coordinator_walks_during_pending_put_without_persisting_or_acknowledgin let mut revision = DataUsageCacheRevision::Missing; let outcome = tokio::time::timeout( Duration::from_secs(30), - run_data_scanner_cycle_with_budget(&ctx, &store, &mut cycle_info, &mut revision, 1, Arc::clone(&budget), true), + run_data_scanner_cycle_with_budget( + &ctx, + &store, + &mut cycle_info, + &mut revision, + 1, + Arc::clone(&budget), + ScannerCycleScheduling { + requires_full_scan: true, + service_cohort: None, + }, + ), ) .await .expect("the coordinator must finish its namespace walk while a PUT is pending"); @@ -1265,7 +1276,18 @@ async fn coordinator_walks_during_pending_put_without_persisting_or_acknowledgin let retry_budget = ScannerCycleBudget::new_with_progress_tracking(&ctx, ScannerCycleBudgetConfig::default()); let outcome = tokio::time::timeout( Duration::from_secs(30), - run_data_scanner_cycle_with_budget(&ctx, &store, &mut cycle_info, &mut revision, 1, Arc::clone(&retry_budget), true), + run_data_scanner_cycle_with_budget( + &ctx, + &store, + &mut cycle_info, + &mut revision, + 1, + Arc::clone(&retry_budget), + ScannerCycleScheduling { + requires_full_scan: true, + service_cohort: None, + }, + ), ) .await .expect("the same cycle must converge after the pending PUT drains"); @@ -8516,6 +8538,56 @@ async fn test_wait_for_next_scanner_cycle_wakes_for_dirty_usage() { crate::scanner_io::clear_dirty_usage_buckets_for_tests(); } +#[tokio::test(start_paused = true)] +#[serial] +async fn service_cohort_aging_preserves_explicit_cycle_wait() { + crate::scanner_io::clear_dirty_usage_buckets_for_tests(); + let config = ScannerRuntimeConfig { + cycle_interval: Duration::from_secs(3600), + cycle_interval_source: ScannerRuntimeConfigSource::Env, + ..Default::default() + }; + let observed = ScannerCycleObservedGenerations::for_wait( + &config, + None, + crate::scanner_io::dirty_usage_generation(), + crate::runtime_config::scanner_runtime_config_generation(), + crate::scanner_io::scanner_maintenance_generation(), + ); + assert_eq!(observed.dirty_usage, None); + let inventory = HashMap::from([( + crate::data_usage_define::DataUsageCacheSource::new(0, 0), + vec![crate::storage_api::scanner_io::BucketInfo { + name: "waiting-bootstrap".to_string(), + ..Default::default() + }], + )]); + let mut cohort = crate::scanner_io::ScannerServiceCohort::default(); + cohort.refresh(&inventory); + let ctx = CancellationToken::new(); + let mut wait = Box::pin(wait_for_next_scanner_cycle( + &ctx, + config.cycle_interval, + observed.dirty_usage, + observed.runtime_config, + observed.maintenance, + || false, + )); + assert!(matches!(futures::poll!(&mut wait), Poll::Pending)); + for _ in 0..59 { + tokio::time::advance(Duration::from_secs(60)).await; + cohort.refresh(&inventory); + crate::scanner_io::record_dirty_usage_bucket("hot"); + assert!( + matches!(futures::poll!(&mut wait), Poll::Pending), + "aging/dirty must not shorten the explicit hour" + ); + } + tokio::time::advance(Duration::from_secs(60)).await; + assert_eq!(wait.await, ScannerCycleWakeReason::Timer); + crate::scanner_io::clear_dirty_usage_buckets_for_tests(); +} + #[tokio::test] #[serial] async fn test_wait_for_next_scanner_cycle_sees_unattempted_dirty_usage() { diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 09e959a3b..b0bac1740 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -286,6 +286,7 @@ pub struct ScannerBucketScanPlan { /// Includes mutation generations even when the set planner uses a structural digest. bucket_coverage_digest: DataUsageScanPlanDigest, requires_full_scan: bool, + service_cohort: Option>>, // Cache work must invalidate on namespace completion even when its scoped baseline remains reusable. execution_digest: DataUsageScanPlanDigest, leader_epoch: u64, @@ -973,6 +974,7 @@ impl ScannerCycleResult { mod cache; mod dirty_usage; mod guards; +pub(crate) use guards::ScannerServiceCohort; mod io_cache; mod io_cycle; #[cfg(test)] diff --git a/crates/scanner/src/scanner_io/guards.rs b/crates/scanner/src/scanner_io/guards.rs index 628a59755..580fccb81 100644 --- a/crates/scanner/src/scanner_io/guards.rs +++ b/crates/scanner/src/scanner_io/guards.rs @@ -14,6 +14,312 @@ /// scan concurrency accounting: gauge recorders, RAII guards, and worker limits. use super::*; +const SCANNER_SERVICE_COHORT_MAX_MEMBERS: usize = 4096; +const SCANNER_SERVICE_COHORT_MAX_NAME_BYTES: usize = 128 * 1024; + +static SERVICE_COHORT_METRICS_OWNER: StdMutex> = StdMutex::new(std::sync::Weak::new()); + +struct ScannerCohortMetricsOwner(Arc<()>); + +impl Default for ScannerCohortMetricsOwner { + fn default() -> Self { + let owner = Arc::new(()); + let mut current = SERVICE_COHORT_METRICS_OWNER + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + *current = Arc::downgrade(&owner); + write_service_cohort_metrics(0, 0.0, false); + Self(owner) + } +} + +impl Drop for ScannerCohortMetricsOwner { + fn drop(&mut self) { + let mut current = SERVICE_COHORT_METRICS_OWNER + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + if current.ptr_eq(&Arc::downgrade(&self.0)) { + *current = std::sync::Weak::new(); + write_service_cohort_metrics(0, 0.0, false); + } + } +} + +fn write_service_cohort_metrics(waiting: usize, oldest: f64, overflowed: bool) { + metrics::gauge!("rustfs_scanner_service_cohort_waiting").set(waiting as f64); + metrics::gauge!("rustfs_scanner_service_cohort_oldest_wait_seconds").set(oldest); + metrics::gauge!("rustfs_scanner_service_cohort_capacity_fallback").set(if overflowed { 1.0 } else { 0.0 }); +} + +struct ScannerCohortWait { + order: u64, + queued_at: Instant, + admitted: bool, + present: bool, +} + +/// Leader-local admission order, never evidence of completed scan coverage. +/// Retains at most 4096 members and 128 KiB of name payload, including the +/// cursor shared with its last member. Candidate selection borrows at most +/// 4096 inventory entries; the existing full inventory is not bounded here. +pub(crate) struct ScannerServiceCohort { + members: HashMap, ScannerCohortWait>>, + cursor: Option<(DataUsageCacheSource, Arc)>, + next_order: u64, + max_members: usize, + max_name_bytes: usize, + overflowed: bool, + metrics_owner: ScannerCohortMetricsOwner, + waiting: usize, + oldest_wait_at_refresh: f64, + #[cfg(test)] + metric_members_examined: usize, +} + +impl Default for ScannerServiceCohort { + fn default() -> Self { + Self { + members: HashMap::new(), + cursor: None, + next_order: 0, + max_members: SCANNER_SERVICE_COHORT_MAX_MEMBERS, + max_name_bytes: SCANNER_SERVICE_COHORT_MAX_NAME_BYTES, + overflowed: false, + metrics_owner: ScannerCohortMetricsOwner::default(), + waiting: 0, + oldest_wait_at_refresh: 0.0, + #[cfg(test)] + metric_members_examined: 0, + } + } +} + +impl ScannerServiceCohort { + pub(crate) fn refresh(&mut self, inventory: &HashMap>) { + let count = inventory + .values() + .fold(0usize, |count, buckets| count.saturating_add(buckets.len())); + let name_bytes = inventory + .values() + .flatten() + .fold(0usize, |bytes, bucket| bytes.saturating_add(bucket.name.len())); + self.overflowed = count > self.max_members || name_bytes > self.max_name_bytes; + for wait in self.members.values_mut().flat_map(HashMap::values_mut) { + wait.present = false; + } + for (source, buckets) in inventory { + for bucket in buckets { + if let Some(wait) = self + .members + .get_mut(source) + .and_then(|members| members.get_mut(bucket.name.as_str())) + { + wait.present = true; + } + } + } + self.members.retain(|_, buckets| { + buckets.retain(|_, wait| wait.present); + !buckets.is_empty() + }); + if self.members.values().flat_map(HashMap::values).all(|wait| wait.admitted) { + self.members.clear(); + self.next_order = 0; + } + if count == 0 { + self.cursor = None; + } + let mut member_count = self.members.values().map(HashMap::len).sum::(); + let mut retained_bytes = self + .members + .values() + .flat_map(HashMap::keys) + .map(|name| name.len()) + .sum::(); + let mut incoming = self.admission_candidates(inventory, true); + if incoming.is_empty() { + incoming = self.admission_candidates(inventory, false); + } + for (pool, set, bucket) in incoming { + if member_count >= self.max_members { + break; + } + if retained_bytes.saturating_add(bucket.len()) > self.max_name_bytes { + continue; + } + let Some(next_order) = self.next_order.checked_add(1) else { + self.overflowed = true; + break; + }; + let source = DataUsageCacheSource::new(pool, set); + let members = self.members.entry(source).or_default(); + if members.contains_key(bucket) { + continue; + } + let name: Arc = bucket.into(); + retained_bytes += name.len(); + member_count += 1; + members.insert( + name.clone(), + ScannerCohortWait { + order: self.next_order, + queued_at: Instant::now(), + admitted: false, + present: true, + }, + ); + self.cursor = Some((source, name)); + self.next_order = next_order; + } + self.refresh_metrics(); + } + + fn admission_candidates<'a>( + &self, + inventory: &'a HashMap>, + after_cursor: bool, + ) -> Vec<(usize, usize, &'a str)> { + let mut candidates = std::collections::BinaryHeap::new(); + for (source, buckets) in inventory { + for bucket in buckets { + let key = (source.pool_index, source.set_index, bucket.name.as_str()); + let after = self + .cursor + .as_ref() + .is_none_or(|(source, name)| key > (source.pool_index, source.set_index, name.as_ref())); + if after != after_cursor + || bucket.name.len() > self.max_name_bytes + || self + .members + .get(source) + .is_some_and(|members| members.contains_key(bucket.name.as_str())) + { + continue; + } + if candidates.len() < self.max_members { + candidates.push(key); + } else if candidates.peek().is_some_and(|last| key < *last) { + candidates.pop(); + candidates.push(key); + } + } + } + candidates.into_sorted_vec() + } + + pub(crate) fn order_set_indices(&self, sets: &[Arc]) -> Vec { + let ranks = self + .members + .iter() + .map(|(source, buckets)| { + ( + *source, + buckets + .values() + .filter(|wait| !wait.admitted) + .map(|wait| wait.order) + .min() + .unwrap_or(u64::MAX), + ) + }) + .collect::>(); + let mut indices = (0..sets.len()).collect::>(); + indices.sort_by_key(|index| { + ( + ranks + .get(&DataUsageCacheSource::new(sets[*index].pool_index, sets[*index].set_index)) + .copied() + .unwrap_or(u64::MAX), + *index, + ) + }); + indices + } + + pub(crate) fn order_buckets(&self, source: DataUsageCacheSource, buckets: &mut [BucketInfo]) { + let rank = |bucket: &str| { + self.members + .get(&source) + .and_then(|members| members.get(bucket)) + .filter(|wait| !wait.admitted) + .map_or(u64::MAX, |wait| wait.order) + }; + // Stable sorting preserves the existing dispatch order in the tail. + buckets.sort_by_key(|bucket| rank(&bucket.name)); + } + + pub(crate) fn record_admitted(&mut self, source: DataUsageCacheSource, bucket: &str) { + let Some(wait) = self.members.get_mut(&source).and_then(|members| members.get_mut(bucket)) else { + return; + }; + if wait.admitted { + return; + } + wait.admitted = true; + self.waiting -= 1; + if self.waiting == 0 { + self.oldest_wait_at_refresh = 0.0; + } + self.record_metrics(); + } + + #[cfg(test)] + pub(super) fn admitted_members(&self) -> Vec<(DataUsageCacheSource, String)> { + self.members + .iter() + .flat_map(|(source, buckets)| { + buckets + .iter() + .filter(|(_, wait)| wait.admitted) + .map(|(bucket, _)| (*source, bucket.to_string())) + }) + .collect() + } + + fn refresh_metrics(&mut self) { + // Oldest age is an inventory-refresh snapshot, not a per-admission + // scan of the cohort. Clear it immediately when no waiters remain. + let (mut waiting, mut oldest) = (0usize, 0.0f64); + for wait in self.members.values().flat_map(HashMap::values) { + #[cfg(test)] + { + self.metric_members_examined += 1; + } + if !wait.admitted { + waiting += 1; + oldest = oldest.max(wait.queued_at.elapsed().as_secs_f64()); + } + } + self.waiting = waiting; + self.oldest_wait_at_refresh = oldest; + self.record_metrics(); + } + + fn record_metrics(&self) { + // Serialize owner replacement, publication and retirement. A retired + // scanner must neither publish nor clear a replacement's gauges. + let current = SERVICE_COHORT_METRICS_OWNER + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + if current.ptr_eq(&Arc::downgrade(&self.metrics_owner.0)) { + write_service_cohort_metrics(self.waiting, self.oldest_wait_at_refresh, self.overflowed); + } + } +} + +pub(super) async fn wait_for_bucket_scan_permit( + semaphore: &Arc, + ctx: &CancellationToken, + complete: &CancellationToken, +) -> Option { + tokio::select! { + biased; + _ = complete.cancelled() => None, + _ = ctx.cancelled() => None, + permit = semaphore.clone().acquire_owned() => permit.ok(), + } +} + pub(super) fn bucket_usage_scan_order( buckets: &[BucketInfo], old_cache: &DataUsageCache, @@ -288,6 +594,305 @@ mod tests { use rustfs_scanner_metrics::metrics::{ScannerWorkSource, global_metrics}; use tokio::sync::oneshot; + #[derive(Default)] + struct RecordedGauge(AtomicU64); + + impl metrics::GaugeFn for RecordedGauge { + fn increment(&self, value: f64) { + self.set(f64::from_bits(self.0.load(Ordering::Relaxed)) + value); + } + fn decrement(&self, value: f64) { + self.increment(-value); + } + fn set(&self, value: f64) { + self.0.store(value.to_bits(), Ordering::Relaxed); + } + } + + #[derive(Default)] + struct CohortGaugeRecorder(StdMutex>>); + + impl metrics::Recorder for CohortGaugeRecorder { + fn describe_counter(&self, _: metrics::KeyName, _: Option, _: metrics::SharedString) {} + fn describe_gauge(&self, _: metrics::KeyName, _: Option, _: metrics::SharedString) {} + fn describe_histogram(&self, _: metrics::KeyName, _: Option, _: metrics::SharedString) {} + fn register_counter(&self, _: &metrics::Key, _: &metrics::Metadata<'_>) -> metrics::Counter { + metrics::Counter::noop() + } + fn register_histogram(&self, _: &metrics::Key, _: &metrics::Metadata<'_>) -> metrics::Histogram { + metrics::Histogram::noop() + } + fn register_gauge(&self, key: &metrics::Key, _: &metrics::Metadata<'_>) -> metrics::Gauge { + metrics::Gauge::from_arc( + self.0 + .lock() + .expect("gauge recorder") + .entry(key.name().to_string()) + .or_default() + .clone(), + ) + } + } + + impl CohortGaugeRecorder { + fn value(&self, name: &str) -> f64 { + f64::from_bits(self.0.lock().expect("gauge recorder")[name].0.load(Ordering::Relaxed)) + } + } + + #[test] + #[serial_test::serial] + fn service_cohort_metrics_retire_only_the_current_owner() { + let recorder = CohortGaugeRecorder::default(); + metrics::with_local_recorder(&recorder, || { + let mut old = ScannerServiceCohort::default(); + old.refresh(&cohort_inventory(&["old"])); + let mut current = ScannerServiceCohort { + max_members: 1, + ..Default::default() + }; + current.refresh(&cohort_inventory(&["a", "b"])); + for wait in current.members.values_mut().flat_map(HashMap::values_mut) { + wait.queued_at = Instant::now() - Duration::from_secs(60); + } + current.refresh_metrics(); + old.refresh(&cohort_inventory(&["old", "more"])); + drop(old); + assert_eq!(recorder.value("rustfs_scanner_service_cohort_waiting"), 1.0); + assert_eq!(recorder.value("rustfs_scanner_service_cohort_capacity_fallback"), 1.0); + assert!(recorder.value("rustfs_scanner_service_cohort_oldest_wait_seconds") >= 60.0); + drop(current); + for metric in [ + "rustfs_scanner_service_cohort_waiting", + "rustfs_scanner_service_cohort_oldest_wait_seconds", + "rustfs_scanner_service_cohort_capacity_fallback", + ] { + assert_eq!(recorder.value(metric), 0.0, "owner retirement must clear {metric}"); + } + }); + } + + #[test] + #[serial_test::serial] + fn service_cohort_admission_metrics_do_not_rescan_a_full_window() { + let recorder = CohortGaugeRecorder::default(); + metrics::with_local_recorder(&recorder, || { + let source = DataUsageCacheSource::new(0, 0); + let names = (0..SCANNER_SERVICE_COHORT_MAX_MEMBERS) + .map(|index| format!("bucket-{index:04}")) + .collect::>(); + let inventory = cohort_inventory(&names.iter().map(String::as_str).collect::>()); + let mut cohort = ScannerServiceCohort::default(); + cohort.refresh(&inventory); + assert_eq!(cohort.metric_members_examined, SCANNER_SERVICE_COHORT_MAX_MEMBERS); + assert_eq!(cohort.waiting, SCANNER_SERVICE_COHORT_MAX_MEMBERS); + for (index, name) in names.iter().enumerate() { + cohort.record_admitted(source, name); + for _ in 0..10 { + cohort.record_admitted(source, name); + cohort.record_admitted(source, "untracked-overflow-name"); + cohort.record_admitted(DataUsageCacheSource::new(99, 0), name); + } + assert_eq!(cohort.waiting, SCANNER_SERVICE_COHORT_MAX_MEMBERS - index - 1); + assert_eq!( + cohort.metric_members_examined, SCANNER_SERVICE_COHORT_MAX_MEMBERS, + "tracked, repeated and overflow admissions must not scan cohort members" + ); + } + assert_eq!(recorder.value("rustfs_scanner_service_cohort_waiting"), 0.0); + assert_eq!(recorder.value("rustfs_scanner_service_cohort_oldest_wait_seconds"), 0.0); + cohort.refresh(&inventory); + assert_eq!( + cohort.metric_members_examined, + 2 * SCANNER_SERVICE_COHORT_MAX_MEMBERS, + "one inventory refresh performs one metrics traversal" + ); + }); + } + + #[tokio::test] + #[serial_test::serial] + async fn service_cohort_queued_permit_cancel_and_drop_return_the_same_capacity() { + let semaphore = Arc::new(Semaphore::new(1)); + let active = Arc::new(AtomicUsize::new(0)); + let mut cohort = ScannerServiceCohort::default(); + cohort.refresh(&cohort_inventory(&["waiting"])); + let gauge_reset = DiskBucketScanGaugeReset::new("cohort-wait".to_string(), "0".to_string()); + record_disk_bucket_scans_queued(1, "cohort-wait", "0"); + record_disk_bucket_scans_active(0, "cohort-wait", "0"); + for cancel in [true, false] { + let held = semaphore + .clone() + .acquire_owned() + .await + .expect("hold the sole permit as a barrier"); + let ctx = CancellationToken::new(); + let complete = CancellationToken::new(); + let mut waiter = Box::pin(wait_for_bucket_scan_permit(&semaphore, &ctx, &complete)); + assert!( + futures::poll!(&mut waiter).is_pending(), + "the production wait must actually enqueue behind the barrier" + ); + assert_eq!(semaphore.available_permits(), 0); + if cancel { + ctx.cancel(); + assert!(waiter.as_mut().await.is_none()); + } + drop(waiter); + assert_eq!(semaphore.available_permits(), 0, "cancelling a waiter must not release the held permit"); + assert!(cohort.admitted_members().is_empty()); + drop(held); + assert_eq!(semaphore.available_permits(), 1, "no queued waiter may leak or steal released capacity"); + } + let ctx = CancellationToken::new(); + let complete = CancellationToken::new(); + let permit = wait_for_bucket_scan_permit(&semaphore, &ctx, &complete) + .await + .expect("same semaphore remains usable"); + let active_guard = DiskBucketScanActiveGuard::new(active.clone(), "cohort-wait".to_string(), "0".to_string()); + assert_eq!(active.load(Ordering::Relaxed), 1); + drop(active_guard); + drop(permit); + drop(gauge_reset); + assert_eq!(active.load(Ordering::Relaxed), 0); + assert_eq!(semaphore.available_permits(), 1); + let state = global_metrics() + .scanner_runtime_details_report() + .disk_bucket_scan_states + .into_iter() + .find(|state| state.pool == "cohort-wait" && state.set == "0") + .expect("fixture gauges"); + assert_eq!((state.queued, state.active), (0, 0)); + assert!(cohort.admitted_members().is_empty(), "permit ownership alone does not admit a bucket"); + } + + fn cohort_inventory(names: &[&str]) -> HashMap> { + HashMap::from([( + DataUsageCacheSource::new(0, 0), + names + .iter() + .map(|name| BucketInfo { + name: (*name).to_string(), + ..Default::default() + }) + .collect(), + )]) + } + + #[test] + #[serial_test::serial] + fn service_cohort_visits_fixed_members_within_service_round_bound() { + let inventory = cohort_inventory(&["a", "b", "c", "d", "e"]); + let source = DataUsageCacheSource::new(0, 0); + let mut cohort = ScannerServiceCohort::default(); + let mut admitted = HashSet::new(); + for _ in 0..3 { + cohort.refresh(&inventory); + let mut buckets = inventory[&source].clone(); + cohort.order_buckets(source, &mut buckets); + for bucket in buckets.iter().take(2) { + admitted.insert(bucket.name.clone()); + cohort.record_admitted(source, &bucket.name); + } + } + assert_eq!(admitted.len(), 5, "ceil(5/2) service rounds must include every member"); + } + + #[test] + #[serial_test::serial] + fn service_cohort_keeps_waiting_bootstrap_ahead_of_new_work() { + let source = DataUsageCacheSource::new(0, 0); + let mut cohort = ScannerServiceCohort::default(); + cohort.refresh(&cohort_inventory(&["a-hot", "z-bootstrap"])); + cohort.record_admitted(source, "a-hot"); + let queued_at = cohort.members[&source]["z-bootstrap"].queued_at; + let inventory = cohort_inventory(&["a-hot", "aaa-new-bootstrap", "z-bootstrap"]); + for _ in 0..10 { + cohort.refresh(&inventory); + let mut buckets = inventory[&source].clone(); + cohort.order_buckets(source, &mut buckets); + assert_eq!(buckets[0].name, "z-bootstrap"); + assert_eq!(buckets[1].name, "aaa-new-bootstrap"); + assert_eq!(cohort.members[&source]["z-bootstrap"].queued_at, queued_at); + } + } + + #[test] + #[serial_test::serial] + fn service_cohort_overflow_preserves_waiters_and_rotates_finished_windows() { + let source = DataUsageCacheSource::new(0, 0); + let mut cohort = ScannerServiceCohort { + max_members: 2, + max_name_bytes: 4, + ..Default::default() + }; + cohort.refresh(&cohort_inventory(&["aa", "bb"])); + cohort.record_admitted(source, "aa"); + for names in [["aa", "bb", "c"], ["aa", "bb", "d"]] { + let inventory = cohort_inventory(&names); + cohort.refresh(&inventory); + assert!(cohort.overflowed); + assert_eq!(cohort.members[&source].len(), 2); + assert!(cohort.members[&source].contains_key("aa")); + let mut fallback = inventory[&source].clone(); + fallback.reverse(); + cohort.order_buckets(source, &mut fallback); + assert_eq!(fallback[0].name, "bb", "overflow must not discard a waiting member's priority"); + assert_eq!(fallback.len(), 3, "unknown tail must remain dispatchable"); + } + cohort.record_admitted(source, "bb"); + let inventory = cohort_inventory(&["aa", "bb", "c", "d"]); + cohort.refresh(&inventory); + assert_eq!( + cohort.members[&source].keys().map(AsRef::as_ref).collect::>(), + HashSet::from(["c", "d"]) + ); + cohort.record_admitted(source, "c"); + cohort.record_admitted(source, "d"); + cohort.refresh(&inventory); + assert!( + cohort.members[&source].contains_key("aa"), + "finite inventory must wrap after the last window" + ); + cohort.refresh(&cohort_inventory(&["bb"])); + assert!(!cohort.overflowed); + assert_eq!(cohort.members[&source].len(), 1); + cohort.next_order = u64::MAX; + cohort.refresh(&cohort_inventory(&["bb", "c"])); + assert!(cohort.overflowed); + assert_eq!(cohort.members[&source].len(), 1); + } + + #[test] + #[serial_test::serial] + fn service_cohort_bounds_names_and_does_not_reset_duplicate_dirty_age() { + let source = DataUsageCacheSource::new(0, 0); + let mut cohort = ScannerServiceCohort { + max_members: 2, + max_name_bytes: 4, + ..Default::default() + }; + cohort.refresh(&cohort_inventory(&["aa", "bb", "long-name"])); + let queued_at = cohort.members[&source]["bb"].queued_at; + for _ in 0..10 { + cohort.refresh(&cohort_inventory(&["aa", "aa", "bb", "long-name"])); + assert_eq!(cohort.members.values().map(HashMap::len).sum::(), 2); + assert_eq!( + cohort + .members + .values() + .flat_map(HashMap::keys) + .map(|name| name.len()) + .sum::(), + 4 + ); + assert_eq!(cohort.members[&source]["bb"].queued_at, queued_at); + } + cohort.refresh(&cohort_inventory(&[])); + assert!(cohort.members.is_empty()); + assert!(cohort.cursor.is_none()); + } + fn active_bucket_drive_count(source: ScannerWorkSource, bucket: &str, drive: &str) -> u64 { global_metrics() .scanner_runtime_details_report() diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index fcae372ca..aad6403c6 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -117,6 +117,7 @@ impl ScannerIOCache for SetDisks { digest: scan_plan_digest, bucket_coverage_digest, requires_full_scan, + service_cohort, execution_digest, leader_epoch, tier_registry_generation, @@ -500,7 +501,13 @@ impl ScannerIOCache for SetDisks { let mut permutes = buckets.clone(); permutes.shuffle(&mut rand::rng()); - let scan_order = bucket_usage_scan_order(&permutes, &old_cache, &dirty_usage_buckets); + let mut scan_order = bucket_usage_scan_order(&permutes, &old_cache, &dirty_usage_buckets); + if let Some(cohort) = &service_cohort { + cohort + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .order_buckets(source, &mut scan_order); + } for bucket in scan_order.iter() { if let Some(c) = old_cache.find(&bucket.name) { @@ -558,6 +565,7 @@ impl ScannerIOCache for SetDisks { let remaining_bucket_work = Arc::new(AtomicUsize::new(buckets.len())); let bucket_work_complete = CancellationToken::new(); for (disk, worker_mode) in workers { + let service_cohort_clone = service_cohort.clone(); let bucket_rx_mutex_clone = bucket_rx_mutex.clone(); let bucket_tx_clone = bucket_tx.clone(); let remaining_bucket_work_clone = remaining_bucket_work.clone(); @@ -587,6 +595,18 @@ impl ScannerIOCache for SetDisks { let remote_session_id = uuid::Uuid::new_v4(); let mut remote_session_sequence = 0_u64; loop { + // Do not prefetch a FIFO member into an independently + // scheduled permit waiter: that can reorder admissions. + let permit_wait_start = Instant::now(); + let Some(_permit) = + wait_for_bucket_scan_permit(&disk_scan_semaphore_clone, &ctx_clone, &bucket_work_complete_clone).await + else { + break; + }; + if ctx_clone.is_cancelled() || budget_clone.budget_elapsed() { + break; + } + let permit_wait_elapsed = permit_wait_start.elapsed(); let bucket = tokio::select! { _ = bucket_work_complete_clone.cancelled() => break, _ = ctx_clone.cancelled() => break, @@ -600,41 +620,27 @@ impl ScannerIOCache for SetDisks { let mut work_guard = BucketWorkGuard::new(remaining_bucket_work_clone.clone(), bucket_work_complete_clone.clone()); - let permit_wait = ctx_clone.clone(); - let permit_wait_start = Instant::now(); - let _permit = tokio::select! { - permit = disk_scan_semaphore_clone.clone().acquire_owned() => match permit { - Ok(permit) => permit, - Err(_) => { - decrement_disk_bucket_scans_queued( - &queued_disk_bucket_scans_clone, - &pool_label_clone, - &set_label_clone, - ); - break; - }, - }, - _ = permit_wait.cancelled() => { - decrement_disk_bucket_scans_queued( - &queued_disk_bucket_scans_clone, - &pool_label_clone, - &set_label_clone, - ); - break; - }, - }; metrics::histogram!( METRIC_SCANNER_DISK_SCAN_WAIT_SECONDS, "pool" => pool_label_clone.clone(), "set" => set_label_clone.clone() ) - .record(permit_wait_start.elapsed().as_secs_f64()); + .record(permit_wait_elapsed.as_secs_f64()); decrement_disk_bucket_scans_queued(&queued_disk_bucket_scans_clone, &pool_label_clone, &set_label_clone); let _active_guard = DiskBucketScanActiveGuard::new( active_disk_bucket_scans_clone.clone(), pool_label_clone.clone(), set_label_clone.clone(), ); + if ctx_clone.is_cancelled() || budget_clone.budget_elapsed() { + break; + } + if let Some(cohort) = &service_cohort_clone { + cohort + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .record_admitted(source, &bucket.name); + } debug!( target: "rustfs::scanner::io", diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index da571d1ad..8cf9a7ffe 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -73,6 +73,7 @@ where scan_scope: ScannerBucketScanScope::default(), persisted_usage_baseline: None, requires_full_scan: true, + service_cohort: None, #[cfg(test)] resolved_scope_observer: None, }; @@ -90,6 +91,7 @@ pub(crate) struct ScannerCycleRequest { pub(crate) persisted_usage_baseline: Option, /// Scheduled maintenance must visit clean buckets even with a valid dirty scope. pub(crate) requires_full_scan: bool, + pub(crate) service_cohort: Option>>, #[cfg(test)] pub(crate) resolved_scope_observer: Option>, } @@ -184,6 +186,7 @@ where scan_scope, persisted_usage_baseline, requires_full_scan, + service_cohort, #[cfg(test)] resolved_scope_observer, } = request; @@ -275,6 +278,12 @@ where } bucket_plan_complete &= buckets_by_source.keys().copied().collect::>() == *expected_sources; bucket_plan_complete &= scanner_bucket_inventory_is_complete(&all_buckets, &buckets_by_source); + if bucket_plan_complete && let Some(cohort) = &service_cohort { + cohort + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .refresh(&buckets_by_source); + } let structural_scan_plan_digest = scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before)); let scan_plan_digest = scanner_bucket_work_digest(structural_scan_plan_digest, scan_mode, requires_full_scan); @@ -399,7 +408,32 @@ where let first_err_mutex: Arc>> = Arc::new(Mutex::new(None)); let mut wait_futs = Vec::new(); - for (results_index, set) in set_disks.iter().enumerate() { + let set_order = service_cohort.as_ref().map_or_else( + || (0..set_disks.len()).collect::>(), + |cohort| { + cohort + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .order_set_indices(&set_disks) + }, + ); + for results_index in set_order { + let set = &set_disks[results_index]; + // Acquire in dispatch order, not in independently scheduled tasks. + // A whole set still shares the existing parent budget; this is not + // a per-bucket quantum or a cross-source completion guarantee. + let permit_wait_start = Instant::now(); + let permit = tokio::select! { + biased; + _ = child_token.cancelled() => break, + permit = set_scan_semaphore.clone().acquire_owned() => match permit { + Ok(permit) => permit, + Err(_) => break, + }, + }; + if child_token.is_cancelled() || budget.budget_elapsed() { + break; + } let results_index_clone = results_index; // Clone the Arc to move it into the spawned task let set_clone: Arc = Arc::clone(set); @@ -414,7 +448,6 @@ where let scan_mode_clone = scan_mode; let results_mutex_clone = results_mutex.clone(); let first_err_mutex_clone = first_err_mutex.clone(); - let set_scan_semaphore_clone = set_scan_semaphore.clone(); let queued_set_scans_clone = queued_set_scans.clone(); let active_set_scans_clone = active_set_scans.clone(); @@ -437,6 +470,7 @@ where digest: structural_scan_plan_digest, bucket_coverage_digest, requires_full_scan, + service_cohort: service_cohort.clone(), execution_digest, leader_epoch, tier_registry_generation, @@ -448,15 +482,10 @@ where }; // Spawn task to run the scanner let scanner_fut = tokio::spawn(async move { - let permit_wait = child_token_clone.clone(); - let permit_wait_start = Instant::now(); - let _permit = tokio::select! { - permit = set_scan_semaphore_clone.acquire_owned() => match permit { - Ok(permit) => permit, - Err(_) => return, - }, - _ = permit_wait.cancelled() => return, - }; + let _permit = permit; + if child_token_clone.is_cancelled() || budget_clone.budget_elapsed() { + return; + } metrics::histogram!( METRIC_SCANNER_SET_SCAN_WAIT_SECONDS, "pool" => pool_label.clone(), diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index bd4add6df..58811e947 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -40,6 +40,7 @@ use time::OffsetDateTime; use uuid::Uuid; mod scoped_entry_fallback; +mod service_cohort; #[derive(Clone)] struct FixedWorkloadProvider { @@ -398,6 +399,7 @@ async fn scoped_scan_production_entry_preserves_deep_and_full_maintenance_work() persisted_usage_baseline: baseline, requires_full_scan, resolved_scope_observer: Some(observer), + service_cohort: None, }, ), ) @@ -470,6 +472,7 @@ async fn scoped_scan_same_cycle_maintenance_rewalks_after_root_delivery_failure( persisted_usage_baseline: None, requires_full_scan: false, resolved_scope_observer: None, + service_cohort: None, }, ), ) @@ -520,6 +523,7 @@ async fn scoped_scan_same_cycle_maintenance_rewalks_after_root_delivery_failure( persisted_usage_baseline: None, requires_full_scan, resolved_scope_observer: None, + service_cohort: None, }, ), ) @@ -1261,6 +1265,7 @@ async fn set_snapshot_reuse_requires_execution_identity_and_fences_stale_writers ctx.clone(), ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default()), ScannerBucketScanPlan { + service_cohort: None, buckets: Vec::new(), all_buckets: Arc::new(Vec::new()), scope: ScannerBucketScanScope::default(), diff --git a/crates/scanner/src/scanner_io/tests/scoped_entry_fallback.rs b/crates/scanner/src/scanner_io/tests/scoped_entry_fallback.rs index 88a7d91e2..ae04cc03c 100644 --- a/crates/scanner/src/scanner_io/tests/scoped_entry_fallback.rs +++ b/crates/scanner/src/scanner_io/tests/scoped_entry_fallback.rs @@ -127,6 +127,7 @@ async fn run_entry(store: &Arc, cycle: u64, selected: Option<&str>, exp scan_scope: ScannerBucketScanScope::default(), persisted_usage_baseline: root_before.0.clone().map(Bytes::from), requires_full_scan: false, + service_cohort: None, resolved_scope_observer: Some(observer), }, ), diff --git a/crates/scanner/src/scanner_io/tests/service_cohort.rs b/crates/scanner/src/scanner_io/tests/service_cohort.rs new file mode 100644 index 000000000..a2d5ffe3a --- /dev/null +++ b/crates/scanner/src/scanner_io/tests/service_cohort.rs @@ -0,0 +1,242 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::*; +use crate::data_usage_define::{DATA_USAGE_OBJ_NAME_PATH, read_config_with_revision}; + +async fn create_cohort_bucket(store: &ECStore, bucket: &str) { + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("fixture bucket"); + for set in store.all_set_disks() { + let mut reader = ScannerPutObjReader::from_vec(b"cohort".to_vec()); + set.put_object( + bucket, + "initial", + &mut reader, + &ScannerObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + .expect("fixture object and all rename tails should persist"); + } +} + +async fn run_cohort_cycle( + store: &Arc, + cohort: Arc>, + cycle: u64, + budget: Arc, +) -> (ScannerCycleResult, Option) { + let root_before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("root before candidate"); + let dirty_before = dirty_usage_buckets_for_tests(); + let generation_before = dirty_usage_generation(); + let (updates, mut receiver) = mpsc::channel(1); + let result = tokio::time::timeout( + Duration::from_secs(30), + nsscanner_with_storage_status_scoped( + store.as_ref(), + ScannerCycleRequest { + ctx: budget.token(), + budget, + updates, + want_cycle: cycle, + leader_epoch: 11, + scan_mode: HealScanMode::Normal, + scan_scope: ScannerBucketScanScope::default(), + persisted_usage_baseline: None, + requires_full_scan: false, + service_cohort: Some(cohort), + resolved_scope_observer: None, + }, + ), + ) + .await + .expect("cohort cycle should finish") + .expect("cohort cycle should return its status"); + let usage = receiver.recv().await; + assert!(receiver.recv().await.is_none()); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("root after candidate"), + root_before + ); + assert_eq!( + dirty_usage_buckets_for_tests(), + dirty_before, + "candidate production must not ACK pending work" + ); + assert_eq!(dirty_usage_generation(), generation_before); + (result, usage) +} + +#[tokio::test] +#[serial] +async fn service_cohort_production_dispatch_services_waiters_across_sources() { + let (_dir, store) = setup_two_pool_scanner_store().await; + clear_dirty_usage_buckets_for_tests(); + let hot = format!("a-hot-{}", Uuid::new_v4().simple()); + let bootstrap = format!("z-bootstrap-{}", Uuid::new_v4().simple()); + create_cohort_bucket(&store, &hot).await; + create_cohort_bucket(&store, &bootstrap).await; + let cohort = Arc::new(StdMutex::new(ScannerServiceCohort::default())); + let expected = store + .all_set_disks() + .iter() + .flat_map(|set| { + let source = DataUsageCacheSource::new(set.pool_index, set.set_index); + [(source, hot.clone()), (source, bootstrap.clone())] + }) + .collect::>(); + let mut seen = HashSet::new(); + for cycle in 1..=4 { + // Both newly bootstrapped names and repeated dirty work sort ahead of + // the original bootstrap in the old dirty-first policy. + if cycle > 1 { + create_cohort_bucket(&store, &format!("b-new-{cycle}")).await; + } + record_dirty_usage_bucket(&hot); + let ctx = CancellationToken::new(); + let budget = ScannerCycleBudget::new_with_progress_tracking( + &ctx, + ScannerCycleBudgetConfig { + max_objects: Some(1), + ..Default::default() + }, + ); + run_cohort_cycle(&store, cohort.clone(), cycle, budget.clone()).await; + assert!(budget.budget_elapsed()); + assert_eq!( + budget.progress().0, + 1, + "each round must reach one real object, not just mark an admission" + ); + let admitted = cohort + .lock() + .expect("cohort lock") + .admitted_members() + .into_iter() + .collect::>(); + let newly_admitted = admitted.difference(&seen).cloned().collect::>(); + assert_eq!( + newly_admitted.len(), + 1, + "serial parent object budget must stop before another bucket admission" + ); + seen.extend(newly_admitted); + } + assert!( + expected.is_subset(&seen), + "ongoing dirty/new bootstrap must not displace the original cohort" + ); + + // Admission fairness is not completed coverage: prior budgeted prefixes + // were observed under changing plans. A clean tail must not certify them. + let ctx = CancellationToken::new(); + let (result, usage) = + run_cohort_cycle(&store, cohort, 5, ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default())).await; + assert_eq!(result.status, ScannerCycleStatus::Incomplete); + assert!(usage.is_none(), "neither source has a complete mixed-plan baseline to publish"); + for set in store.all_set_disks() { + let mut cache = DataUsageCache::default(); + cache + .load(set, &path_join_buf(&[&hot, DATA_USAGE_CACHE_NAME])) + .await + .expect("retained hot prefix"); + assert!(!cache.info.snapshot_complete); + assert!(cache.info.scan_progress.is_some()); + assert!(cache.info.scan_plan_digest.is_none(), "mixed coverage must remain non-authoritative"); + } + clear_dirty_usage_buckets_for_tests(); +} + +#[tokio::test] +#[serial] +async fn service_cohort_fresh_complete_aggregate_preserves_reordered_sources() { + let (_dir, store) = setup_two_pool_scanner_store().await; + clear_dirty_usage_buckets_for_tests(); + for bucket in ["cohort-first", "cohort-second"] { + create_cohort_bucket(&store, bucket).await; + } + let sets = store.all_set_disks(); + let listing = store + .list_bucket_for_scanner(&BucketOptions::default()) + .await + .expect("fresh inventory"); + let inventory = listing + .set_buckets + .into_iter() + .map(|set| (DataUsageCacheSource::new(set.pool_index, set.set_index), set.buckets)) + .collect::>(); + let cohort = Arc::new(StdMutex::new(ScannerServiceCohort::default())); + { + let mut cohort = cohort.lock().expect("cohort lock"); + cohort.refresh(&inventory); + for bucket in &inventory[&DataUsageCacheSource::new(0, 0)] { + cohort.record_admitted(DataUsageCacheSource::new(0, 0), &bucket.name); + } + assert_eq!(cohort.order_set_indices(&sets), vec![1, 0]); + } + let ctx = CancellationToken::new(); + let (result, usage) = + run_cohort_cycle(&store, cohort, 1, ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default())).await; + assert_eq!(result.status, ScannerCycleStatus::Complete); + let usage = usage.expect("fresh complete aggregate"); + assert_eq!(usage.objects_total_count, 4); + assert!(usage.buckets_usage.values().all(|bucket| bucket.objects_count == 2)); + assert_eq!(usage.usage_snapshot_set_states.len(), 2); + assert_eq!( + usage + .usage_snapshot_set_states + .iter() + .map(|set| (set.pool_index, set.set_index)) + .collect::>(), + HashSet::from([(0, 0), (1, 0)]) + ); + clear_dirty_usage_buckets_for_tests(); +} + +#[tokio::test] +#[serial] +async fn service_cohort_cancelled_dispatch_does_not_consume_waiters_or_leak_permits() { + let (_dir, store) = setup_two_pool_scanner_store().await; + clear_dirty_usage_buckets_for_tests(); + let bucket = format!("cancel-{}", Uuid::new_v4().simple()); + create_cohort_bucket(&store, &bucket).await; + let cohort = Arc::new(StdMutex::new(ScannerServiceCohort::default())); + let ctx = CancellationToken::new(); + let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default()); + ctx.cancel(); + run_cohort_cycle(&store, cohort.clone(), 1, budget).await; + assert!(cohort.lock().expect("cohort lock").admitted_members().is_empty()); + let report = rustfs_scanner_metrics::metrics::global_metrics().scanner_runtime_details_report(); + assert!(report.active_bucket_drive_scans.is_empty()); + + let ctx = CancellationToken::new(); + let (result, usage) = + run_cohort_cycle(&store, cohort, 1, ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default())).await; + assert_eq!( + result.status, + ScannerCycleStatus::Complete, + "cancellation must not block a subsequent scan" + ); + assert_eq!(usage.expect("retry aggregate").objects_total_count, 2); + clear_dirty_usage_buckets_for_tests(); +}