feat(scanner): preserve bounded bootstrap admission fairness (#7260)

feat(scanner): retain bounded bootstrap admission fairness

Keep a leader-local bounded cohort across scanner retries and preserve
waiting bucket priority during dirty arrivals and capacity overflow.
Order source permits in the dispatcher without changing result identity,
parent budgets, explicit cycle timing, or persistent coverage evidence.

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-06 16:45:57 +08:00
committed by GitHub
parent cb3100a252
commit 0ee5408b94
9 changed files with 1034 additions and 48 deletions
+34 -10
View File
@@ -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<Arc<StdMutex<crate::scanner_io::ScannerServiceCohort>>>,
}
#[cfg(test)]
async fn run_data_scanner_cycle<S>(
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<S>(
cycle_revision: &mut DataUsageCacheRevision,
leader_epoch: u64,
cycle_budget: Arc<ScannerCycleBudget>,
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(),
)
+74 -2
View File
@@ -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() {
+2
View File
@@ -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<Arc<StdMutex<ScannerServiceCohort>>>,
// 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)]
+605
View File
@@ -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<std::sync::Weak<()>> = 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<DataUsageCacheSource, HashMap<Arc<str>, ScannerCohortWait>>,
cursor: Option<(DataUsageCacheSource, Arc<str>)>,
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<DataUsageCacheSource, Vec<BucketInfo>>) {
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::<usize>();
let mut retained_bytes = self
.members
.values()
.flat_map(HashMap::keys)
.map(|name| name.len())
.sum::<usize>();
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<str> = 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<DataUsageCacheSource, Vec<BucketInfo>>,
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<SetDisks>]) -> Vec<usize> {
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::<HashMap<_, _>>();
let mut indices = (0..sets.len()).collect::<Vec<_>>();
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<Semaphore>,
ctx: &CancellationToken,
complete: &CancellationToken,
) -> Option<tokio::sync::OwnedSemaphorePermit> {
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<HashMap<String, Arc<RecordedGauge>>>);
impl metrics::Recorder for CohortGaugeRecorder {
fn describe_counter(&self, _: metrics::KeyName, _: Option<metrics::Unit>, _: metrics::SharedString) {}
fn describe_gauge(&self, _: metrics::KeyName, _: Option<metrics::Unit>, _: metrics::SharedString) {}
fn describe_histogram(&self, _: metrics::KeyName, _: Option<metrics::Unit>, _: 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::<Vec<_>>();
let inventory = cohort_inventory(&names.iter().map(String::as_str).collect::<Vec<_>>());
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<DataUsageCacheSource, Vec<BucketInfo>> {
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<&str>>(),
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::<usize>(), 2);
assert_eq!(
cohort
.members
.values()
.flat_map(HashMap::keys)
.map(|name| name.len())
.sum::<usize>(),
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()
+31 -25
View File
@@ -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",
+40 -11
View File
@@ -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<Bytes>,
/// Scheduled maintenance must visit clean buckets even with a valid dirty scope.
pub(crate) requires_full_scan: bool,
pub(crate) service_cohort: Option<Arc<StdMutex<ScannerServiceCohort>>>,
#[cfg(test)]
pub(crate) resolved_scope_observer: Option<tokio::sync::oneshot::Sender<ScannerBucketScanScope>>,
}
@@ -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::<HashSet<_>>() == *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<Mutex<Option<Error>>> = 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::<Vec<_>>(),
|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<SetDisks> = 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(),
+5
View File
@@ -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(),
@@ -127,6 +127,7 @@ async fn run_entry(store: &Arc<ECStore>, 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),
},
),
@@ -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<ECStore>,
cohort: Arc<StdMutex<ScannerServiceCohort>>,
cycle: u64,
budget: Arc<ScannerCycleBudget>,
) -> (ScannerCycleResult, Option<DataUsageInfo>) {
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::<HashSet<_>>();
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::<HashSet<_>>();
let newly_admitted = admitted.difference(&seen).cloned().collect::<Vec<_>>();
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::<HashMap<_, _>>();
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<_>>(),
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();
}