mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-06 03:59:14 +00:00
fix(scanner): bind bucket cache reuse to scan work requirements
Carry stable scan mode and full-maintenance requirements in the existing opaque bucket digest before local and remote cache admission. Different requirements cannot replay a same-cycle Normal cache after root delivery failure; matching requirements remain reusable for the same intent. Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -273,6 +273,7 @@ pub struct ScannerBucketScanPlan {
|
||||
digest: DataUsageScanPlanDigest,
|
||||
/// Includes mutation generations even when the set planner uses a structural digest.
|
||||
bucket_coverage_digest: DataUsageScanPlanDigest,
|
||||
requires_full_scan: bool,
|
||||
leader_epoch: u64,
|
||||
tier_registry_generation: u64,
|
||||
/// Epoch captured once for the whole scanner cycle. `None` is retained
|
||||
@@ -338,6 +339,29 @@ fn scanner_bucket_inventory_is_complete(
|
||||
covered.len() == inventory.len()
|
||||
}
|
||||
|
||||
// Bind known work requirements before both local and remote cache admission.
|
||||
// Matching requirements remain reusable for the same intent; this is not a
|
||||
// new deadline or a durable generation for newly due maintenance.
|
||||
fn scanner_bucket_work_digest(
|
||||
scan_plan_digest: DataUsageScanPlanDigest,
|
||||
scan_mode: HealScanMode,
|
||||
requires_full_scan: bool,
|
||||
) -> DataUsageScanPlanDigest {
|
||||
if scan_mode == HealScanMode::Normal && !requires_full_scan {
|
||||
return scan_plan_digest;
|
||||
}
|
||||
let mut hasher = Sha256::new();
|
||||
hasher.update(b"scanner-bucket-work-v1");
|
||||
hasher.update(scan_plan_digest.0);
|
||||
hasher.update([match scan_mode {
|
||||
HealScanMode::Unknown => 0,
|
||||
HealScanMode::Normal => 1,
|
||||
HealScanMode::Deep => 2,
|
||||
}]);
|
||||
hasher.update([u8::from(requires_full_scan || scan_mode == HealScanMode::Deep)]);
|
||||
DataUsageScanPlanDigest(hasher.finalize().into())
|
||||
}
|
||||
|
||||
fn scanner_bucket_cache_digest(
|
||||
scan_plan_digest: DataUsageScanPlanDigest,
|
||||
dirty_generation: Option<u64>,
|
||||
|
||||
@@ -116,6 +116,7 @@ impl ScannerIOCache for SetDisks {
|
||||
scope,
|
||||
digest: scan_plan_digest,
|
||||
bucket_coverage_digest,
|
||||
requires_full_scan,
|
||||
leader_epoch,
|
||||
tier_registry_generation,
|
||||
publication_epoch,
|
||||
@@ -124,6 +125,7 @@ impl ScannerIOCache for SetDisks {
|
||||
pending_maintenance_work,
|
||||
cache_cycle_floor,
|
||||
} = scan_plan;
|
||||
let bucket_work_digest = scanner_bucket_work_digest(bucket_coverage_digest, scan_mode, requires_full_scan);
|
||||
let pool_label = self.pool_index.to_string();
|
||||
let set_label = self.set_index.to_string();
|
||||
|
||||
@@ -635,7 +637,7 @@ impl ScannerIOCache for SetDisks {
|
||||
|
||||
let cache_name = path_join_buf(&[&bucket.name, DATA_USAGE_CACHE_NAME]);
|
||||
let bucket_scan_plan_digest =
|
||||
scanner_bucket_cache_digest(bucket_coverage_digest, dirty_usage_buckets_clone.get(&bucket.name).copied());
|
||||
scanner_bucket_cache_digest(bucket_work_digest, dirty_usage_buckets_clone.get(&bucket.name).copied());
|
||||
|
||||
if let Some(server_epoch) = remote_server_epoch {
|
||||
let request_sequence = remote_session_sequence;
|
||||
|
||||
@@ -433,6 +433,7 @@ where
|
||||
scope: scan_scope.clone(),
|
||||
digest: scan_plan_digest,
|
||||
bucket_coverage_digest,
|
||||
requires_full_scan,
|
||||
leader_epoch,
|
||||
tier_registry_generation,
|
||||
publication_epoch,
|
||||
|
||||
@@ -1330,6 +1330,56 @@ fn dirty_bucket_cache_digest_changes_with_generation() {
|
||||
assert!(!cache_snapshot_is_current(&cache, "photos", source, 11, 0, second));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scoped_scan_bucket_work_proof_fences_same_cycle_cache() {
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let structural_plan = DataUsageScanPlanDigest([9; 32]);
|
||||
let normal_plan = scanner_bucket_work_digest(structural_plan, HealScanMode::Normal, false);
|
||||
assert_eq!(normal_plan, structural_plan, "ordinary work keeps the existing digest contract");
|
||||
for (scan_mode, full) in [(HealScanMode::Deep, false), (HealScanMode::Normal, true)] {
|
||||
let requested_plan = scanner_bucket_work_digest(structural_plan, scan_mode, full);
|
||||
assert_ne!(requested_plan, normal_plan);
|
||||
let mut cache = DataUsageCache {
|
||||
info: DataUsageCacheInfo {
|
||||
name: "cold".to_string(),
|
||||
next_cycle: 7,
|
||||
leader_epoch: 11,
|
||||
last_update: Some(SystemTime::UNIX_EPOCH),
|
||||
source: Some(source),
|
||||
snapshot_complete: true,
|
||||
scan_plan_digest: Some(normal_plan),
|
||||
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
cache.replace("cold", "", DataUsageEntry::default());
|
||||
assert!(cache_snapshot_is_current(&cache, "cold", source, 7, 11, normal_plan));
|
||||
assert!(matches!(
|
||||
current_cache_root_or_prepare(&mut cache, "cold", source, 7, 11, requested_plan, true),
|
||||
DataUsageCacheScanState::Prepared {
|
||||
outcome: DataUsageCachePrepareOutcome::Reset,
|
||||
..
|
||||
}
|
||||
));
|
||||
assert!(cache.cache.is_empty(), "different work requirements must enter a fresh walk");
|
||||
cache.replace("cold", "", DataUsageEntry::default());
|
||||
cache.info.snapshot_complete = true;
|
||||
cache.info.last_update = Some(SystemTime::UNIX_EPOCH);
|
||||
assert!(
|
||||
matches!(
|
||||
current_cache_root_or_prepare(&mut cache, "cold", source, 7, 11, requested_plan, true),
|
||||
DataUsageCacheScanState::Current(_)
|
||||
),
|
||||
"completed matching work may satisfy the same intent retry"
|
||||
);
|
||||
}
|
||||
assert_eq!(
|
||||
scanner_bucket_work_digest(structural_plan, HealScanMode::Deep, false),
|
||||
scanner_bucket_work_digest(structural_plan, HealScanMode::Deep, true)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_cache_lock_resource_is_scoped_to_cache_source() {
|
||||
let cache_name = "photos/.usage-cache.bin";
|
||||
|
||||
@@ -415,6 +415,106 @@ async fn scoped_scan_production_entry_preserves_deep_and_full_maintenance_work()
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn scoped_scan_same_cycle_maintenance_rewalks_after_root_delivery_failure() {
|
||||
for (scan_mode, requires_full_scan) in [(HealScanMode::Deep, false), (HealScanMode::Normal, true)] {
|
||||
let (_temp_dir, store) = setup_two_pool_scanner_store().await;
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
for bucket in ["hot-bucket", "cold-bucket"] {
|
||||
store
|
||||
.make_bucket(bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("bucket should be created");
|
||||
let mut reader = ScannerPutObjReader::from_vec(b"initial".to_vec());
|
||||
store.pools[0].disk_set[0]
|
||||
.put_object(bucket, "initial", &mut reader, &ScannerObjectOptions::default())
|
||||
.await
|
||||
.expect("initial object should persist");
|
||||
}
|
||||
let ctx = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default());
|
||||
let (updates, receiver) = mpsc::channel(1);
|
||||
drop(receiver);
|
||||
let failed = tokio::time::timeout(
|
||||
Duration::from_secs(30),
|
||||
nsscanner_with_storage_status_scoped(
|
||||
store.as_ref(),
|
||||
ScannerCycleRequest {
|
||||
ctx,
|
||||
budget,
|
||||
updates,
|
||||
want_cycle: 7,
|
||||
leader_epoch: 11,
|
||||
scan_mode: HealScanMode::Normal,
|
||||
scan_scope: ScannerBucketScanScope::default(),
|
||||
persisted_usage_baseline: None,
|
||||
requires_full_scan: false,
|
||||
resolved_scope_observer: None,
|
||||
},
|
||||
),
|
||||
)
|
||||
.await
|
||||
.expect("normal scan should finish")
|
||||
.expect_err("root delivery must fail after bucket cache persistence");
|
||||
assert!(failed.to_string().contains("receiver closed"), "{failed}");
|
||||
let cache_name = path_join_buf(&["cold-bucket", DATA_USAGE_CACHE_NAME]);
|
||||
let mut cached = DataUsageCache::default();
|
||||
cached
|
||||
.load(store.pools[0].disk_set[0].clone(), &cache_name)
|
||||
.await
|
||||
.expect("normal bucket cache should have committed");
|
||||
assert!(cached.info.snapshot_complete);
|
||||
assert_eq!(cached.info.next_cycle, 7);
|
||||
assert_eq!(
|
||||
cached
|
||||
.checked_flatten("cold-bucket")
|
||||
.expect("cached root should be valid")
|
||||
.objects,
|
||||
1
|
||||
);
|
||||
|
||||
let mut reader = ScannerPutObjReader::from_vec(b"maintenance".to_vec());
|
||||
store.pools[0].disk_set[0]
|
||||
.put_object("cold-bucket", "new", &mut reader, &ScannerObjectOptions::default())
|
||||
.await
|
||||
.expect("new cold object should persist");
|
||||
record_dirty_usage_bucket("hot-bucket");
|
||||
let ctx = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default());
|
||||
let (updates, mut receiver) = mpsc::channel(1);
|
||||
let result = tokio::time::timeout(
|
||||
Duration::from_secs(30),
|
||||
nsscanner_with_storage_status_scoped(
|
||||
store.as_ref(),
|
||||
ScannerCycleRequest {
|
||||
ctx,
|
||||
budget,
|
||||
updates,
|
||||
want_cycle: 7,
|
||||
leader_epoch: 11,
|
||||
scan_mode,
|
||||
scan_scope: ScannerBucketScanScope::default(),
|
||||
persisted_usage_baseline: None,
|
||||
requires_full_scan,
|
||||
resolved_scope_observer: None,
|
||||
},
|
||||
),
|
||||
)
|
||||
.await
|
||||
.expect("maintenance scan should finish")
|
||||
.expect("maintenance scan should succeed");
|
||||
assert_eq!(result.status, ScannerCycleStatus::Complete);
|
||||
let snapshot = receiver.recv().await.expect("maintenance snapshot should be published");
|
||||
assert_eq!(snapshot.scanner_cycle, Some(7));
|
||||
assert_eq!(
|
||||
snapshot.buckets_usage["cold-bucket"].objects_count, 2,
|
||||
"{scan_mode:?}/full={requires_full_scan} must not replay the same-cycle Normal root"
|
||||
);
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn data_usage_publish_fails_when_receiver_is_closed() {
|
||||
let (updates, receiver) = mpsc::channel(1);
|
||||
|
||||
Reference in New Issue
Block a user