fix(scanner): retain scoped partial coverage across dirty plans

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-05 16:59:00 +08:00
parent d5a1530409
commit a0c32e7c6b
12 changed files with 588 additions and 45 deletions
+141
View File
@@ -400,6 +400,52 @@ impl DataUsageScanCheckpoint {
}
}
/// Durable scope of a bucket checkpoint, independent of namespace mutation counters.
#[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct DataUsageScanIdentity {
pub version: u16,
pub bucket_incarnation: uuid::Uuid,
pub set_layout: DataUsageScanPlanDigest,
pub publication_epoch: u64,
pub tier_registry_generation: u64,
}
impl Serialize for DataUsageScanIdentity {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
let mut map = serializer.serialize_map(Some(5))?;
map.serialize_entry("version", &self.version)?;
map.serialize_entry("bucket_incarnation", &self.bucket_incarnation)?;
map.serialize_entry("set_layout", &self.set_layout)?;
map.serialize_entry("publication_epoch", &self.publication_epoch)?;
map.serialize_entry("tier_registry_generation", &self.tier_registry_generation)?;
map.end()
}
}
impl DataUsageScanIdentity {
pub(crate) fn is_valid(&self) -> bool {
self.version == 1 && !self.bucket_incarnation.is_nil()
}
}
/// A forward coverage sweep may span budgets, but not authorize mixed mutation generations.
#[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct DataUsageScanProgress {
pub started_plan: DataUsageScanPlanDigest,
pub requested_plan: DataUsageScanPlanDigest,
}
impl Serialize for DataUsageScanProgress {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
let mut map = serializer.serialize_map(Some(2))?;
map.serialize_entry("started_plan", &self.started_plan)?;
map.serialize_entry("requested_plan", &self.requested_plan)?;
map.end()
}
}
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct DataUsageEntryInfo {
pub name: String,
@@ -470,6 +516,10 @@ pub struct DataUsageCacheInfo {
#[serde(default)]
pub scan_checkpoint: Option<DataUsageScanCheckpoint>,
#[serde(default)]
pub scan_identity: Option<DataUsageScanIdentity>,
#[serde(default)]
pub scan_progress: Option<DataUsageScanProgress>,
#[serde(default)]
pub pending_heals: Vec<PendingScannerHeal>,
#[serde(default)]
pub object_lock: Option<Arc<ObjectLockConfiguration>>,
@@ -513,6 +563,8 @@ impl Serialize for DataUsageCacheInfo {
// Keep this metadata map-encoded so older readers can ignore fields
// appended by newer scanner versions during rolling upgrades.
let field_count = 16
+ usize::from(self.scan_identity.is_some())
+ usize::from(self.scan_progress.is_some())
+ usize::from(self.tier_registry_generation.is_some())
+ usize::from(!self.size_reconciliation.is_empty())
+ usize::from(self.lkg_snapshot_complete)
@@ -531,6 +583,12 @@ impl Serialize for DataUsageCacheInfo {
state.serialize_entry("failed_objects", &self.failed_objects)?;
state.serialize_entry("scan_resume_after", &self.scan_resume_after)?;
state.serialize_entry("scan_checkpoint", &self.scan_checkpoint)?;
if let Some(identity) = self.scan_identity {
state.serialize_entry("scan_identity", &identity)?;
}
if let Some(progress) = self.scan_progress {
state.serialize_entry("scan_progress", &progress)?;
}
state.serialize_entry("pending_heals", &self.pending_heals)?;
state.serialize_entry("object_lock", &self.object_lock)?;
state.serialize_entry("source", &self.source)?;
@@ -695,6 +753,89 @@ impl DataUsageCache {
}
}
pub(crate) fn prepare_bucket_checkpoint(
&mut self,
name: &str,
next_cycle: u64,
leader_epoch: u64,
source: DataUsageCacheSource,
scan_plan_digest: DataUsageScanPlanDigest,
identity: DataUsageScanIdentity,
) -> DataUsageCachePrepareOutcome {
if self.info.next_cycle > next_cycle {
return DataUsageCachePrepareOutcome::RejectedNewerCycle;
}
if self.info.leader_epoch > leader_epoch {
return DataUsageCachePrepareOutcome::RejectedNewerLeader;
}
let reusable = identity.is_valid()
&& name != DATA_USAGE_ROOT
&& self.info.name == name
&& self.info.source == Some(source)
&& self.info.leader_epoch == leader_epoch
&& self.info.cache_key_format == DATA_USAGE_CACHE_KEY_FORMAT
&& self.info.scan_identity == Some(identity)
&& self.info.tier_registry_generation == Some(identity.tier_registry_generation)
&& (self.cache.is_empty() || self.checked_flatten_complete_scope(name).is_some());
if !reusable {
let keep_debts = self.info.name == name
&& self
.info
.scan_identity
.is_none_or(|previous| previous.bucket_incarnation == identity.bucket_incarnation);
let (pending_heals, size_reconciliation) = if keep_debts {
(
std::mem::take(&mut self.info.pending_heals),
std::mem::take(&mut self.info.size_reconciliation),
)
} else {
(Vec::new(), HashMap::new())
};
*self = Self::default();
self.info.pending_heals = pending_heals;
self.info.size_reconciliation = size_reconciliation;
}
let cursor_is_valid = match (&self.info.scan_checkpoint, &self.info.scan_resume_after) {
(None, None) => true,
(Some(checkpoint), Some(resume)) => {
checkpoint.version == DATA_USAGE_SCAN_CHECKPOINT_VERSION
&& checkpoint.resume_after == *resume
&& resume.strip_prefix(name).is_some_and(|suffix| suffix.starts_with('/'))
&& self.find(resume).is_some()
}
_ => false,
};
if !cursor_is_valid {
self.info.scan_progress = None;
}
self.info.name = name.to_owned();
self.info.next_cycle = next_cycle;
self.info.leader_epoch = leader_epoch;
self.info.source = Some(source);
self.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
self.info.tier_registry_generation = Some(identity.tier_registry_generation);
self.info.scan_identity = Some(identity);
self.info.snapshot_complete = false;
if let Some(progress) = &mut self.info.scan_progress {
progress.requested_plan = scan_plan_digest;
} else {
self.info.scan_progress = Some(DataUsageScanProgress {
started_plan: scan_plan_digest,
requested_plan: scan_plan_digest,
});
self.info.scan_resume_after = None;
self.info.scan_checkpoint = None;
}
// Old readers do not understand coverage sweeps. An absent plan makes
// their existing prepare path rebuild instead of promoting mixed data.
self.info.scan_plan_digest = None;
if reusable {
DataUsageCachePrepareOutcome::Reused
} else {
DataUsageCachePrepareOutcome::Reset
}
}
fn ensure_cache_save_metrics_registered() {
CACHE_SAVE_METRICS_ONCE.call_once(|| {
describe_counter!(
@@ -728,6 +728,14 @@ async fn scan_and_persist_local_bucket(
DataUsageCacheReuseOptions {
require_source: true,
tier_registry_generation: Some(tier_registry_generation),
checkpoint_identity: crate::scanner_io::scanner_bucket_checkpoint_identity(
&set,
&bucket,
expected_publication_epoch,
tier_registry_generation,
)
.await
.ok(),
},
);
match scan_state {
@@ -821,6 +829,21 @@ async fn scan_and_persist_local_bucket(
ScannerDiskScanOutcome::Partial(cache) => (cache, Some(RemoteScannerFrameResult::Partial)),
ScannerDiskScanOutcome::NamespaceNotFound(cache) => (cache, Some(RemoteScannerFrameResult::NamespaceNotFound)),
};
if let Some(expected) = cache.info.scan_identity
&& crate::scanner_io::scanner_bucket_checkpoint_identity(
&set,
&bucket,
expected_publication_epoch,
tier_registry_generation,
)
.await
.ok()
!= Some(expected)
{
return Err(RemoteScannerServerError::retry_bucket(
"remote scanner checkpoint identity changed during scanning",
));
}
if guard.is_lock_lost() {
return Err(RemoteScannerServerError::worker(
+40 -4
View File
@@ -1741,7 +1741,13 @@ impl FolderScanner {
source: FolderScanSource::Existing,
}));
let has_queued_folders = !queued_folders.is_empty();
let resume_order = order_queued_folders_for_resume(&mut queued_folders, scan_resume_after);
let forward_sweep = self.old_cache.info.scan_progress.is_some();
let forward_resume_after = forward_sweep.then(|| scan_resume_after.map(str::to_owned)).flatten();
let resume_order = if forward_sweep {
order_queued_folders_for_resume(&mut queued_folders, None)
} else {
order_queued_folders_for_resume(&mut queued_folders, scan_resume_after)
};
if checkpoint_tracks_child_order && has_queued_folders {
match resume_order {
FolderResumeOrder::Used => global_metrics().record_scanner_checkpoint_used(),
@@ -1758,6 +1764,18 @@ impl FolderScanner {
let mut folder_item = queued_folder.folder;
let h = hash_path(&folder_item.name);
if forward_sweep
&& !into.compacted
&& forward_resume_after.as_deref().is_some_and(|resume| {
folder_item.name.as_str() <= resume
&& !matches!(folder_resume_match(&folder_item.name, resume), Some(FolderResumeMatch::Descendant))
})
&& self.old_cache.find(&folder_item.name).is_some()
{
self.new_cache.copy_with_children(&self.old_cache, &h, &folder_item.parent);
into.add_child(&h);
continue;
}
match queued_folder.source {
FolderScanSource::New => {
@@ -1789,7 +1807,7 @@ impl FolderScanner {
}
}
FolderScanSource::Existing => {
if !into.compacted && self.old_cache.is_compacted(&h) {
if !forward_sweep && !into.compacted && self.old_cache.is_compacted(&h) {
let next_cycle = self.old_cache.info.next_cycle as u32;
if !h.mod_(next_cycle, data_usage_update_dir_cycles()) {
// Transfer and add as child...
@@ -2250,7 +2268,10 @@ impl FolderScanner {
self.new_cache.replace_hashed(&this_hash, &folder.parent, into);
}
// Keep independently accounted children while the sweep cursor may
// reference them; the hard cardinality compaction below still applies.
if !into.compacted
&& self.old_cache.info.scan_progress.is_none()
&& self.new_cache.info.name != folder.name
&& let Some(mut flat) = self.new_cache.size_recursive(&this_hash.key())
{
@@ -2425,7 +2446,22 @@ pub async fn scan_data_folder(
let unresolved_objects = root.failed_objects > 0
|| !new_cache.info.failed_objects.is_empty()
|| !new_cache.info.size_reconciliation.is_empty();
new_cache.info.snapshot_complete = !unresolved_objects;
let mixed_coverage = new_cache
.info
.scan_progress
.is_some_and(|progress| progress.started_plan != progress.requested_plan);
new_cache.info.snapshot_complete = !unresolved_objects && !mixed_coverage;
if let Some(progress) = &mut new_cache.info.scan_progress {
if new_cache.info.snapshot_complete {
new_cache.info.scan_plan_digest = Some(progress.requested_plan);
new_cache.info.scan_progress = None;
} else {
// Retain observations, then verify from the beginning under
// the latest plan. A clean tail cannot certify an old prefix.
progress.started_plan = progress.requested_plan;
new_cache.info.scan_plan_digest = None;
}
}
let had_scan_checkpoint = cache.info.scan_checkpoint.is_some() || new_cache.info.scan_checkpoint.is_some();
new_cache.info.scan_resume_after = None;
new_cache.info.scan_checkpoint = None;
@@ -2434,7 +2470,7 @@ pub async fn scan_data_folder(
}
close_disk_guard.close().await;
if unresolved_objects {
if unresolved_objects || mixed_coverage {
Err(ScannerError::PartialCache(Box::new(new_cache.clone())))
} else {
Ok(new_cache.clone())
@@ -234,6 +234,156 @@ fn checkpoint_fixture_compaction_preserves_aggregate_not_child_enumeration() {
);
}
fn bound_checkpoint() -> (DataUsageCache, crate::DataUsageScanIdentity) {
let identity = crate::DataUsageScanIdentity {
version: 1,
bucket_incarnation: Uuid::from_u128(1),
set_layout: DataUsageScanPlanDigest([41; 32]),
publication_epoch: 0,
tier_registry_generation: 7,
};
let mut cache = DataUsageCache::default();
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
cache.replace("bucket", "", DataUsageEntry::default());
cache.replace(
"bucket/static",
"bucket",
DataUsageEntry {
objects: 3,
..Default::default()
},
);
cache.info.scan_resume_after = Some("bucket/static".into());
cache.info.scan_checkpoint = Some(DataUsageScanCheckpoint::new(
"bucket/static".into(),
DataUsageScanCheckpointReason::Objects,
));
(cache, identity)
}
#[test]
fn checkpoint_fixture_roundtrip_retains_verified_scope_but_old_reader_rebuilds() {
let (cache, identity) = bound_checkpoint();
let mut cache = decode_fixture(&cache.marshal_msg().expect("encode bound progress")).expect("read bound progress");
let next_plan = DataUsageScanPlanDigest([42; 32]);
assert_eq!(
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, next_plan, identity),
crate::DataUsageCachePrepareOutcome::Reused
);
assert_eq!(retained(&cache), 3);
assert_eq!(cache.info.scan_identity, Some(identity));
assert_eq!(
cache.info.scan_progress,
Some(crate::DataUsageScanProgress {
started_plan: PLAN,
requested_plan: next_plan
})
);
assert!(cache.info.scan_plan_digest.is_none());
let mut old_wire = serde_json::to_value(&cache).expect("map-encoded compatibility fixture");
let old_info = old_wire["info"].as_object_mut().expect("cache info is a map");
old_info.remove("scan_identity");
old_info.remove("scan_progress");
let mut old_view: DataUsageCache = serde_json::from_value(old_wire).expect("old writer drops unknown metadata");
assert_eq!(
old_view.prepare_for_scan("bucket", 11, 7, SOURCE, next_plan, true),
crate::DataUsageCachePrepareOutcome::Reset
);
assert!(old_view.cache.is_empty());
assert!(!old_view.info.snapshot_complete);
}
#[test]
fn checkpoint_fixture_identity_changes_and_future_state_fail_closed() {
let (cache, identity) = bound_checkpoint();
for next_identity in [
crate::DataUsageScanIdentity {
bucket_incarnation: Uuid::from_u128(2),
..identity
},
crate::DataUsageScanIdentity {
set_layout: DataUsageScanPlanDigest([9; 32]),
..identity
},
crate::DataUsageScanIdentity {
publication_epoch: 1,
..identity
},
crate::DataUsageScanIdentity {
tier_registry_generation: 8,
..identity
},
] {
let mut next = cache.clone();
assert_eq!(
next.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, next_identity),
crate::DataUsageCachePrepareOutcome::Reset
);
assert!(next.cache.is_empty());
assert!(next.info.scan_checkpoint.is_none());
assert!(!next.info.snapshot_complete);
}
for (source, epoch) in [(crate::DataUsageCacheSource::new(1, 0), 7), (SOURCE, 8)] {
let mut next = cache.clone();
assert_eq!(
next.prepare_bucket_checkpoint("bucket", 11, epoch, source, PLAN, identity),
crate::DataUsageCachePrepareOutcome::Reset
);
assert!(next.cache.is_empty());
}
for (cycle, epoch, expected) in [
(10, 7, crate::DataUsageCachePrepareOutcome::RejectedNewerCycle),
(11, 6, crate::DataUsageCachePrepareOutcome::RejectedNewerLeader),
] {
let mut next = cache.clone();
assert_eq!(next.prepare_bucket_checkpoint("bucket", cycle, epoch, SOURCE, PLAN, identity), expected);
assert_eq!(
serde_json::to_value(&next).expect("current cache"),
serde_json::to_value(&cache).expect("saved cache")
);
}
for invalid in [
crate::DataUsageScanIdentity { version: 2, ..identity },
crate::DataUsageScanIdentity {
bucket_incarnation: Uuid::nil(),
..identity
},
] {
let mut next = cache.clone();
crate::scanner_io::current_cache_root_or_prepare_with_generation(
&mut next,
"bucket",
SOURCE,
11,
7,
PLAN,
crate::scanner_io::DataUsageCacheReuseOptions {
checkpoint_identity: Some(invalid),
..Default::default()
},
);
assert!(next.cache.is_empty(), "unsupported identity must not retain coverage");
}
}
#[test]
fn checkpoint_fixture_corrupt_cursor_restarts_validation_without_claiming_completion() {
let (cache, identity) = bound_checkpoint();
for resume in ["other/static", "bucket/missing"] {
let mut next = cache.clone();
next.info.scan_resume_after = Some(resume.into());
next.info.scan_checkpoint = Some(DataUsageScanCheckpoint::new(resume.into(), DataUsageScanCheckpointReason::Objects));
let plan = DataUsageScanPlanDigest([43; 32]);
next.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, plan, identity);
assert_eq!(retained(&next), 3, "observations may survive an invalid cursor");
assert!(next.info.scan_resume_after.is_none());
assert!(next.info.scan_checkpoint.is_none());
assert_eq!(next.info.scan_progress.expect("new verification sweep").started_plan, plan);
assert!(!next.info.snapshot_complete);
assert!(next.info.scan_plan_digest.is_none());
}
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_save_reload_resume() {
@@ -242,23 +392,44 @@ async fn checkpoint_fixture_save_reload_resume() {
#[tokio::test]
#[serial]
async fn checkpoint_fixture_hot_digest_diagnostic() {
async fn checkpoint_fixture_hot_digest_retains_partial_progress() {
run_checkpoint_fixture(true).await;
}
async fn write_checkpoint_object(root: &std::path::Path, object: &str, versions: &[(Option<Uuid>, i64)]) {
let mut metadata = FileMeta::new();
for (index, (version_id, size)) in versions.iter().enumerate() {
let mut info = FileInfo::new(object, 4, 2);
info.volume = "bucket".into();
info.version_id = *version_id;
info.versioned = version_id.is_some();
info.size = *size;
info.mod_time = Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(i64::try_from(index).expect("fixture index")));
metadata.add_version(info).expect("fixture version");
}
write_test_object_metadata_bytes(root, "bucket", object, &metadata.marshal_msg().expect("fixture metadata")).await;
}
async fn run_checkpoint_fixture(change_digest: bool) {
let (scanner, root) = build_test_scanner().await;
let _guard = TestGuard {
temp_dir: Some(root.clone()),
};
for index in 0..STATIC_OBJECTS {
write_test_object_metadata(&root, "bucket", &format!("static/{index:04}")).await;
write_checkpoint_object(&root, &format!("static/{index:04}"), &[(None, 1)]).await;
}
let identity = crate::DataUsageScanIdentity {
version: 1,
bucket_incarnation: Uuid::from_u128(1),
set_layout: DataUsageScanPlanDigest([41; 32]),
publication_epoch: 0,
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
};
let store = FixtureStore::new();
let mut previous = 0;
let mut visited = 0;
for round in 0..3_u8 {
write_test_object_metadata(&root, "bucket", "hot/current").await;
write_checkpoint_object(&root, "hot/current", &[(None, 1)]).await;
let mut cache = DataUsageCache::default();
let revisions = cache
.load_with_revisions(store.clone(), CACHE_NAME)
@@ -278,6 +449,7 @@ async fn run_checkpoint_fixture(change_digest: bool) {
crate::scanner_io::DataUsageCacheReuseOptions {
require_source: true,
tier_registry_generation: None,
checkpoint_identity: Some(identity),
},
);
let prepared = retained(&cache);
@@ -326,13 +498,12 @@ async fn run_checkpoint_fixture(change_digest: bool) {
eprintln!(
"checkpoint_fixture round={round} hot_digest={change_digest} visited_total={visited} before={previous} prepared={prepared} scanned={scanned} reloaded={reloaded} diagnosis={diagnosis:?}"
);
if !change_digest || std::env::var_os("RUSTFS_CHECKPOINT_REQUIRE_PROGRESS").is_some() {
assert_eq!(
diagnosis,
CoverageDiagnosis::Progress,
"visited growth must produce durable static coverage"
);
}
assert_eq!(
diagnosis,
CoverageDiagnosis::Progress,
"visited growth must produce durable static coverage"
);
assert!(loaded.info.scan_plan_digest.is_none(), "old readers must rebuild an uncertified sweep");
crate::remote_scanner::checkpoint_fixture_partial_return(budget.progress(), budget.entries_visited()).await;
previous = reloaded;
}
@@ -383,28 +554,84 @@ async fn run_checkpoint_fixture(change_digest: bool) {
assert!(result.is_err(), "pre-scan cancellation must not produce a complete root");
assert_eq!(budget.reason(), None, "parent cancellation is not object budget exhaustion");
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new(&parent, Default::default());
let result = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget,
vec![scanner.local_disk.clone()],
loaded,
None,
HealScanMode::Normal,
)
store.reject_save.store(false, Ordering::SeqCst);
write_checkpoint_object(&root, "static/0000", &[(Some(Uuid::from_u128(2)), 7), (Some(Uuid::from_u128(3)), 3)]).await;
tokio::fs::remove_dir_all(root.join("bucket/static/0001"))
.await
.expect("unbounded scan must complete after durable partial progress");
let ScannerDiskScanOutcome::Complete(cache) = result else {
panic!("unbounded fixture must produce a complete disk cache");
};
assert!(cache.info.snapshot_complete);
assert!(cache.info.scan_checkpoint.is_none());
assert_eq!(
cache.checked_flatten("bucket").expect("complete bucket root").objects,
usize::try_from(STATIC_OBJECTS + 1).expect("fixture object count fits usize")
);
.expect("remove previously scanned fixture object");
write_checkpoint_object(&root, "hot/later", &[(None, 1)]).await;
let final_plan = crate::scanner_io::checkpoint_fixture_bucket_digest(PLAN, Some(3));
let mut saw_mixed_sweep_end = false;
for _ in 0..32 {
let mut cache = DataUsageCache::default();
let revisions = cache
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("reload every bounded round");
crate::scanner_io::current_cache_root_or_prepare_with_generation(
&mut cache,
"bucket",
SOURCE,
11,
7,
final_plan,
crate::scanner_io::DataUsageCacheReuseOptions {
require_source: true,
tier_registry_generation: Some(identity.tier_registry_generation),
checkpoint_identity: Some(identity),
},
);
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(
&parent,
ScannerCycleBudgetConfig {
max_objects: Some(4),
..Default::default()
},
);
let outcome = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget.clone(),
vec![scanner.local_disk.clone()],
cache,
None,
HealScanMode::Normal,
)
.await
.expect("bounded sweep outcome");
let (cache, complete) = match outcome {
ScannerDiskScanOutcome::Complete(cache) => (cache, true),
ScannerDiskScanOutcome::Partial(cache) => (cache, false),
ScannerDiskScanOutcome::NamespaceNotFound(_) => panic!("fixture namespace exists"),
};
cache
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect("save each bounded sweep");
let saved = store.strict_load().await;
if complete {
assert!(
saw_mixed_sweep_end,
"a clean tail must first finish as partial before a new validation sweep"
);
assert!(saved.info.snapshot_complete);
assert!(saved.info.scan_progress.is_none());
assert!(saved.info.scan_checkpoint.is_none());
assert_eq!(saved.info.scan_plan_digest, Some(final_plan));
let total = saved.checked_flatten("bucket").expect("complete bucket root");
assert_eq!((total.objects, total.versions, total.size), (25, 2, 34));
assert_eq!(saved.checked_flatten("bucket/static").expect("static subtree").objects, 23);
assert_eq!(saved.checked_flatten("bucket/hot").expect("hot subtree").objects, 2);
return;
}
assert!(!saved.info.snapshot_complete);
assert!(saved.info.scan_plan_digest.is_none());
if !budget.budget_elapsed() {
saw_mixed_sweep_end = true;
}
}
panic!("finite stable fixture must converge using the same four-object budget without an unbounded final sweep");
}
+33
View File
@@ -271,6 +271,8 @@ pub struct ScannerBucketScanPlan {
all_buckets: Arc<Vec<BucketInfo>>,
scope: ScannerBucketScanScope,
digest: DataUsageScanPlanDigest,
/// Includes mutation generations even when the set planner uses a structural digest.
bucket_coverage_digest: DataUsageScanPlanDigest,
leader_epoch: u64,
tier_registry_generation: u64,
/// Epoch captured once for the whole scanner cycle. `None` is retained
@@ -730,6 +732,37 @@ pub(crate) async fn scanner_set_disk_inventory(set: &SetDisks) -> Vec<Arc<Disk>>
disks
}
pub(crate) async fn scanner_bucket_checkpoint_identity(
set: &SetDisks,
bucket: &str,
publication_epoch: u64,
tier_registry_generation: u64,
) -> Result<crate::DataUsageScanIdentity> {
let bucket_incarnation = set.bucket_incarnation_id_from_disk(bucket).await?;
let disks = set
.format
.erasure
.sets
.get(set.set_index)
.filter(|disks| !disks.is_empty())
.ok_or_else(|| Error::other("scanner checkpoint set layout is absent"))?;
if set.format.id.is_nil() || disks.iter().any(uuid::Uuid::is_nil) {
return Err(Error::other("scanner checkpoint set layout has a nil identity"));
}
let mut digest = Sha256::new();
digest.update(set.format.id.as_bytes());
for disk in disks {
digest.update(disk.as_bytes());
}
Ok(crate::DataUsageScanIdentity {
version: 1,
bucket_incarnation,
set_layout: crate::DataUsageScanPlanDigest(digest.finalize().into()),
publication_epoch,
tier_registry_generation,
})
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum ScannerCycleDeferReason {
ActivityBaselineUnavailable,
+15 -1
View File
@@ -112,6 +112,7 @@ pub(crate) fn current_cache_root_entry_with_generation(
let metadata_is_current = cache.info.name == name
&& cache.info.source == Some(source)
&& cache.info.snapshot_complete
&& cache.info.scan_progress.is_none()
&& cache.info.scan_plan_digest == Some(scan_plan_digest)
&& cache.info.last_update.is_some()
&& cache.info.next_cycle == next_cycle
@@ -137,6 +138,7 @@ pub(crate) enum DataUsageCacheScanState {
pub(crate) struct DataUsageCacheReuseOptions {
pub(crate) require_source: bool,
pub(crate) tier_registry_generation: Option<u64>,
pub(crate) checkpoint_identity: Option<crate::DataUsageScanIdentity>,
}
#[cfg(test)]
@@ -159,6 +161,7 @@ pub(crate) fn current_cache_root_or_prepare(
DataUsageCacheReuseOptions {
require_source,
tier_registry_generation: None,
checkpoint_identity: None,
},
)
}
@@ -172,6 +175,12 @@ pub(crate) fn current_cache_root_or_prepare_with_generation(
scan_plan_digest: DataUsageScanPlanDigest,
options: DataUsageCacheReuseOptions,
) -> DataUsageCacheScanState {
if cache.info.next_cycle <= next_cycle
&& cache.info.leader_epoch <= leader_epoch
&& cache.info.scan_identity != options.checkpoint_identity
{
cache.info.scan_plan_digest = None;
}
if options.tier_registry_generation.is_some_and(|generation| {
cache.info.next_cycle <= next_cycle
&& cache.info.leader_epoch <= leader_epoch
@@ -193,7 +202,12 @@ pub(crate) fn current_cache_root_or_prepare_with_generation(
Ok(Some(root)) => DataUsageCacheScanState::Current(Box::new(root)),
current => DataUsageCacheScanState::Prepared {
invalid_current: current.err(),
outcome: cache.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, options.require_source),
outcome: match options.checkpoint_identity.filter(crate::DataUsageScanIdentity::is_valid) {
Some(identity) => {
cache.prepare_bucket_checkpoint(name, next_cycle, leader_epoch, source, scan_plan_digest, identity)
}
None => cache.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, options.require_source),
},
},
}
}
+27 -1
View File
@@ -115,6 +115,7 @@ impl ScannerIOCache for SetDisks {
all_buckets,
scope,
digest: scan_plan_digest,
bucket_coverage_digest,
leader_epoch,
tier_registry_generation,
publication_epoch,
@@ -634,7 +635,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(scan_plan_digest, dirty_usage_buckets_clone.get(&bucket.name).copied());
scanner_bucket_cache_digest(bucket_coverage_digest, dirty_usage_buckets_clone.get(&bucket.name).copied());
if let Some(server_epoch) = remote_server_epoch {
let request_sequence = remote_session_sequence;
@@ -880,6 +881,16 @@ impl ScannerIOCache for SetDisks {
continue;
}
};
// Lack of an authoritative legacy identity disables the new
// checkpoint protocol; the existing full rebuild remains available.
let checkpoint_identity = scanner_bucket_checkpoint_identity(
&store_clone_clone,
&bucket.name,
expected_publication_epoch_clone,
tier_registry_generation,
)
.await
.ok();
let scan_state = current_cache_root_or_prepare_with_generation(
&mut cache,
&bucket.name,
@@ -890,6 +901,7 @@ impl ScannerIOCache for SetDisks {
DataUsageCacheReuseOptions {
require_source: require_cache_source,
tier_registry_generation: Some(tier_registry_generation),
checkpoint_identity,
},
);
let outcome = match scan_state {
@@ -1039,6 +1051,20 @@ impl ScannerIOCache for SetDisks {
}
}
};
if let Some(expected) = checkpoint_identity
&& scanner_bucket_checkpoint_identity(
&store_clone_clone,
&bucket.name,
expected_publication_epoch_clone,
tier_registry_generation,
)
.await
.ok()
!= Some(expected)
{
record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await;
continue;
}
let scan_outcome = match scan_result {
Ok(scan_outcome) => scan_outcome,
Err(e) => {
@@ -263,6 +263,8 @@ where
bucket_plan_complete &= scanner_bucket_inventory_is_complete(&all_buckets, &buckets_by_source);
let scan_plan_digest =
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before));
let bucket_coverage_digest =
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_snapshot_digest(&activity_before));
let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list));
let scan_scope = resolve_scanner_bucket_scan_scope(
store,
@@ -411,6 +413,7 @@ where
all_buckets: Arc::clone(&all_buckets),
scope: scan_scope.clone(),
digest: scan_plan_digest,
bucket_coverage_digest,
leader_epoch,
tier_registry_generation,
publication_epoch,
@@ -771,6 +771,7 @@ fn current_cache_root_with_new_tier_generation_resets_old_cache() {
DataUsageCacheReuseOptions {
require_source: false,
tier_registry_generation: Some(2),
checkpoint_identity: None,
},
);
+33
View File
@@ -123,6 +123,39 @@ async fn setup_two_pool_scanner_store() -> (tempfile::TempDir, Arc<ECStore>) {
(temp_dir, store)
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_bucket_identity_uses_its_set_instance_owner() {
let (_first_dir, first) = setup_two_pool_scanner_store().await;
first
.make_bucket("checkpoint-identity", &MakeBucketOptions::default())
.await
.expect("first instance bucket");
let first_identity = scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "checkpoint-identity", 0, 7)
.await
.expect("first durable identity");
let (_second_dir, second) = setup_two_pool_scanner_store().await;
second
.make_bucket("checkpoint-identity", &MakeBucketOptions::default())
.await
.expect("second instance bucket");
let second_identity = scanner_bucket_checkpoint_identity(&second.pools[0].disk_set[0], "checkpoint-identity", 0, 7)
.await
.expect("second durable identity");
assert_ne!(first_identity.bucket_incarnation, second_identity.bucket_incarnation);
assert_eq!(
scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "checkpoint-identity", 0, 7)
.await
.expect("first owner remains bound"),
first_identity
);
assert!(
scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "missing-checkpoint-bucket", 0, 7)
.await
.is_err()
);
}
#[tokio::test]
#[serial]
async fn scanner_cache_locks_block_same_source_workers() {