fix(scanner): verify coverage receipts and scan strength

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-05 19:21:05 +08:00
parent b09a4f3706
commit 489c3276a4
10 changed files with 543 additions and 26 deletions
+135 -13
View File
@@ -14,6 +14,7 @@
use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration};
use serde::{Deserialize, Serialize, ser::SerializeMap};
use sha2::{Digest, Sha256};
use std::{
collections::{HashMap, HashSet},
future::Future,
@@ -409,23 +410,27 @@ pub struct DataUsageScanIdentity {
pub set_layout: DataUsageScanPlanDigest,
pub publication_epoch: u64,
pub tier_registry_generation: u64,
pub scan_mode: HealScanMode,
}
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))?;
let mut map = serializer.serialize_map(Some(6))?;
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.serialize_entry("scan_mode", &self.scan_mode)?;
map.end()
}
}
impl DataUsageScanIdentity {
pub(crate) fn is_valid(&self) -> bool {
self.version == 1 && !self.bucket_incarnation.is_nil()
self.version == 1
&& !self.bucket_incarnation.is_nil()
&& matches!(self.scan_mode, HealScanMode::Normal | HealScanMode::Deep)
}
}
@@ -446,6 +451,35 @@ impl Serialize for DataUsageScanProgress {
}
}
#[derive(Clone, Debug, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct DataUsageScanCoverageReceipt {
pub through: String,
pub digest: [u8; 32],
}
impl Serialize for DataUsageScanCoverageReceipt {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
let mut map = serializer.serialize_map(Some(2))?;
map.serialize_entry("through", &self.through)?;
map.serialize_entry("digest", &self.digest)?;
map.end()
}
}
struct CheckpointDigestWriter(Sha256);
impl std::io::Write for CheckpointDigestWriter {
fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
self.0.update(bytes);
Ok(bytes.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct DataUsageEntryInfo {
pub name: String,
@@ -520,6 +554,8 @@ pub struct DataUsageCacheInfo {
#[serde(default)]
pub scan_progress: Option<DataUsageScanProgress>,
#[serde(default)]
pub scan_coverage_receipt: Option<DataUsageScanCoverageReceipt>,
#[serde(default)]
pub pending_heals: Vec<PendingScannerHeal>,
#[serde(default)]
pub object_lock: Option<Arc<ObjectLockConfiguration>>,
@@ -565,6 +601,7 @@ impl Serialize for DataUsageCacheInfo {
let field_count = 16
+ usize::from(self.scan_identity.is_some())
+ usize::from(self.scan_progress.is_some())
+ usize::from(self.scan_coverage_receipt.is_some())
+ usize::from(self.tier_registry_generation.is_some())
+ usize::from(!self.size_reconciliation.is_empty())
+ usize::from(self.lkg_snapshot_complete)
@@ -589,6 +626,9 @@ impl Serialize for DataUsageCacheInfo {
if let Some(progress) = self.scan_progress {
state.serialize_entry("scan_progress", &progress)?;
}
if let Some(receipt) = &self.scan_coverage_receipt {
state.serialize_entry("scan_coverage_receipt", receipt)?;
}
state.serialize_entry("pending_heals", &self.pending_heals)?;
state.serialize_entry("object_lock", &self.object_lock)?;
state.serialize_entry("source", &self.source)?;
@@ -777,7 +817,14 @@ impl DataUsageCache {
&& 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 && self.info.scan_progress.is_none() && self.info.scan_plan_digest == Some(scan_plan_digest) {
if reusable
&& self.info.snapshot_complete
&& self.info.scan_progress.is_none()
&& self.info.scan_checkpoint.is_none()
&& self.info.scan_resume_after.is_none()
&& self.info.scan_coverage_receipt.is_none()
&& self.info.scan_plan_digest == Some(scan_plan_digest)
{
return self.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, true);
}
if !reusable {
@@ -798,16 +845,10 @@ impl DataUsageCache {
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,
};
let cursor_is_valid = (self.info.scan_checkpoint.is_none()
&& self.info.scan_resume_after.is_none()
&& self.info.scan_coverage_receipt.is_none())
|| self.validated_scan_frontier().is_some();
if !cursor_is_valid {
self.info.scan_progress = None;
}
@@ -828,6 +869,7 @@ impl DataUsageCache {
});
self.info.scan_resume_after = None;
self.info.scan_checkpoint = None;
self.info.scan_coverage_receipt = None;
}
// Old readers do not understand coverage sweeps. An absent plan makes
// their existing prepare path rebuild instead of promoting mixed data.
@@ -839,6 +881,86 @@ impl DataUsageCache {
}
}
fn coverage_prefix_digest(&self, through: &str) -> Result<[u8; 32], serde_json::Error> {
let mut writer = CheckpointDigestWriter(Sha256::new());
serde_json::to_writer(
&mut writer,
&(
&self.info.name,
self.info.scan_identity,
self.info.source,
self.info.leader_epoch,
self.info.cache_key_format,
self.info.scan_progress.map(|progress| progress.started_plan),
through,
),
)?;
let mut prefix = self
.cache
.iter()
.filter(|(key, _)| {
let ancestor = through
.strip_prefix(key.as_str())
.is_some_and(|suffix| suffix.starts_with('/'));
let descendant = key.strip_prefix(through).is_some_and(|suffix| suffix.starts_with('/'));
(key.as_str() <= through && !ancestor) || descendant
})
.collect::<Vec<_>>();
prefix.sort_unstable_by(|(left, _), (right, _)| left.cmp(right));
for (key, entry) in prefix {
let mut value = serde_json::to_value(entry)?;
value.sort_all_objects();
if let Some(children) = value.get_mut("children").and_then(serde_json::Value::as_array_mut) {
children.sort_unstable_by(|left, right| left.as_str().cmp(&right.as_str()));
}
serde_json::to_writer(&mut writer, &(key, value))?;
}
Ok(writer.0.finalize().into())
}
pub(crate) fn validated_scan_frontier(&self) -> Option<&str> {
let receipt = self.info.scan_coverage_receipt.as_ref()?;
let checkpoint = self.info.scan_checkpoint.as_ref()?;
(self.info.scan_progress.is_some()
&& self.info.scan_identity.is_some_and(|identity| identity.is_valid())
&& self.info.source.is_some()
&& receipt.through.len() <= 16 * 1024
&& checkpoint.version == DATA_USAGE_SCAN_CHECKPOINT_VERSION
&& checkpoint.resume_after == receipt.through
&& self.info.scan_resume_after.as_deref() == Some(receipt.through.as_str())
&& receipt
.through
.strip_prefix(&self.info.name)
.is_some_and(|suffix| suffix.starts_with('/'))
&& self.find(&receipt.through).is_some()
&& self.coverage_prefix_digest(&receipt.through).ok() == Some(receipt.digest))
.then_some(receipt.through.as_str())
}
/// Seal only the frontier supplied by completed traversal, never a restored cursor.
pub(crate) fn seal_scan_frontier(&mut self, frontier: Option<&str>) -> Result<(), serde_json::Error> {
if self.info.scan_progress.is_none() {
self.info.scan_coverage_receipt = None;
return Ok(());
}
let frontier = frontier.filter(|path| path.len() <= 16 * 1024 && self.find(path).is_some());
self.info.scan_coverage_receipt = match frontier {
Some(through) => Some(DataUsageScanCoverageReceipt {
through: through.to_owned(),
digest: self.coverage_prefix_digest(through)?,
}),
None => None,
};
self.info.scan_resume_after = frontier.map(str::to_owned);
let reason = self
.info
.scan_checkpoint
.as_ref()
.map_or(DataUsageScanCheckpointReason::Unknown, |checkpoint| checkpoint.reason);
self.info.scan_checkpoint = frontier.map(|through| DataUsageScanCheckpoint::new(through.to_owned(), reason));
Ok(())
}
fn ensure_cache_save_metrics_registered() {
CACHE_SAVE_METRICS_ONCE.call_once(|| {
describe_counter!(
@@ -733,6 +733,7 @@ async fn scan_and_persist_local_bucket(
&bucket,
expected_publication_epoch,
tier_registry_generation,
scan_mode,
)
.await
.ok(),
@@ -835,6 +836,7 @@ async fn scan_and_persist_local_bucket(
&bucket,
expected_publication_epoch,
tier_registry_generation,
scan_mode,
)
.await
.ok()
+44 -4
View File
@@ -707,6 +707,9 @@ pub struct FolderScanner {
/// next scan and cannot mix generations in one aggregate.
tier_registry: TierRegistrySnapshot,
pending_heals_changed: bool,
coverage_frontier: Option<String>,
resume_frontier: Option<String>,
coverage_gap: bool,
pending_size_reconciliation_keys: HashSet<String>,
pending_size_reconciliation_scopes: HashSet<String>,
pending_size_reconciliation_truncated: bool,
@@ -906,6 +909,16 @@ impl FolderScanner {
self.update_cache.info.scan_checkpoint = Some(checkpoint);
}
fn record_completed_child(&mut self, folder: &str, healthy: bool) {
if self.old_cache.info.scan_progress.is_some() {
self.coverage_gap |= !healthy;
if !self.coverage_gap {
self.coverage_frontier = Some(folder.to_owned());
}
}
self.record_scan_resume_hint(folder);
}
fn record_scan_resume_hint_if_not_ancestor(&mut self, folder: &str) {
let keep_existing = self
.new_cache
@@ -1471,6 +1484,7 @@ impl FolderScanner {
// (e.g. in the get_size error branch below). This branch only accounts
// for subsequent skips of already-failed paths.
if self.should_skip_failed(&item.path) {
self.coverage_gap |= self.old_cache.info.scan_progress.is_some();
continue;
}
@@ -1484,6 +1498,7 @@ impl FolderScanner {
let failure_action = classify_get_size_failure(&item, &e);
if failure_action != GetSizeFailureAction::Skip {
self.coverage_gap |= self.old_cache.info.scan_progress.is_some();
// Track failed objects to prevent infinite retry loops
into.failed_objects += 1;
self.record_failed(&item.path);
@@ -1576,6 +1591,9 @@ impl FolderScanner {
abandoned_children.remove(&path_join_buf(&[&item.bucket, &item.object_path()]));
apply_scanner_size_summary(into, &sz);
if !sz.size_reconciliation.is_empty() {
self.coverage_gap |= self.old_cache.info.scan_progress.is_some();
}
self.apply_size_reconciliation(&sz);
into.objects += 1;
object_count += 1;
@@ -1620,6 +1638,7 @@ impl FolderScanner {
}
if self.is_erasure_mode && found_erasure_data_directory && !found_object_metadata {
self.coverage_gap |= self.old_cache.info.scan_progress.is_some();
found_object_metadata = true;
let metadata_path = path_join_buf(&[&dir_path, STORAGE_FORMAT_FILE]);
@@ -1742,7 +1761,7 @@ impl FolderScanner {
}));
let has_queued_folders = !queued_folders.is_empty();
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 forward_resume_after = self.resume_frontier.clone();
let resume_order = if forward_sweep {
order_queued_folders_for_resume(&mut queued_folders, None)
} else {
@@ -1828,7 +1847,7 @@ impl FolderScanner {
// In compacted mode child totals are accumulated directly into the parent entry.
let fut = Box::pin(self.scan_folder(ctx.clone(), folder_item.clone(), into));
fut.await.map_err(|e| ScannerError::Other(e.to_string()))?;
self.record_scan_resume_hint(&folder_item.name);
self.record_completed_child(&folder_item.name, into.failed_objects == 0);
self.send_update_for_entry(&this_hash, &folder.parent, into).await;
tokio::task::yield_now().await;
} else {
@@ -1853,12 +1872,13 @@ impl FolderScanner {
error = %e,
"Scanner child folder scan failed"
);
self.coverage_gap |= forward_sweep;
continue;
}
tokio::task::yield_now().await;
into.add_child(&h);
self.record_scan_resume_hint(&folder_item.name);
self.record_completed_child(&folder_item.name, dst.failed_objects == 0);
// We scanned a folder, optionally send update.
self.update_cache.delete_recursive(&h);
self.update_cache.copy_with_children(&self.new_cache, &h, &folder_item.parent);
@@ -2380,6 +2400,7 @@ pub async fn scan_data_folder(
cache.fold_retired_tiers(&tier_registry.names);
cache.info.tier_registry_generation = Some(tier_registry.generation);
let resume_frontier = cache.validated_scan_frontier().map(str::to_owned);
// Create folder scanner
let mut scanner = FolderScanner {
root: base_path,
@@ -2409,6 +2430,9 @@ pub async fn scan_data_folder(
local_disk,
tier_registry,
pending_heals_changed: false,
coverage_frontier: resume_frontier.clone(),
resume_frontier,
coverage_gap: false,
pending_size_reconciliation_keys: HashSet::new(),
pending_size_reconciliation_scopes: HashSet::new(),
pending_size_reconciliation_truncated: false,
@@ -2439,11 +2463,13 @@ pub async fn scan_data_folder(
match scanner.scan_folder(ctx.clone(), folder, &mut root).await {
Ok(()) => {
// Get the new cache and finalize it
let coverage_gap = scanner.coverage_gap;
let new_cache = scanner.as_mut_new_cache();
new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN);
new_cache.info.last_update = Some(SystemTime::now());
new_cache.info.next_cycle = cache.info.next_cycle;
let unresolved_objects = root.failed_objects > 0
let unresolved_objects = coverage_gap
|| root.failed_objects > 0
|| !new_cache.info.failed_objects.is_empty()
|| !new_cache.info.size_reconciliation.is_empty();
let mixed_coverage = new_cache
@@ -2465,6 +2491,7 @@ pub async fn scan_data_folder(
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;
new_cache.info.scan_coverage_receipt = None;
if had_scan_checkpoint {
global_metrics().record_scanner_checkpoint_cleared();
}
@@ -2484,6 +2511,7 @@ pub async fn scan_data_folder(
if root_has_progress {
scanner.carry_forward_old_children(&root_hash, &mut root);
}
let coverage_frontier = scanner.coverage_frontier.clone();
let new_cache = scanner.as_mut_new_cache();
if root_has_progress {
new_cache.replace_hashed(&root_hash, &None, &root);
@@ -2498,6 +2526,18 @@ pub async fn scan_data_folder(
if root_has_progress {
set_scan_checkpoint(new_cache, checkpoint_reason_from_budget(budget.reason()));
}
new_cache.seal_scan_frontier(coverage_frontier.as_deref())?;
if new_cache.info.scan_progress.is_some() {
if let Some(checkpoint) = &new_cache.info.scan_checkpoint {
global_metrics().record_scanner_checkpoint_set(
checkpoint.version,
checkpoint.resume_after.clone(),
checkpoint.reason.as_str(),
);
} else {
global_metrics().record_scanner_checkpoint_cleared();
}
}
close_disk_guard.close().await;
return Err(ScannerError::PartialCache(Box::new(new_cache.clone())));
}
@@ -346,6 +346,9 @@ async fn build_test_scanner() -> (FolderScanner, std::path::PathBuf) {
refresh_failed: false,
},
pending_heals_changed: false,
coverage_frontier: None,
resume_frontier: None,
coverage_gap: false,
pending_size_reconciliation_keys: HashSet::new(),
pending_size_reconciliation_scopes: HashSet::new(),
pending_size_reconciliation_truncated: false,
@@ -241,6 +241,7 @@ fn bound_checkpoint() -> (DataUsageCache, crate::DataUsageScanIdentity) {
set_layout: DataUsageScanPlanDigest([41; 32]),
publication_epoch: 0,
tier_registry_generation: 7,
scan_mode: HealScanMode::Normal,
};
let mut cache = DataUsageCache::default();
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
@@ -258,6 +259,9 @@ fn bound_checkpoint() -> (DataUsageCache, crate::DataUsageScanIdentity) {
"bucket/static".into(),
DataUsageScanCheckpointReason::Objects,
));
cache
.seal_scan_frontier(Some("bucket/static"))
.expect("completed fixture prefix receipt");
(cache, identity)
}
@@ -284,6 +288,7 @@ fn checkpoint_fixture_roundtrip_retains_verified_scope_but_old_reader_rebuilds()
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");
old_info.remove("scan_coverage_receipt");
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),
@@ -300,6 +305,7 @@ fn checkpoint_fixture_unchanged_complete_plan_keeps_existing_rescan_policy() {
cache.info.scan_plan_digest = Some(PLAN);
cache.info.scan_resume_after = None;
cache.info.scan_checkpoint = None;
cache.info.scan_coverage_receipt = None;
cache.info.snapshot_complete = true;
assert_eq!(
cache.prepare_bucket_checkpoint("bucket", 12, 7, SOURCE, PLAN, identity),
@@ -338,6 +344,10 @@ fn checkpoint_fixture_identity_changes_and_future_state_fail_closed() {
tier_registry_generation: 8,
..identity
},
crate::DataUsageScanIdentity {
scan_mode: HealScanMode::Deep,
..identity
},
] {
let mut next = cache.clone();
assert_eq!(
@@ -409,6 +419,332 @@ fn checkpoint_fixture_corrupt_cursor_restarts_validation_without_claiming_comple
}
}
#[test]
fn checkpoint_fixture_receipt_binds_covered_prefix_not_unvisited_suffix() {
let (mut cache, _) = bound_checkpoint();
cache.replace(
"bucket/z-unvisited",
"bucket",
DataUsageEntry {
objects: 99,
..Default::default()
},
);
assert_eq!(cache.validated_scan_frontier(), Some("bucket/static"));
let saved = decode_fixture(&cache.marshal_msg().expect("persist receipt")).expect("load receipt");
assert_eq!(saved.validated_scan_frontier(), Some("bucket/static"));
cache.cache.get_mut("bucket/static").expect("covered prefix").objects = 100;
assert!(
cache.validated_scan_frontier().is_none(),
"altered covered content must invalidate the receipt"
);
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_existing_uncovered_cursor_cannot_skip_to_complete() {
let (scanner, root) = build_test_scanner().await;
let _guard = TestGuard {
temp_dir: Some(root.clone()),
};
for (object, size) in [
("a-done/object", 1),
("b-pending/first", 7),
("b-pending/second", 1),
("z-stale/object", 2),
] {
write_checkpoint_object(&root, object, &[(None, size)]).await;
}
let identity = crate::DataUsageScanIdentity {
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
..bound_checkpoint().1
};
for tamper_receipt_path in [false, true] {
let store = FixtureStore::new();
let mut cache = DataUsageCache::default();
let revisions = cache
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("empty fixture revisions");
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
cache.replace("bucket", "", DataUsageEntry::default());
for (prefix, objects) in [("a-done", 1), ("b-pending", 99), ("z-stale", 99)] {
cache.replace(
&format!("bucket/{prefix}"),
"bucket",
DataUsageEntry {
objects,
size: objects,
compacted: true,
..Default::default()
},
);
}
cache
.seal_scan_frontier(Some("bucket/a-done"))
.expect("actual completed prefix receipt");
cache.info.scan_resume_after = Some("bucket/z-stale".into());
cache.info.scan_checkpoint = Some(DataUsageScanCheckpoint::new(
"bucket/z-stale".into(),
DataUsageScanCheckpointReason::Objects,
));
if tamper_receipt_path {
cache.info.scan_coverage_receipt.as_mut().expect("receipt").through = "bucket/z-stale".into();
}
cache
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect("persist corrupted existing cursor");
let mut loaded = store.strict_load().await;
let revisions = loaded
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("corrupt cursor CAS revision");
assert!(loaded.validated_scan_frontier().is_none());
loaded.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
assert!(loaded.info.scan_resume_after.is_none());
loaded.info.skip_healing = true;
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(
&parent,
ScannerCycleBudgetConfig {
max_objects: Some(2),
..Default::default()
},
);
let outcome = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget.clone(),
vec![scanner.local_disk.clone()],
loaded,
None,
HealScanMode::Normal,
)
.await
.expect("scan must revisit the prefix");
let ScannerDiskScanOutcome::Partial(cache) = outcome else {
panic!("uncovered suffix must not become complete")
};
assert_eq!(budget.progress().0, 2);
assert_eq!(cache.checked_flatten("bucket/b-pending").expect("revisited prefix").size, 7);
assert!(!cache.info.snapshot_complete);
cache
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect("persist verified partial");
assert!(!store.strict_load().await.info.snapshot_complete);
}
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_failed_child_prevents_receipt_advancing_past_gap() {
let (scanner, root) = build_test_scanner().await;
let _guard = TestGuard {
temp_dir: Some(root.clone()),
};
for object in ["a-good", "b-skipped", "c-later"] {
write_checkpoint_object(&root, object, &[(None, 1)]).await;
}
let identity = crate::DataUsageScanIdentity {
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
..bound_checkpoint().1
};
let mut cache = DataUsageCache::default();
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
cache.info.skip_healing = true;
cache.info.failed_objects.insert(
root.join("bucket/b-skipped/xl.meta").to_string_lossy().into_owned(),
FolderScanner::now_secs(),
);
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(
&parent,
ScannerCycleBudgetConfig {
max_objects: Some(2),
..Default::default()
},
);
let result = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget,
vec![scanner.local_disk.clone()],
cache,
None,
HealScanMode::Normal,
)
.await
.expect("scan with a known failed child");
let ScannerDiskScanOutcome::Partial(cache) = result else {
panic!("skipped failure is not complete")
};
assert_eq!(cache.validated_scan_frontier(), Some("bucket/a-good"));
assert!(!cache.info.failed_objects.is_empty());
assert!(!cache.info.snapshot_complete);
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_complete_sampling_partial_resumes_with_fixed_budget() {
check_complete_sampling_resumption(HealScanMode::Normal).await;
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_normal_partial_reenters_prefix_for_deep_scan() {
check_complete_sampling_resumption(HealScanMode::Deep).await;
}
async fn check_complete_sampling_resumption(resume_mode: HealScanMode) {
let (scanner, root) = build_test_scanner().await;
let _guard = TestGuard {
temp_dir: Some(root.clone()),
};
for index in 0..9 {
write_checkpoint_object(&root, &format!("prefix/{index:04}"), &[(None, 1)]).await;
}
let identity = crate::DataUsageScanIdentity {
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
..bound_checkpoint().1
};
let mut cache = DataUsageCache::default();
cache.prepare_bucket_checkpoint("bucket", 11, 7, SOURCE, PLAN, identity);
cache.info.skip_healing = true;
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new(&parent, Default::default());
// Seed an existing complete baseline; every recovery attempt below is bounded.
let baseline = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget,
vec![scanner.local_disk.clone()],
cache,
None,
HealScanMode::Normal,
)
.await
.expect("initial complete baseline");
let ScannerDiskScanOutcome::Complete(mut cache) = baseline else { panic!("baseline is complete") };
let mut deep_current = cache.clone();
let deep_identity = crate::DataUsageScanIdentity {
scan_mode: HealScanMode::Deep,
..identity
};
let state = crate::scanner_io::current_cache_root_or_prepare_with_generation(
&mut deep_current,
"bucket",
SOURCE,
11,
7,
PLAN,
crate::scanner_io::DataUsageCacheReuseOptions {
checkpoint_identity: Some(deep_identity),
..Default::default()
},
);
assert!(
matches!(state, crate::scanner_io::DataUsageCacheScanState::Prepared { .. }),
"Normal complete is not Deep Current"
);
cache.prepare_bucket_checkpoint("bucket", 12, 7, SOURCE, PLAN, identity);
assert!(cache.info.scan_progress.is_none(), "complete unchanged baseline uses existing sampling");
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(
&parent,
ScannerCycleBudgetConfig {
max_objects: Some(3),
..Default::default()
},
);
let outcome = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget,
vec![scanner.local_disk.clone()],
cache,
None,
HealScanMode::Normal,
)
.await
.expect("sampling interruption");
let ScannerDiskScanOutcome::Partial(cache) = outcome else {
panic!("sampling must exhaust the three-object budget")
};
assert!(cache.info.scan_progress.is_none());
assert_eq!(cache.info.scan_plan_digest, Some(PLAN));
assert!(cache.info.scan_checkpoint.is_some());
let store = FixtureStore::new();
let mut loaded = DataUsageCache::default();
let revisions = loaded
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("fixture revision");
cache
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect("persist sampling partial");
if resume_mode == HealScanMode::Deep {
write_checkpoint_object(&root, "prefix/0000", &[(None, 7)]).await;
}
let resumed_identity = crate::DataUsageScanIdentity {
scan_mode: resume_mode,
..identity
};
for round in 0..16 {
let mut loaded = DataUsageCache::default();
let revisions = loaded
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("reload partial each recovery round");
loaded.prepare_bucket_checkpoint("bucket", 12, 7, SOURCE, PLAN, resumed_identity);
assert!(loaded.info.scan_progress.is_some(), "sampling partial must enter forward validation");
loaded.info.skip_healing = true;
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(
&parent,
ScannerCycleBudgetConfig {
max_objects: Some(3),
..Default::default()
},
);
let result = scanner
.local_disk
.clone()
.nsscanner_disk(budget.token(), budget, vec![scanner.local_disk.clone()], loaded, None, resume_mode)
.await
.expect("bounded recovery scan");
let (cache, complete) = match result {
ScannerDiskScanOutcome::Partial(cache) => (cache, false),
ScannerDiskScanOutcome::Complete(cache) => (cache, true),
_ => panic!("fixture namespace remains present"),
};
if round == 0 && resume_mode == HealScanMode::Deep {
assert_eq!(cache.find("bucket/prefix/0000").expect("Deep revisits earlier prefix").size, 7);
}
cache
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect("save bounded recovery");
if complete {
let root = store.strict_load().await.checked_flatten("bucket").expect("certified root");
assert_eq!(root.objects, 9);
assert_eq!(root.size, if resume_mode == HealScanMode::Deep { 15 } else { 9 });
return;
}
}
panic!("sampling interruption must recover with the unchanged three-object budget");
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_save_reload_resume() {
@@ -449,6 +785,7 @@ async fn run_checkpoint_fixture(change_digest: bool) {
set_layout: DataUsageScanPlanDigest([41; 32]),
publication_epoch: 0,
tier_registry_generation: crate::runtime_tier_registry_for_cycle(11, 7).await.generation,
scan_mode: HealScanMode::Normal,
};
let store = FixtureStore::new();
let mut previous = 0;
+2
View File
@@ -737,6 +737,7 @@ pub(crate) async fn scanner_bucket_checkpoint_identity(
bucket: &str,
publication_epoch: u64,
tier_registry_generation: u64,
scan_mode: HealScanMode,
) -> Result<crate::DataUsageScanIdentity> {
let bucket_incarnation = set.bucket_incarnation_id_from_disk(bucket).await?;
let disks = set
@@ -760,6 +761,7 @@ pub(crate) async fn scanner_bucket_checkpoint_identity(
set_layout: crate::DataUsageScanPlanDigest(digest.finalize().into()),
publication_epoch,
tier_registry_generation,
scan_mode,
})
}
+3
View File
@@ -113,6 +113,9 @@ pub(crate) fn current_cache_root_entry_with_generation(
&& cache.info.source == Some(source)
&& cache.info.snapshot_complete
&& cache.info.scan_progress.is_none()
&& cache.info.scan_checkpoint.is_none()
&& cache.info.scan_resume_after.is_none()
&& cache.info.scan_coverage_receipt.is_none()
&& cache.info.scan_plan_digest == Some(scan_plan_digest)
&& cache.info.last_update.is_some()
&& cache.info.next_cycle == next_cycle
@@ -888,6 +888,7 @@ impl ScannerIOCache for SetDisks {
&bucket.name,
expected_publication_epoch_clone,
tier_registry_generation,
scan_mode,
)
.await
.ok();
@@ -1057,6 +1058,7 @@ impl ScannerIOCache for SetDisks {
&bucket.name,
expected_publication_epoch_clone,
tier_registry_generation,
scan_mode,
)
.await
.ok()
+10 -8
View File
@@ -131,26 +131,28 @@ async fn checkpoint_fixture_bucket_identity_uses_its_set_instance_owner() {
.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 first_identity =
scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "checkpoint-identity", 0, 7, HealScanMode::Normal)
.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");
let second_identity =
scanner_bucket_checkpoint_identity(&second.pools[0].disk_set[0], "checkpoint-identity", 0, 7, HealScanMode::Normal)
.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)
scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "checkpoint-identity", 0, 7, HealScanMode::Normal)
.await
.expect("first owner remains bound"),
first_identity
);
assert!(
scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "missing-checkpoint-bucket", 0, 7)
scanner_bucket_checkpoint_identity(&first.pools[0].disk_set[0], "missing-checkpoint-bucket", 0, 7, HealScanMode::Normal)
.await
.is_err()
);