use super::*;
use crate::data_usage_define::{DATA_USAGE_OBJ_NAME_PATH, read_config_with_revision};
use crate::storage_api::owner::EcstoreDiskAPI;
type DriveIdentities = HashMap<String, (Uuid, DataUsageCacheSource)>;
type WalkCounts = HashMap<(String, String, String), u64>;
async fn drive_identities(store: &ECStore) -> DriveIdentities {
let mut identities = HashMap::new();
let mut ids = HashSet::new();
for set in store.all_set_disks() {
let source = DataUsageCacheSource::new(set.pool_index, set.set_index);
for disk in scanner_set_disk_inventory(set.as_ref()).await {
let id = EcstoreDiskAPI::get_disk_id(disk.as_ref())
.await
.expect("fixture disk identity should be readable")
.expect("fixture disk must have a durable identity");
assert!(!id.is_nil());
assert!(ids.insert(id), "fixture disk identities must be unique");
let path = crate::ScannerDiskExt::path(disk.as_ref()).to_string_lossy().into_owned();
assert!(identities.insert(path, (id, source)).is_none());
}
}
assert_eq!(identities.len(), 8);
identities
}
fn walk_counts(drives: &DriveIdentities) -> WalkCounts {
rustfs_scanner_metrics::metrics::global_metrics()
.scanner_runtime_details_report()
.bucket_drive_results
.into_iter()
.filter(|result| drives.contains_key(&result.drive))
.map(|result| ((result.bucket, result.drive, result.result), result.count))
.collect()
}
async fn put_and_settle(store: &ECStore, bucket: &str, object: &str) {
let set = &store.pools[0].disk_set[0];
let mut reader = ScannerPutObjReader::from_vec(b"object".to_vec());
set.put_object(bucket, object, &mut reader, &ScannerObjectOptions::default())
.await
.expect("fixture object should persist");
let lock = set.new_ns_lock(bucket, object).await.expect("fixture namespace lock");
let _settled = lock
.get_write_lock(Duration::from_secs(30))
.await
.expect("quorum-ACK rename tail must settle before taking the activity baseline");
}
async fn create_bucket(store: &ECStore, bucket: &str) {
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("fixture bucket should be created");
put_and_settle(store, bucket, "initial").await;
}
async fn persist_baseline(store: &Arc<ECStore>, baseline: &DataUsageInfo) {
let mut baseline = baseline.clone();
baseline.usage_snapshot_converged = Some(true);
crate::save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&baseline).expect("baseline should encode"),
)
.await
.expect("fixture baseline should persist");
}
async fn run_entry(
store: &Arc<ECStore>,
cycle: u64,
selected: Option<&str>,
expect_walks: bool,
expect_activation: bool,
expect_prefix_scope: bool,
) -> DataUsageInfo {
let drives = drive_identities(store).await;
let inventory = store
.list_bucket_for_scanner(&BucketOptions::default())
.await
.expect("fixture inventory should be complete");
assert!(inventory.topology_complete);
let expected_walks = if expect_walks {
inventory
.set_buckets
.into_iter()
.flat_map(|set| {
let source = DataUsageCacheSource::new(set.pool_index, set.set_index);
set.buckets.into_iter().map(move |bucket| ((source, bucket.name), 1_u64))
})
.filter(|((_, bucket), _)| selected.is_none_or(|selected| bucket == selected))
.collect::<HashMap<_, _>>()
} else {
HashMap::new()
};
let root_before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("root baseline should be readable");
let dirty_before = dirty_usage_buckets_for_tests();
let generation_before = dirty_usage_generation();
let before = walk_counts(&drives);
let ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default());
let (updates, mut receiver) = mpsc::channel(1);
let (observer, observed) = tokio::sync::oneshot::channel();
let result = tokio::time::timeout(
Duration::from_secs(30),
nsscanner_with_storage_status_scoped(
store.as_ref(),
ScannerCycleRequest {
ctx,
budget,
updates,
want_cycle: cycle,
leader_epoch: 11,
scan_mode: HealScanMode::Normal,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: root_before.0.clone().map(Bytes::from),
observed_usage_candidate: None,
requires_full_scan: false,
service_cohort: None,
resolved_scope_observer: Some(observer),
},
),
)
.await
.expect("entry cycle should finish within the fixture deadline")
.expect("entry cycle should succeed");
assert_eq!(result.status, ScannerCycleStatus::Complete);
let activation_preflight = result.segment_reuse_activation_preflight;
let scope = observed.await.expect("production resolver should report its decision");
assert_eq!(
scope.selected_buckets.as_deref(),
selected.map(|name| HashSet::from([name.to_string()])).as_ref()
);
if let Some(selected) = selected {
assert_eq!(
scope.prefix_scope_for(selected).is_some(),
expect_prefix_scope,
"resolved prefix scope must match activation replay for cycle {cycle}"
);
}
let usage = receiver.recv().await.expect("one candidate should be delivered");
assert!(receiver.recv().await.is_none(), "there must be exactly one terminal candidate");
assert!(usage.usage_snapshot_complete);
assert!(!usage.usage_snapshot_partial);
assert_eq!(usage.scanner_cycle, Some(cycle));
assert_eq!(
drive_identities(store).await,
drives,
"drive identities must not change during the oracle"
);
let after = walk_counts(&drives);
let mut actual = HashMap::new();
for key in before.keys() {
assert!(after.contains_key(key), "metrics eviction would invalidate this exact-delta oracle");
}
for ((bucket, drive, outcome), count) in after {
let previous = before
.get(&(bucket.clone(), drive.clone(), outcome.clone()))
.copied()
.unwrap_or(0);
let delta = count.checked_sub(previous).expect("fixture counters must not reset");
if delta > 0 {
assert_eq!(outcome, "success", "no error or partial walker is expected");
*actual.entry((drives[&drive].1, bucket)).or_insert(0_u64) += delta;
}
}
assert_eq!(
actual, expected_walks,
"each listed source/bucket must have exactly the expected real walks"
);
assert!(activation_preflight.production_activation);
assert_eq!(activation_preflight.scanner_segment_reuse_activated, expect_activation);
let activation_blockers = activation_preflight.fail_closed_blockers().collect::<Vec<_>>();
if expect_activation {
assert_eq!(activation_blockers, Vec::<&str>::new());
} else if selected.is_some() && expect_walks {
assert!(
!activation_blockers.contains(&"missing_cold_zero_walk_oracle"),
"a complete scoped reuse cycle must carry the cold zero-walk oracle: cycle={cycle} selected={selected:?} blockers={activation_blockers:?}"
);
} else {
assert!(
activation_blockers.contains(&"missing_cold_zero_walk_oracle"),
"unscoped or same-cycle cache reuse must not claim the cold zero-walk oracle: cycle={cycle} selected={selected:?} expect_walks={expect_walks} blockers={activation_blockers:?}"
);
}
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("root after scan"),
root_before,
"producing a candidate must not replace the coordinator-owned root baseline"
);
assert_eq!(dirty_usage_generation(), generation_before);
assert!(
dirty_usage_buckets_for_tests() == dirty_before,
"candidate delivery must not ACK pending dirty buckets"
);
usage
}
fn record_segment_dirty_usage(bucket: &str) {
for producer in crate::segment_invalidation::SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION {
record_dirty_usage_object_from_producer(bucket, "hot-segment/object", producer);
}
}
fn replay_segment_dirty_usage(bucket: &str) {
replay_dirty_usage(
bucket,
ScannerDurableDirtyUsageReplayScope::TopLevelEntries {
entries: BTreeSet::from(["hot-segment".to_string()]),
},
);
}
fn replay_whole_bucket_dirty_usage(bucket: &str) {
replay_dirty_usage(bucket, ScannerDurableDirtyUsageReplayScope::WholeBucket);
}
fn replay_dirty_usage(bucket: &str, scope: ScannerDurableDirtyUsageReplayScope) {
replay_durable_dirty_usage_producer_record(
&encode_durable_dirty_usage_producer_replay_record(vec![ScannerDurableDirtyUsageReplayEntry {
bucket: bucket.to_string(),
generation: dirty_usage_generation(),
scope,
producers: crate::segment_invalidation::SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION
.into_iter()
.collect(),
}])
.expect("durable segment replay should encode"),
)
.expect("durable segment replay should restore producer authority");
}
fn run_scoped_entry_fallback_test<F, Fut>(thread_name: &'static str, test_fn: F)
where
F: FnOnce() -> Fut + Send + 'static,
Fut: std::future::Future<Output = ()> + 'static,
{
let handle = std::thread::Builder::new()
.name(thread_name.to_string())
.stack_size(32 * 1024 * 1024)
.spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("scoped entry fallback runtime should build");
runtime.block_on(test_fn());
})
.expect("scoped entry fallback test thread should spawn");
if let Err(payload) = handle.join() {
std::panic::resume_unwind(payload);
}
}
#[test]
#[serial]
fn scoped_entry_fallback_distinguishes_planned_scope_from_real_cold_walks() {
run_scoped_entry_fallback_test(
"scanner-scoped-entry-planned-scope",
scoped_entry_fallback_distinguishes_planned_scope_from_real_cold_walks_case,
);
}
async fn scoped_entry_fallback_distinguishes_planned_scope_from_real_cold_walks_case() {
let (_dir, store) = setup_two_pool_scanner_store().await;
clear_dirty_usage_buckets_for_tests();
let hot = format!("hot-{}", Uuid::new_v4().simple());
let cold = format!("cold-{}", Uuid::new_v4().simple());
create_bucket(&store, &hot).await;
create_bucket(&store, &cold).await;
record_segment_dirty_usage(&hot);
replay_segment_dirty_usage(&hot);
let baseline = run_entry(&store, 1, None, true, false, false).await;
persist_baseline(&store, &baseline).await;
run_entry(&store, 1, Some(&hot), false, false, false).await;
let usage = run_entry(&store, 2, Some(&hot), true, true, false).await;
persist_baseline(&store, &usage).await;
acknowledge_dirty_usage_generation(scanner_activity_epoch(), dirty_usage_generation())
.expect("durable whole-cycle publication should acknowledge the initial producer window");
put_and_settle(&store, &hot, "hot-segment/object").await;
record_dirty_usage_object_from_producer(
&hot,
"hot-segment/object",
crate::segment_invalidation::SegmentInvalidationProducerIdentity::PutObject,
);
replay_segment_dirty_usage(&hot);
let usage = run_entry(&store, 3, Some(&hot), true, true, true).await;
assert_eq!(usage.buckets_usage[&hot].objects_count, 2);
assert_eq!(usage.buckets_usage[&cold].objects_count, 1);
assert_eq!(usage.objects_total_count, 3);
persist_baseline(&store, &usage).await;
acknowledge_dirty_usage_generation(scanner_activity_epoch(), dirty_usage_generation())
.expect("durable prefix publication should acknowledge the typed suffix");
for index in 0..=MAX_DIRTY_USAGE_TOP_LEVEL_ENTRIES_PER_BUCKET {
record_dirty_usage_object_from_producer(
&hot,
&format!("overflow-{index}/object"),
crate::segment_invalidation::SegmentInvalidationProducerIdentity::PutObject,
);
}
replay_whole_bucket_dirty_usage(&hot);
let usage = run_entry(&store, 4, Some(&hot), true, true, false).await;
assert_eq!(usage.objects_total_count, 3);
persist_baseline(&store, &usage).await;
acknowledge_dirty_usage_generation(scanner_activity_epoch(), dirty_usage_generation())
.expect("durable whole-bucket fallback should acknowledge the overflow window");
record_dirty_usage_object_from_producer(
&hot,
"hot-segment/object",
crate::segment_invalidation::SegmentInvalidationProducerIdentity::Unknown,
);
let usage = run_entry(&store, 5, Some(&hot), true, false, false).await;
assert_eq!(usage.objects_total_count, 3);
clear_dirty_usage_buckets_for_tests();
}
#[test]
#[serial]
fn scoped_entry_fallback_rejects_invalid_persisted_baseline_at_the_walker() {
run_scoped_entry_fallback_test(
"scanner-scoped-entry-invalid-baseline",
scoped_entry_fallback_rejects_invalid_persisted_baseline_at_the_walker_case,
);
}
async fn scoped_entry_fallback_rejects_invalid_persisted_baseline_at_the_walker_case() {
let (_dir, store) = setup_two_pool_scanner_store().await;
clear_dirty_usage_buckets_for_tests();
let hot = format!("hot-{}", Uuid::new_v4().simple());
let cold = format!("cold-{}", Uuid::new_v4().simple());
create_bucket(&store, &hot).await;
create_bucket(&store, &cold).await;
record_dirty_usage_bucket(&hot);
let baseline = run_entry(&store, 1, None, true, false, false).await;
for (index, kind) in [
"malformed",
"unconverged",
"missing-set",
"wrong-source",
"mixed-plan",
"wrong-epoch",
]
.into_iter()
.enumerate()
{
let mut candidate = baseline.clone();
candidate.usage_snapshot_converged = Some(true);
match kind {
"unconverged" => candidate.usage_snapshot_converged = Some(false),
"missing-set" => {
candidate.usage_snapshot_set_states.pop();
}
"wrong-source" => candidate.usage_snapshot_set_states[0].set_index = 99,
"mixed-plan" => candidate.usage_snapshot_set_states[1].scan_plan_digest = Some([0xA5; 32]),
"wrong-epoch" => candidate.usage_snapshot_set_states[0].scanner_epoch = Some(10),
"malformed" => {}
_ => unreachable!(),
}
let bytes = if kind == "malformed" {
b"{broken".to_vec()
} else {
serde_json::to_vec(&candidate).expect("candidate JSON")
};
crate::save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes)
.await
.expect("negative baseline should persist");
let usage = run_entry(
&store,
u64::try_from(index).expect("fixture cycle index should fit") + 2,
None,
true,
false,
false,
)
.await;
assert_eq!(usage.objects_total_count, 2, "{kind}");
assert_eq!(usage.buckets_usage[&cold].objects_count, 1, "{kind}");
}
clear_dirty_usage_buckets_for_tests();
}
#[test]
#[serial]
fn scoped_entry_fallback_covers_overflow_and_new_bucket_inventory() {
run_scoped_entry_fallback_test(
"scanner-scoped-entry-overflow-inventory",
scoped_entry_fallback_covers_overflow_and_new_bucket_inventory_case,
);
}
async fn scoped_entry_fallback_covers_overflow_and_new_bucket_inventory_case() {
let (_dir, store) = setup_two_pool_scanner_store().await;
clear_dirty_usage_buckets_for_tests();
let hot = format!("hot-{}", Uuid::new_v4().simple());
create_bucket(&store, &hot).await;
record_dirty_usage_bucket(&hot);
let baseline = run_entry(&store, 1, None, true, false, false).await;
persist_baseline(&store, &baseline).await;
for index in 0..=crate::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES {
record_dirty_usage_bucket(&format!("overflow-{index}"));
}
assert!(dirty_usage_buckets_for_tests().len() > crate::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES);
let usage = run_entry(&store, 2, None, true, false, false).await;
assert_eq!(usage.objects_total_count, 1);
clear_dirty_usage_buckets_for_tests();
record_dirty_usage_bucket(&hot);
let new_bucket = format!("new-{}", Uuid::new_v4().simple());
create_bucket(&store, &new_bucket).await;
let usage = run_entry(&store, 3, None, true, false, false).await;
assert_eq!(usage.objects_total_count, 2);
assert_eq!(usage.buckets_usage[&new_bucket].objects_count, 1);
clear_dirty_usage_buckets_for_tests();
}