Compare commits

...

33 Commits

Author SHA1 Message Date
Zhengchao An 982ed2888f Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-06 03:32:00 +08:00
Zhengchao An e1608fbd9c test(odm): exercise overflow and invalid cursors reliably (#7236) 2026-09-06 03:09:39 +08:00
Zhengchao An eb1b17802c test(odm): provide the source region in access fixture (#7235) 2026-09-06 02:34:56 +08:00
houseme eec6d3a76d Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-06 02:29:11 +08:00
Zhengchao An 112f70914d fix(build): scope migration helpers to their features (#7234)
* test(odm): keep listing header import test scoped

* fix(build): gate GCS-only migration HTTP helpers
2026-09-06 02:24:04 +08:00
houseme bbca907781 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-06 01:50:19 +08:00
houseme 51982be687 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-06 01:33:21 +08:00
houseme b8f0119320 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-05 23:53:05 +08:00
houseme 20ca988e85 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-05 23:33:13 +08:00
houseme 045770d09d fix(scanner): satisfy cache prefix sort lint
Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 22:34:20 +08:00
houseme 278cbaefd5 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-05 22:16:44 +08:00
houseme 982bebf6a5 Merge remote-tracking branch 'origin/main' into houseme/fix/scanner-heal-v2-scanner-integration
Resolved scanner cache publication conflicts after the main branch added
execution identity fencing and post-lease activity proof coverage.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 22:10:39 +08:00
houseme 1aafe477ec test(scanner): verify joint checkpoint coverage metadata
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 21:05:13 +08:00
houseme c33a42c5ab test(scanner): supply explicit coverage in publication fixtures
Keep the confirmed-empty namespace fixture authoritative under the required
coverage contract and qualify the bucket cache metadata test type.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:55:19 +08:00
houseme e3fd1ae19a fix(scanner): fence same-cycle caches with full activity coverage
Keep structural baseline identity separate from the full activity coverage
required by bucket admission and set publication. Require complete set
coverage proofs while retaining revision CAS and epoch regression checks.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:55:19 +08:00
houseme 6c1d401215 test(scanner): reproduce same-cycle dirty aggregate replay
Cover a Normal-to-Normal retry with a new dirty bucket generation after
bucket persistence and root delivery failure.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:53:20 +08:00
houseme c80ff445c0 fix(scanner): fence set snapshot reuse with the scan work proof
Prevent same-cycle set publication from replacing freshly scanned maintenance
results with an older Normal aggregate. Recognize uniform completed
maintenance baselines when planning later ordinary dirty-bucket work.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:53:20 +08:00
houseme 3d4c05ddf9 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>
2026-09-05 20:52:15 +08:00
houseme f9dd0388fa fix(scanner): refresh scope safety independently of idle backoff
Inspect maintenance on multi-disk startup and refresh changed or failed
evidence even when explicit bitrot configuration disables idle backoff.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:50:13 +08:00
houseme 527860d71e fix(scanner): keep maintenance cycles outside dirty bucket scopes
Force complete bucket scope for deep scans and scheduled maintenance while
preserving the existing planner for verified ordinary dirty work.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:50:12 +08:00
houseme 7d2c120073 test(scanner): use valid modification times in checkpoint fixtures
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:49:50 +08:00
houseme 489c3276a4 fix(scanner): verify coverage receipts and scan strength
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:49:50 +08:00
houseme b09a4f3706 fix(scanner): keep stable snapshot rescan behavior
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:49:50 +08:00
houseme a0c32e7c6b fix(scanner): retain scoped partial coverage across dirty plans
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:49:50 +08:00
houseme d5a1530409 fix(scanner): require complete publication coverage
Refs rustfs/backlog#2261 and rustfs/backlog#2240.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:48:01 +08:00
houseme e7b0349a4f Merge remote-tracking branch 'origin/main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-05 20:43:53 +08:00
houseme bdbdca07c8 fix(deps): preserve supported hotpath focus expressions
Keep the profiler runtime before its regex-lite compatibility regression.
Track the opt-in validation required to remove this constraint in backlog.

Refs rustfs/backlog#2302.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 19:50:46 +08:00
houseme 53efaa2b8f chore(deps): refresh profiling dependencies for the next batch
Update hotpath and its macro crate to the compatible patch release before
the next dependency-ready implementation tasks.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 19:10:54 +08:00
houseme ec9672a397 Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b4-base 2026-09-05 18:54:50 +08:00
houseme cee35f7e54 Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b3-base 2026-09-05 16:43:37 +08:00
houseme ef7e7afd8c Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b3-base 2026-09-05 16:35:05 +08:00
houseme 652ebb12c6 fix(ecstore): remove duplicate local rename implementation
Keep the canonical commit module after concurrent storage changes merged.
The control-write and rollback changes are already present there.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 16:28:01 +08:00
houseme 38c03d9d5d chore(deps): refresh scanner heal batch dependency baseline
Regenerate compatible lockfile selections before the next implementation
batch. Cargo upgrade leaves direct requirements unchanged.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 16:22:13 +08:00
21 changed files with 2207 additions and 154 deletions
@@ -20,9 +20,10 @@
//! journal (`count_requests`) carries the assertion in every one of them.
use super::common::{BoxError, OdmTestEnv, RawResponse, SeedObject, start_configured_env};
use crate::fake_s3_target::Operation;
use crate::fake_s3_target::{FaultAction, Operation};
use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
use bytes::Bytes;
use futures::{StreamExt, TryStreamExt};
use std::time::Duration;
type TestResult = Result<(), BoxError>;
@@ -145,14 +146,38 @@ async fn test_odm_range_burst_overflows_the_pull_queue_without_failing_clients()
.await?;
let body = payload(128 * 1024);
let blocker = "queue/blocker.bin";
env.seed_source(SOURCE_BUCKET, &[SeedObject::new(blocker, body.clone())]);
// The one-chunk range completes immediately; its full background pull
// occupies the only slot while the remaining requests fill the queue.
env.source.inject_for_key(
Operation::GetObject,
blocker,
FaultAction::SlowSendBody {
chunk_bytes: 1024,
delay: Duration::from_millis(100),
},
2,
);
let response = env
.raw_object_request(http::Method::GET, bucket, blocker, &[("range", "bytes=0-1023")])
.await?;
assert_eq!(response.status, 206);
assert_eq!(response.body, body.slice(0..1024));
env.wait_for_status_counter(bucket, "/inflight_pulls", 1, SETTLE).await?;
let keys: Vec<String> = (0..REQUESTS).map(|index| format!("queue/object-{index:03}.bin")).collect();
let seeds: Vec<SeedObject> = keys.iter().map(|key| SeedObject::new(key.clone(), body.clone())).collect();
env.seed_source(SOURCE_BUCKET, &seeds);
let responses: Vec<RawResponse> = futures::future::try_join_all(
// Bound source connections below the fixture's limit while still
// submitting all 100 requests to the eight-slot background queue.
let responses: Vec<RawResponse> = futures::stream::iter(
keys.iter()
.map(|key| env.raw_object_request(http::Method::GET, bucket, key, &[("range", "bytes=0-1023")])),
)
.buffered(16)
.try_collect()
.await?;
for (key, response) in keys.iter().zip(&responses) {
assert_eq!(response.status, 206, "{key}: {}", String::from_utf8_lossy(&response.body));
@@ -168,6 +193,15 @@ async fn test_odm_range_burst_overflows_the_pull_queue_without_failing_clients()
.wait_for_status_counter(bucket, "/counters/pull_failures_total/queue_full", 1, SETTLE)
.await?;
assert!(queue_full > 0, "a 100-deep burst must overflow an 8-slot queue");
let queue_full = usize::try_from(queue_full)?;
assert!(queue_full <= REQUESTS);
env.wait_for_status_counter(
bucket,
"/counters/pulled_objects_total/background",
u64::try_from(REQUESTS + 1 - queue_full)?,
SETTLE,
)
.await?;
let ranged_reads: usize = keys.iter().map(|key| source_get_count(&env, key)).sum();
assert!(
@@ -175,9 +209,6 @@ async fn test_odm_range_burst_overflows_the_pull_queue_without_failing_clients()
"every reader is served from the source: {ranged_reads} GETs for {REQUESTS} readers"
);
let dropped = keys.iter().filter(|key| source_get_count(&env, key) == 1).count();
assert!(
dropped > 0,
"the overflowed keys are the ones with no backfill GET, but every key got one"
);
assert_eq!(dropped, queue_full, "only overflowed keys remain without a background GET");
Ok(())
}
@@ -265,16 +265,13 @@ async fn list_through_rejects_a_tampered_continuation_token() -> TestResult {
let decoded = String::from_utf8(base64_simd::STANDARD.decode_to_vec(token.as_bytes())?)?;
assert!(decoded.contains("\"t\":\"odm-list\""), "the merged token is an envelope: {decoded}");
let tampered = base64_simd::STANDARD.encode_to_string(decoded.replace("\"v\":1", "\"v\":2").as_bytes());
let rejected = env
.raw_list_objects_v2(bucket, &format!("continuation-token={tampered}"))
.await?;
assert_eq!(
rejected.status,
400,
"a bumped token version is a client error: {}",
String::from_utf8_lossy(&rejected.body)
);
let tampered = base64_simd::STANDARD.encode_to_string(decoded.replace("\"v\":1", "\"v\":3").as_bytes());
assert_ne!(tampered, token, "the test must change the token version");
let query = serde_urlencoded::to_string([("continuation-token", tampered.as_str())])?;
let rejected = env.raw_list_objects_v2(bucket, &query).await?;
let error_body = String::from_utf8_lossy(&rejected.body);
assert_eq!(rejected.status, 400, "a bumped token version is a client error: {}", error_body);
assert!(error_body.contains("<Code>InvalidArgument</Code>"), "{error_body}");
Ok(())
}
+6
View File
@@ -4498,6 +4498,12 @@ impl SetDisks {
&self.ctx
}
/// Read the persisted bucket identity through this set's metadata owner.
/// Missing or non-authoritative legacy identities remain errors.
pub async fn bucket_incarnation_id_from_disk(&self, bucket: &str) -> Result<Uuid> {
metadata_sys::get_bucket_incarnation_id_in(&self.ctx, bucket).await
}
/// Admit one short scanner cache publication under this set's instance
/// movement fence. The caller must hold the returned guard through its
/// final conditional cache write; no scan-round work belongs under it.
+298
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,
@@ -400,6 +401,85 @@ 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,
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(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()
&& matches!(self.scan_mode, HealScanMode::Normal | HealScanMode::Deep)
}
}
/// 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, 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,
@@ -470,6 +550,12 @@ 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 scan_coverage_receipt: Option<DataUsageScanCoverageReceipt>,
#[serde(default)]
pub pending_heals: Vec<PendingScannerHeal>,
#[serde(default)]
pub object_lock: Option<Arc<ObjectLockConfiguration>>,
@@ -481,6 +567,11 @@ pub struct DataUsageCacheInfo {
pub snapshot_complete: bool,
#[serde(default)]
pub scan_plan_digest: Option<DataUsageScanPlanDigest>,
/// Full activity and inventory scope of a set scan; only a complete
/// snapshot proves coverage. Bucket caches bind this scope into their
/// opaque scan plan digest instead.
#[serde(default)]
pub scan_coverage_digest: Option<DataUsageScanPlanDigest>,
#[serde(default)]
pub cache_key_format: u16,
/// Registry generation used for the completed/partial scan. This is
@@ -517,6 +608,10 @@ 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.scan_coverage_receipt.is_some())
+ usize::from(self.scan_coverage_digest.is_some())
+ usize::from(self.tier_registry_generation.is_some())
+ usize::from(!self.size_reconciliation.is_empty())
+ usize::from(self.lkg_snapshot_complete)
@@ -536,11 +631,23 @@ 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)?;
}
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)?;
state.serialize_entry("snapshot_complete", &self.snapshot_complete)?;
state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?;
if let Some(coverage) = self.scan_coverage_digest {
state.serialize_entry("scan_coverage_digest", &coverage)?;
}
state.serialize_entry("cache_key_format", &self.cache_key_format)?;
if let Some(generation) = self.tier_registry_generation {
state.serialize_entry("tier_registry_generation", &generation)?;
@@ -703,6 +810,174 @@ 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
&& 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 {
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 = (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;
}
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;
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.
self.info.scan_plan_digest = None;
if reusable {
DataUsageCachePrepareOutcome::Reused
} else {
DataUsageCachePrepareOutcome::Reset
}
}
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_key(|(key, _)| *key);
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!(
@@ -782,6 +1057,29 @@ impl DataUsageCache {
(visited == expected_entries).then_some(entry)
}
pub(crate) fn has_complete_root_inventory(&self, bucket_keys: &HashSet<String>) -> bool {
let Some(root) = self.find(DATA_USAGE_ROOT) else {
return false;
};
// Set roots only connect bucket entries. Scalar data at the root, an
// extra bucket, or an orphan must not disappear during bucket folding.
root.children.len() == bucket_keys.len()
&& bucket_keys.iter().all(|key| root.children.contains(key))
&& root.size == 0
&& root.objects == 0
&& root.versions == 0
&& root.delete_markers == 0
&& root.failed_objects == 0
&& !root.compacted
&& root.obj_sizes.is_empty()
&& root.obj_versions.is_empty()
&& root.replication_stats.is_none()
&& root.all_tier_stats.is_none()
&& root.unknown_tier_stats.is_none()
&& root.tier_accounting_proof.is_none()
&& self.checked_flatten_complete(DATA_USAGE_ROOT).is_some()
}
fn checked_flatten_inner(&self, path: &str) -> Option<(DataUsageEntry, usize)> {
let root_key = hash_path(path).key();
let (root_key, root) = self.cache.get_key_value(&root_key)?;
@@ -29,6 +29,31 @@ use tokio::sync::Mutex;
const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([3; 32]);
#[test]
fn scoped_scan_coverage_metadata_preserves_map_compatibility() {
#[derive(serde::Deserialize)]
struct LegacyInfo {
name: String,
next_cycle: u64,
}
let mut info = DataUsageCacheInfo {
name: DATA_USAGE_ROOT.to_string(),
next_cycle: 7,
..Default::default()
};
let old = serde_json::to_value(&info).expect("legacy metadata should encode");
assert!(old.get("scan_coverage_digest").is_none());
let old: DataUsageCacheInfo = serde_json::from_value(old).expect("missing coverage must remain readable");
assert!(old.scan_coverage_digest.is_none());
info.scan_coverage_digest = Some(TEST_PLAN_DIGEST);
let encoded = rmp_serde::to_vec(&info).expect("coverage metadata should remain map encoded");
let legacy: LegacyInfo = rmp_serde::from_slice(&encoded).expect("old map readers should ignore additive proof fields");
assert_eq!(legacy.name, DATA_USAGE_ROOT);
assert_eq!(legacy.next_cycle, 7);
let decoded: DataUsageCacheInfo = rmp_serde::from_slice(&encoded).expect("new reader should restore the coverage proof");
assert_eq!(decoded.scan_coverage_digest, Some(TEST_PLAN_DIGEST));
}
#[derive(Debug, PartialEq, Eq)]
struct CachePutRecord {
object: String,
@@ -728,6 +728,15 @@ 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,
scan_mode,
)
.await
.ok(),
},
);
match scan_state {
@@ -821,6 +830,22 @@ 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,
scan_mode,
)
.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(
+32 -22
View File
@@ -1086,6 +1086,12 @@ impl ScannerMaintenanceFeatures {
fn needs_regular_cycle(self) -> bool {
self.lifecycle || self.replication || self.inspection_failed
}
fn requires_full_scan(self, observed_generation: Option<u64>, current_generation: u64, wake: ScannerCycleWakeReason) -> bool {
self.needs_regular_cycle()
|| observed_generation != Some(current_generation)
|| !matches!(wake, ScannerCycleWakeReason::DirtyUsage | ScannerCycleWakeReason::ClusterActivity)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
@@ -1283,18 +1289,18 @@ async fn configure_scanner_defaults(
ctx: &CancellationToken,
storeapi: &Arc<impl ScannerStorage>,
) -> (ScannerMaintenanceFeatures, Option<u64>) {
let (features, maintenance_generation) = detect_stable_scanner_maintenance_features(ctx, storeapi)
.await
.unwrap_or_else(|| {
(
ScannerMaintenanceFeatures {
inspection_failed: true,
..Default::default()
},
scanner_maintenance_generation(),
)
});
if storeapi.setup_is_erasure_sd().await {
let (features, maintenance_generation) = detect_stable_scanner_maintenance_features(ctx, storeapi)
.await
.unwrap_or_else(|| {
(
ScannerMaintenanceFeatures {
inspection_failed: true,
..Default::default()
},
scanner_maintenance_generation(),
)
});
// Single-disk keeps the speed-preset-derived default cycle (60s at the
// `default` preset) instead of a special shorter cycle: no measured
// cold-start ILM latency basis for an override, and clean-idle backoff
@@ -1319,7 +1325,7 @@ async fn configure_scanner_defaults(
} else {
set_scanner_default_speed(ScannerSpeed::Default);
set_scanner_default_cycle_secs(None);
(ScannerMaintenanceFeatures::default(), None)
(features, Some(maintenance_generation))
}
}
@@ -1564,7 +1570,7 @@ 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).await
run_data_scanner_cycle_with_budget(ctx, storeapi, cycle_info, cycle_revision, leader_epoch, cycle_budget, true).await
}
#[instrument(skip_all)]
@@ -1576,6 +1582,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,
) -> ScannerCycleOutcome
where
S: ScannerStorage,
@@ -1714,6 +1721,9 @@ where
scan_mode,
scan_scope: crate::scanner_io::ScannerBucketScanScope::default(),
persisted_usage_baseline: usage_persist_baseline.data.clone(),
requires_full_scan,
#[cfg(test)]
resolved_scope_observer: None,
},
)
.await;
@@ -2574,10 +2584,7 @@ where
let mut superseded_backoff = ScannerRetryBackoff::default();
let mut deferred_backoff = ScannerRetryBackoff::default();
let initial_runtime_config = resolve_scanner_runtime_config();
if clean_idle_topology_supported
&& scanner_clean_idle_backoff_configured(&initial_runtime_config)
&& maintenance_generation_seen.is_none()
{
if clean_idle_topology_supported && maintenance_generation_seen.is_none() {
let Some((features, generation)) = detect_stable_scanner_maintenance_features(&ctx, &storeapi).await else {
global_metrics().set_cycle(None).await;
finish_scanner_leader_iteration(false, "stopped", String::new()).await;
@@ -2785,6 +2792,7 @@ where
&mut cycle_revision,
leader_epoch,
cycle_budget.clone(),
true,
),
guard.lock_lost_notified(),
)
@@ -2879,7 +2887,7 @@ where
#[cfg(test)]
notify_scanner_runtime_observed_for_test(&storeapi, pause_backlog_observation);
let runtime_config = resolve_scanner_runtime_config();
if clean_idle_topology_supported && scanner_clean_idle_backoff_configured(&runtime_config) {
if clean_idle_topology_supported {
let current_generation = scanner_maintenance_generation();
if maintenance_generation_seen != Some(current_generation) {
scanner_activity_seen = None;
@@ -3075,6 +3083,11 @@ where
&mut cycle_revision,
leader_epoch,
cycle_budget.clone(),
maintenance_features.requires_full_scan(
maintenance_generation_seen,
scanner_maintenance_generation(),
wake_reason,
),
),
guard.lock_lost_notified(),
)
@@ -3132,10 +3145,7 @@ where
let maintenance_config_changed =
maintenance_generation_seen.is_some_and(|generation| generation != current_maintenance_generation);
let retry_failed_inspection = maintenance_inspection_retry.retry_due(maintenance_features, wake_reason, Instant::now());
if clean_idle_topology_supported
&& scanner_clean_idle_backoff_configured(&runtime_config)
&& (maintenance_config_changed || retry_failed_inspection)
{
if clean_idle_topology_supported && (maintenance_config_changed || retry_failed_inspection) {
let Some((features, generation)) = detect_stable_scanner_maintenance_features(&ctx, &storeapi).await else {
break;
};
+159 -2
View File
@@ -1204,7 +1204,7 @@ 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)),
run_data_scanner_cycle_with_budget(&ctx, &store, &mut cycle_info, &mut revision, 1, Arc::clone(&budget), true),
)
.await
.expect("the coordinator must finish its namespace walk while a PUT is pending");
@@ -1240,7 +1240,7 @@ 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)),
run_data_scanner_cycle_with_budget(&ctx, &store, &mut cycle_info, &mut revision, 1, Arc::clone(&retry_budget), true),
)
.await
.expect("the same cycle must converge after the pending PUT drains");
@@ -4869,6 +4869,28 @@ async fn usage_bootstrap_does_not_overwrite_concurrent_replacement() {
#[serial]
async fn scanner_usage_state_reset_publishes_fenced_bootstrap_marker() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let quota_ledger_path = "config/quota-ledger/reserved-bucket.json";
let quota_ledger = serde_json::to_vec(&serde_json::json!({
"version": 1,
"bucket_incarnation": "00000000-0000-0000-0000-000000000001",
"quota_revision_unix_nanos": 1,
"accounted_usage": 100,
"reservations": {
"00000000-0000-0000-0000-000000000002": {
"object": "pending-object",
"old_size": 0,
"new_size": 64,
"created_at": 1,
"pool_index": 0,
"set_index": 0,
"commit_started": true
}
}
}))
.expect("quota ledger fixture should encode");
save_config(store.clone(), quota_ledger_path, quota_ledger.clone())
.await
.expect("independent quota reservations should persist");
let cycle = CurrentCycle {
current: 41,
next: 42,
@@ -4927,6 +4949,14 @@ async fn scanner_usage_state_reset_publishes_fenced_bootstrap_marker() {
assert!(!data_usage_info_has_persisted_baseline_identity(&usage));
assert_eq!(usage.scanner_epoch, Some(9));
assert_eq!(
read_config(store.clone(), quota_ledger_path)
.await
.expect("quota ledger must remain readable after scanner reset"),
quota_ledger,
"scanner reset must preserve incarnation and outstanding reserved bytes exactly"
);
for path in [
usage_backup_path.as_str(),
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str(),
@@ -8208,6 +8238,76 @@ fn clean_idle_backoff_policy_preserves_explicit_and_maintenance_cycles() {
}
}
#[tokio::test]
#[serial]
async fn scoped_scan_explicit_bitrot_keeps_dirty_planning_without_idle_backoff() {
temp_env::async_with_vars([(ENV_SCANNER_CYCLE, None), (ENV_SCANNER_BITROT_CYCLE_SECS, Some("3600"))], async {
crate::runtime_config::refresh_scanner_runtime_config_for_tests();
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
let (_temp_dir, store) = setup_scanner_cycle_store_with_pool_count(true, 2).await;
let ctx = CancellationToken::new();
let (features, generation) = configure_scanner_defaults(&ctx, &store).await;
let config = resolve_scanner_runtime_config();
assert_eq!(config.cycle_interval_source, ScannerRuntimeConfigSource::Default);
assert_eq!(config.bitrot_cycle_source, ScannerRuntimeConfigSource::Env);
assert!(!scanner_clean_idle_backoff_configured(&config));
assert!(!features.needs_regular_cycle());
assert_eq!(
generation,
Some(scanner_maintenance_generation()),
"multi-disk startup must inspect maintenance independently"
);
let observed = ScannerCycleObservedGenerations::for_wait(&config, None, 7, 0, scanner_maintenance_generation());
assert_eq!(observed.dirty_usage, Some(7), "explicit bitrot still permits ordinary dirty wakeups");
for (wake, full) in [
(ScannerCycleWakeReason::DirtyUsage, false),
(ScannerCycleWakeReason::ClusterActivity, false),
(ScannerCycleWakeReason::Timer, true),
(ScannerCycleWakeReason::ClusterMaintenance, true),
] {
assert_eq!(
features.requires_full_scan(generation, scanner_maintenance_generation(), wake),
full,
"{wake:?}"
);
}
assert!(features.requires_full_scan(None, scanner_maintenance_generation(), ScannerCycleWakeReason::DirtyUsage));
for unsafe_features in [
ScannerMaintenanceFeatures {
lifecycle: true,
..Default::default()
},
ScannerMaintenanceFeatures {
replication: true,
..Default::default()
},
ScannerMaintenanceFeatures {
inspection_failed: true,
..Default::default()
},
] {
assert!(unsafe_features.requires_full_scan(
generation,
scanner_maintenance_generation(),
ScannerCycleWakeReason::DirtyUsage
));
}
crate::scanner_io::record_scanner_maintenance_change("maintenance-proof-change");
assert!(features.requires_full_scan(generation, scanner_maintenance_generation(), ScannerCycleWakeReason::DirtyUsage));
let (refreshed, refreshed_generation) = detect_stable_scanner_maintenance_features(&ctx, &store)
.await
.expect("changed maintenance generation should be inspected");
assert!(!refreshed.requires_full_scan(
Some(refreshed_generation),
scanner_maintenance_generation(),
ScannerCycleWakeReason::DirtyUsage
));
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
})
.await;
crate::runtime_config::refresh_scanner_runtime_config_for_tests();
}
#[test]
fn clean_idle_backoff_requires_activity_probes() {
let default_config = ScannerRuntimeConfig::default();
@@ -8596,6 +8696,63 @@ fn scanner_node_activity(epoch: &str, namespace_generation: u64, maintenance_gen
}
}
#[test]
fn scoped_scan_remote_dirty_coverage_invalidates_local_bucket_current() {
let before = BTreeMap::from([("remote".to_string(), scanner_node_activity("epoch-a", 7, 3))]);
let mut after = before.clone();
let remote = after.get_mut("remote").expect("remote activity should exist");
remote.dirty_usage_generation += 1;
remote.dirty_usage_pending = true;
assert_eq!(scanner_activity_structural_digest(&before), scanner_activity_structural_digest(&after));
let old_plan = crate::scanner_io::checkpoint_fixture_bucket_digest(
DataUsageScanPlanDigest(scanner_activity_snapshot_digest(&before)),
None,
);
let new_plan = crate::scanner_io::checkpoint_fixture_bucket_digest(
DataUsageScanPlanDigest(scanner_activity_snapshot_digest(&after)),
None,
);
assert_ne!(
old_plan, new_plan,
"remote dirty changes must fence Current even without a local bucket hint"
);
let source = DataUsageCacheSource::new(0, 0);
let mut cache = DataUsageCache {
info: crate::DataUsageCacheInfo {
name: "bucket".to_string(),
next_cycle: 7,
leader_epoch: 11,
source: Some(source),
last_update: Some(std::time::SystemTime::UNIX_EPOCH),
snapshot_complete: true,
scan_plan_digest: Some(old_plan),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
..Default::default()
};
cache.replace("bucket", "", DataUsageEntry::default());
assert!(matches!(
crate::scanner_io::current_cache_root_or_prepare_with_generation(
&mut cache,
"bucket",
source,
7,
11,
new_plan,
crate::scanner_io::DataUsageCacheReuseOptions {
require_source: true,
tier_registry_generation: None,
checkpoint_identity: None,
},
),
crate::scanner_io::DataUsageCacheScanState::Prepared {
outcome: DataUsageCachePrepareOutcome::Reset,
..
}
));
}
#[test]
fn post_lease_activity_proof_rejects_a_put_tail_that_finished_before_lease_acquisition() {
let before = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]);
+83 -7
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]);
@@ -1741,7 +1760,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 = self.resume_frontier.clone();
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 +1783,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 +1826,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...
@@ -1810,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 {
@@ -1835,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);
@@ -2250,7 +2288,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())
{
@@ -2359,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,
@@ -2388,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,
@@ -2418,23 +2463,41 @@ 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();
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;
new_cache.info.scan_coverage_receipt = None;
if had_scan_checkpoint {
global_metrics().record_scanner_checkpoint_cleared();
}
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())
@@ -2448,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);
@@ -2462,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,
@@ -234,6 +234,531 @@ 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,
scan_mode: HealScanMode::Normal,
};
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
.seal_scan_frontier(Some("bucket/static"))
.expect("completed fixture prefix receipt");
(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");
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),
crate::DataUsageCachePrepareOutcome::Reset
);
assert!(old_view.cache.is_empty());
assert!(!old_view.info.snapshot_complete);
}
#[test]
fn checkpoint_fixture_optional_coverage_metadata_roundtrips_together() {
let (mut cache, _) = bound_checkpoint();
cache.info.scan_coverage_digest = Some(DataUsageScanPlanDigest([77; 32]));
let encoded = cache.marshal_msg().expect("encode every optional coverage field");
let decoded = DataUsageCache::unmarshal(&encoded).expect("decode combined W03/W04 metadata map");
assert_eq!(decoded.info.scan_identity, cache.info.scan_identity);
assert_eq!(decoded.info.scan_progress, cache.info.scan_progress);
assert_eq!(decoded.info.scan_coverage_receipt, cache.info.scan_coverage_receipt);
assert_eq!(decoded.info.scan_coverage_digest, cache.info.scan_coverage_digest);
assert_eq!(decoded.validated_scan_frontier(), Some("bucket/static"));
assert!(!decoded.info.snapshot_complete);
}
#[test]
fn checkpoint_fixture_unchanged_complete_plan_keeps_existing_rescan_policy() {
let (mut cache, identity) = bound_checkpoint();
cache.info.scan_progress = None;
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),
crate::DataUsageCachePrepareOutcome::Reused
);
assert!(
cache.info.scan_progress.is_none(),
"unchanged complete coverage needs no forced verification sweep"
);
assert_eq!(cache.info.scan_plan_digest, Some(PLAN));
assert_eq!(retained(&cache), 3);
let next = DataUsageScanPlanDigest([44; 32]);
cache.prepare_bucket_checkpoint("bucket", 12, 7, SOURCE, next, identity);
assert_eq!(cache.info.scan_progress.expect("changed plan must be verified").started_plan, next);
assert!(cache.info.scan_plan_digest.is_none());
assert!(!cache.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
},
crate::DataUsageScanIdentity {
scan_mode: HealScanMode::Deep,
..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());
}
}
#[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() {
@@ -242,23 +767,48 @@ 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::from_unix_timestamp(1_700_000_000 + i64::try_from(index).expect("fixture index"))
.expect("non-sentinel fixture modification time"),
);
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,
scan_mode: HealScanMode::Normal,
};
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,9 +828,11 @@ 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);
cache.info.skip_healing = true;
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(
&parent,
@@ -326,13 +878,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 +934,85 @@ 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),
},
);
cache.info.skip_healing = true;
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");
}
+97 -2
View File
@@ -183,6 +183,18 @@ fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<
return None;
}
// Completed maintenance also covers ordinary usage. Keep its exact stored
// proof for cache reuse, and reject mixtures of different set work proofs.
let baseline_plan_digest = DataUsageScanPlanDigest(baseline.usage_snapshot_set_states.first()?.scan_plan_digest?);
if ![
proof.scan_plan_digest,
scanner_bucket_work_digest(proof.scan_plan_digest, HealScanMode::Normal, true),
scanner_bucket_work_digest(proof.scan_plan_digest, HealScanMode::Deep, true),
]
.contains(&baseline_plan_digest)
{
return None;
}
let mut states = HashSet::with_capacity(baseline.usage_snapshot_set_states.len());
for state in &baseline.usage_snapshot_set_states {
let source = DataUsageCacheSource::new(usize::try_from(state.pool_index).ok()?, usize::try_from(state.set_index).ok()?);
@@ -192,13 +204,13 @@ fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<
|| state.tombstone
|| state.scanner_epoch != Some(proof.leader_epoch)
|| state.scanner_cycle.is_none_or(|cycle| cycle > proof.want_cycle)
|| state.scan_plan_digest != Some(proof.scan_plan_digest.0)
|| state.scan_plan_digest != Some(baseline_plan_digest.0)
{
return None;
}
}
(states == *proof.expected_sources).then_some(proof.scan_plan_digest)
(states == *proof.expected_sources).then_some(baseline_plan_digest)
}
fn scoped_scan_scope_from_dirty_buckets(
@@ -271,6 +283,9 @@ 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,
requires_full_scan: bool,
// Cache work must invalidate on namespace completion even when its scoped baseline remains reusable.
execution_digest: DataUsageScanPlanDigest,
leader_epoch: u64,
@@ -314,6 +329,53 @@ fn scanner_bucket_plan_digest(buckets: &[BucketInfo], activity_digest: [u8; 32])
DataUsageScanPlanDigest(hasher.finalize().into())
}
fn scanner_bucket_inventory_is_complete(
all_buckets: &[BucketInfo],
buckets_by_source: &HashMap<DataUsageCacheSource, Vec<BucketInfo>>,
) -> bool {
let inventory = all_buckets
.iter()
.map(|bucket| (bucket.name.as_str(), bucket.created))
.collect::<HashMap<_, _>>();
if inventory.len() != all_buckets.len() || inventory.keys().any(|name| name.is_empty() || *name == DATA_USAGE_ROOT) {
return false;
}
let mut covered = HashSet::with_capacity(inventory.len());
for buckets in buckets_by_source.values() {
let mut set_names = HashSet::with_capacity(buckets.len());
for bucket in buckets {
if !set_names.insert(bucket.name.as_str()) || inventory.get(bucket.name.as_str()) != Some(&bucket.created) {
return false;
}
covered.insert(bucket.name.as_str());
}
}
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>,
@@ -711,6 +773,39 @@ 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,
scan_mode: HealScanMode,
) -> 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,
scan_mode,
})
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum ScannerCycleDeferReason {
ActivityBaselineUnavailable,
+107 -27
View File
@@ -112,6 +112,10 @@ 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_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
@@ -137,6 +141,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 +164,7 @@ pub(crate) fn current_cache_root_or_prepare(
DataUsageCacheReuseOptions {
require_source,
tier_registry_generation: None,
checkpoint_identity: None,
},
)
}
@@ -172,6 +178,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 +205,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),
},
},
}
}
@@ -213,10 +230,87 @@ pub(super) fn cache_snapshot_is_current(
)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) struct ScannerSnapshotIdentity {
pub(super) cycle: u64,
pub(super) leader_epoch: u64,
pub(super) plan_digest: DataUsageScanPlanDigest,
pub(super) coverage_digest: DataUsageScanPlanDigest,
pub(super) tier_registry_generation: Option<u64>,
}
pub(super) struct ScannerSnapshotScope<'a> {
pub(super) sources: &'a HashSet<DataUsageCacheSource>,
pub(super) buckets: &'a [String],
pub(super) identity: ScannerSnapshotIdentity,
}
#[derive(Debug, PartialEq, Eq, thiserror::Error)]
pub(super) enum ScannerSnapshotValidationError {
#[error("scanner snapshot does not cover the expected complete sets")]
IncompleteSets,
#[error("scanner snapshot does not match the requested generation")]
GenerationMismatch,
#[error("scanner snapshot bucket inventory is invalid")]
InvalidInventory,
#[error("scanner snapshot root is incomplete or corrupt")]
InvalidRoot,
}
struct ValidatedScannerSnapshot<'a> {
results: &'a [DataUsageCache],
last_update: SystemTime,
}
impl<'a> ValidatedScannerSnapshot<'a> {
fn validate(
results: &'a [DataUsageCache],
scope: &ScannerSnapshotScope<'_>,
) -> std::result::Result<Self, ScannerSnapshotValidationError> {
if !scanner_results_form_complete_snapshot(results, scope.sources) {
return Err(ScannerSnapshotValidationError::IncompleteSets);
}
let bucket_keys = scope
.buckets
.iter()
.map(|bucket| crate::hash_path(bucket).key())
.collect::<HashSet<_>>();
if bucket_keys.len() != scope.buckets.len()
|| scope
.buckets
.iter()
.any(|bucket| bucket.is_empty() || bucket == DATA_USAGE_ROOT)
{
return Err(ScannerSnapshotValidationError::InvalidInventory);
}
for result in results {
if result.info.next_cycle != scope.identity.cycle
|| result.info.leader_epoch != scope.identity.leader_epoch
|| result.info.scan_plan_digest != Some(scope.identity.plan_digest)
|| result.info.scan_coverage_digest != Some(scope.identity.coverage_digest)
|| result.info.tier_registry_generation != scope.identity.tier_registry_generation
{
return Err(ScannerSnapshotValidationError::GenerationMismatch);
}
if result.info.name != DATA_USAGE_ROOT
|| result.info.cache_key_format != DATA_USAGE_CACHE_KEY_FORMAT
|| !result.has_complete_root_inventory(&bucket_keys)
{
return Err(ScannerSnapshotValidationError::InvalidRoot);
}
}
let last_update = results
.iter()
.filter_map(|result| result.info.last_update)
.max()
.ok_or(ScannerSnapshotValidationError::IncompleteSets)?;
Ok(Self { results, last_update })
}
}
pub(super) fn completed_data_usage_info(
results: &[DataUsageCache],
expected_sources: &HashSet<DataUsageCacheSource>,
all_buckets: &[String],
scope: &ScannerSnapshotScope<'_>,
tier_registry_names: &[String],
bucket_plan_complete: bool,
budget_elapsed: bool,
@@ -229,26 +323,10 @@ pub(super) fn completed_data_usage_info(
if !should_publish_completed_snapshot(completed_set_count, results.len(), budget_elapsed, cancelled) {
return None;
}
if !scanner_results_form_complete_snapshot(results, expected_sources) {
return None;
}
// A generation is comparable across nodes because it is derived from the
// frozen registry names. Cycle and leader fencing remain separate cache
// metadata. Legacy peers omit the generation; an all-legacy result remains
// readable, but mixing legacy and new (or two new generations) would make
// the per-tier accounting ambiguous.
let registry_generation = results.first()?.info.tier_registry_generation;
if results.iter().any(|result| match registry_generation {
Some(generation) => result.info.tier_registry_generation != Some(generation),
None => result.info.tier_registry_generation.is_some(),
}) {
return None;
}
if results.iter().any(|result| result.root().is_none()) {
return None;
}
let validated = ValidatedScannerSnapshot::validate(results, scope).ok()?;
let results = validated.results;
let all_buckets = scope.buckets;
let registry_generation = scope.identity.tier_registry_generation;
let mut total = DataUsageEntry::default();
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
@@ -273,7 +351,7 @@ pub(super) fn completed_data_usage_info(
return None;
}
let merged_last_update = results.iter().filter_map(|result| result.info.last_update).max()?;
let merged_last_update = validated.last_update;
let buckets_usage = bucket_entries
.iter()
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
@@ -300,8 +378,8 @@ pub(super) fn completed_data_usage_info(
usage_snapshot_set_states.sort_by_key(|state| (state.pool_index, state.set_index));
let data_usage_info = DataUsageInfo {
last_update: Some(merged_last_update),
scanner_cycle: Some(results.first()?.info.next_cycle),
scanner_epoch: Some(results.first()?.info.leader_epoch),
scanner_cycle: Some(scope.identity.cycle),
scanner_epoch: Some(scope.identity.leader_epoch),
objects_total_count: u64::try_from(total.objects).ok()?,
versions_total_count: u64::try_from(total.versions).ok()?,
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
@@ -609,6 +687,7 @@ pub(super) async fn persist_and_publish_cache_snapshot(
expected_publication_epoch: u64,
) -> Option<SystemTime> {
let source = cache_snapshot.info.source?;
let coverage_digest = cache_snapshot.info.scan_coverage_digest?;
let execution_digest = cache_snapshot.info.scan_execution_digest?;
let guard = match acquire_scanner_cache_locks(store.as_ref(), DATA_USAGE_CACHE_NAME, source).await {
Ok(guard) => guard,
@@ -674,7 +753,8 @@ pub(super) async fn persist_and_publish_cache_snapshot(
);
return None;
}
if persisted.info.scan_execution_digest == Some(execution_digest)
if persisted.info.scan_coverage_digest == Some(coverage_digest)
&& persisted.info.scan_execution_digest == Some(execution_digest)
&& matches!(
current_cache_root_entry_with_generation(
&persisted,
+44 -13
View File
@@ -39,6 +39,12 @@ pub(super) fn prepare_scoped_set_scan(
else {
return None;
};
// The existing cache does not bind each bucket to a durable incarnation.
// Listing creation times can come from volume metadata, so even Some(time)
// cannot prove that an unselected same-name bucket is the cached bucket.
if all_buckets.iter().any(|bucket| !selected_buckets.contains(&bucket.name)) {
return None;
}
if selected_buckets.is_empty()
|| !old_cache.info.snapshot_complete
|| old_cache.info.last_update.is_none()
@@ -49,7 +55,7 @@ pub(super) fn prepare_scoped_set_scan(
|| old_cache.info.source != Some(generation.source)
|| old_cache.info.scan_plan_digest != Some(baseline_scan_plan_digest)
|| old_cache.info.cache_key_format != DATA_USAGE_CACHE_KEY_FORMAT
|| old_cache.checked_flatten_complete_scope(DATA_USAGE_ROOT).is_none()
|| !old_cache.has_complete_root_inventory(&old_cache.find(DATA_USAGE_ROOT)?.children)
{
return None;
}
@@ -74,21 +80,12 @@ pub(super) fn prepare_scoped_set_scan(
cache: HashMap::new(),
};
cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default());
let root_hash = crate::hash_path(DATA_USAGE_ROOT);
let mut current_bucket_names = HashSet::with_capacity(all_buckets.len());
for bucket in all_buckets {
if !current_bucket_names.insert(bucket.name.as_str()) {
return None;
}
if selected_buckets.contains(&bucket.name) {
cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default());
continue;
}
let bucket_hash = crate::hash_path(&bucket.name);
old_cache.find(&bucket.name)?;
cache.copy_with_children(old_cache, &bucket_hash, &Some(root_hash.clone()));
cache.find(&bucket.name)?;
cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default());
}
Some(PreparedScopedSetScan {
@@ -118,6 +115,8 @@ impl ScannerIOCache for SetDisks {
all_buckets,
scope,
digest: scan_plan_digest,
bucket_coverage_digest,
requires_full_scan,
execution_digest,
leader_epoch,
tier_registry_generation,
@@ -127,6 +126,8 @@ impl ScannerIOCache for SetDisks {
pending_maintenance_work,
cache_cycle_floor,
} = scan_plan;
let scan_plan_digest = scanner_bucket_work_digest(scan_plan_digest, scan_mode, requires_full_scan);
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();
@@ -169,8 +170,9 @@ impl ScannerIOCache for SetDisks {
scan_plan_digest,
},
);
let mut scoped_cache = scoped_scan.map(|prepared| {
let mut scoped_cache = scoped_scan.map(|mut prepared| {
buckets = prepared.buckets;
prepared.cache.info.scan_coverage_digest = Some(bucket_coverage_digest);
prepared.cache
});
if buckets.is_empty() {
@@ -186,6 +188,7 @@ impl ScannerIOCache for SetDisks {
tier_registry_generation: Some(tier_registry_generation),
source: Some(source),
scan_plan_digest: Some(scan_plan_digest),
scan_coverage_digest: Some(bucket_coverage_digest),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
@@ -475,6 +478,7 @@ impl ScannerIOCache for SetDisks {
source: Some(source),
snapshot_complete: false,
scan_plan_digest: Some(scan_plan_digest),
scan_coverage_digest: Some(bucket_coverage_digest),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
lkg_snapshot_complete: old_cache.info.lkg_snapshot_complete,
lkg_next_cycle: old_cache.info.lkg_next_cycle,
@@ -644,7 +648,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(execution_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;
@@ -890,6 +894,17 @@ 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,
scan_mode,
)
.await
.ok();
let scan_state = current_cache_root_or_prepare_with_generation(
&mut cache,
&bucket.name,
@@ -900,6 +915,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 {
@@ -1049,6 +1065,21 @@ 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,
scan_mode,
)
.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) => {
+40 -10
View File
@@ -72,6 +72,9 @@ where
scan_mode,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: None,
requires_full_scan: true,
#[cfg(test)]
resolved_scope_observer: None,
};
nsscanner_with_storage_status_scoped(store, request).await
}
@@ -85,6 +88,10 @@ pub(crate) struct ScannerCycleRequest {
pub(crate) scan_mode: HealScanMode,
pub(crate) scan_scope: ScannerBucketScanScope,
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,
#[cfg(test)]
pub(crate) resolved_scope_observer: Option<tokio::sync::oneshot::Sender<ScannerBucketScanScope>>,
}
struct ScannerBucketScopeResolution<'a> {
@@ -93,6 +100,7 @@ struct ScannerBucketScopeResolution<'a> {
activity_before: &'a crate::scanner::ScannerActivitySnapshot,
dirty_usage_snapshot: &'a DirtyUsageSnapshot,
all_buckets: &'a [BucketInfo],
requires_full_scan: bool,
}
async fn resolve_scanner_bucket_scan_scope<S>(
@@ -103,6 +111,9 @@ async fn resolve_scanner_bucket_scan_scope<S>(
where
S: ScannerStorage,
{
if resolution.requires_full_scan {
return ScannerBucketScanScope::default();
}
if !resolution.requested_scope.is_default()
|| !resolution.dirty_usage_snapshot.covers_all_pending
|| resolution.dirty_usage_snapshot.generation == u64::MAX
@@ -172,6 +183,9 @@ where
scan_mode,
scan_scope,
persisted_usage_baseline,
requires_full_scan,
#[cfg(test)]
resolved_scope_observer,
} = request;
let child_token = ctx.child_token();
let _tier_cycle_guard = begin_tier_registry_cycle(want_cycle, leader_epoch);
@@ -260,13 +274,13 @@ where
}
}
bucket_plan_complete &= buckets_by_source.keys().copied().collect::<HashSet<_>>() == *expected_sources;
let activity_digest = crate::scanner::scanner_activity_snapshot_digest(&activity_before);
let scan_plan_digest =
bucket_plan_complete &= scanner_bucket_inventory_is_complete(&all_buckets, &buckets_by_source);
let structural_scan_plan_digest =
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before));
let mut execution_hasher = Sha256::new();
execution_hasher.update(scan_plan_digest.0);
execution_hasher.update(activity_digest);
let execution_digest = DataUsageScanPlanDigest(execution_hasher.finalize().into());
let scan_plan_digest = scanner_bucket_work_digest(structural_scan_plan_digest, scan_mode, requires_full_scan);
let activity_digest = crate::scanner::scanner_activity_snapshot_digest(&activity_before);
let bucket_coverage_digest = scanner_bucket_plan_digest(&all_buckets, activity_digest);
let execution_digest = scanner_bucket_work_digest(bucket_coverage_digest, scan_mode, requires_full_scan);
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,
@@ -278,14 +292,19 @@ where
expected_sources: &expected_sources,
leader_epoch,
want_cycle,
scan_plan_digest,
scan_plan_digest: structural_scan_plan_digest,
},
activity_before: &activity_before,
dirty_usage_snapshot: &dirty_usage_snapshot,
all_buckets: &all_buckets,
requires_full_scan: requires_full_scan || scan_mode == HealScanMode::Deep,
},
)
.await;
#[cfg(test)]
if let Some(observer) = resolved_scope_observer {
let _ = observer.send(scan_scope.clone());
}
let cache_cycle_floor = Arc::new(AtomicU64::new(want_cycle));
let tier_registry = runtime_tier_registry_for_cycle(want_cycle, leader_epoch).await;
let tier_registry_generation = tier_registry.generation;
@@ -415,7 +434,9 @@ where
buckets: set_buckets,
all_buckets: Arc::clone(&all_buckets),
scope: scan_scope.clone(),
digest: scan_plan_digest,
digest: structural_scan_plan_digest,
bucket_coverage_digest,
requires_full_scan,
execution_digest,
leader_epoch,
tier_registry_generation,
@@ -543,8 +564,17 @@ where
let all_bucket_names = all_buckets.iter().map(|bucket| bucket.name.clone()).collect::<Vec<_>>();
let completed_usage = completed_data_usage_info(
&results,
&expected_sources,
&all_bucket_names,
&ScannerSnapshotScope {
sources: &expected_sources,
buckets: &all_bucket_names,
identity: ScannerSnapshotIdentity {
cycle: want_cycle,
leader_epoch,
plan_digest: scan_plan_digest,
coverage_digest: bucket_coverage_digest,
tier_registry_generation: Some(tier_registry_generation),
},
},
&tier_registry.names,
bucket_plan_complete,
budget_elapsed,
@@ -17,6 +17,33 @@ use crate::data_usage_define::{UNKNOWN_TIER, UnknownTierStats, hash_path};
use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage, TierAccountingProof};
const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]);
const TEST_COVERAGE_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([6; 32]);
#[test]
fn scanner_bucket_inventory_requires_exact_unique_set_union() {
let first = BucketInfo {
name: "first".to_string(),
..Default::default()
};
let second = BucketInfo {
name: "second".to_string(),
..Default::default()
};
let source = DataUsageCacheSource::new(0, 0);
let mut sets = HashMap::from([(source, vec![first.clone()])]);
assert!(scanner_bucket_inventory_is_complete(std::slice::from_ref(&first), &sets));
assert!(!scanner_bucket_inventory_is_complete(&[first.clone(), second.clone()], &sets));
assert!(!scanner_bucket_inventory_is_complete(&[], &sets));
assert!(!scanner_bucket_inventory_is_complete(&[first.clone(), first.clone()], &sets));
sets.insert(source, vec![first.clone(), first.clone()]);
assert!(!scanner_bucket_inventory_is_complete(std::slice::from_ref(&first), &sets));
sets.insert(source, vec![second]);
assert!(!scanner_bucket_inventory_is_complete(std::slice::from_ref(&first), &sets));
let mut recreated = first.clone();
recreated.created = Some(OffsetDateTime::UNIX_EPOCH);
sets.insert(source, vec![recreated]);
assert!(!scanner_bucket_inventory_is_complete(&[first], &sets));
}
#[test]
fn should_publish_completed_snapshot_requires_full_clean_cycle() {
@@ -79,6 +106,7 @@ fn completed_root_cache(bucket: &str, objects: usize, update_secs: u64, source:
source: Some(source),
snapshot_complete: true,
scan_plan_digest: Some(TEST_PLAN_DIGEST),
scan_coverage_digest: Some(TEST_COVERAGE_DIGEST),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
@@ -108,7 +136,139 @@ fn completed_data_usage_info_for_test(
cancelled: bool,
) -> Option<(DataUsageInfo, SystemTime)> {
let expected_sources = results.iter().filter_map(|result| result.info.source).collect::<HashSet<_>>();
completed_data_usage_info(results, &expected_sources, all_buckets, &[], true, budget_elapsed, cancelled)
completed_usage_for_scope(results, &expected_sources, all_buckets, &[], true, budget_elapsed, cancelled)
}
fn completed_usage_for_scope(
results: &[DataUsageCache],
expected_sources: &HashSet<DataUsageCacheSource>,
all_buckets: &[String],
tier_registry_names: &[String],
bucket_plan_complete: bool,
budget_elapsed: bool,
cancelled: bool,
) -> Option<(DataUsageInfo, SystemTime)> {
let first = results.first()?;
completed_data_usage_info(
results,
&ScannerSnapshotScope {
sources: expected_sources,
buckets: all_buckets,
identity: ScannerSnapshotIdentity {
cycle: first.info.next_cycle,
leader_epoch: first.info.leader_epoch,
plan_digest: TEST_PLAN_DIGEST,
coverage_digest: TEST_COVERAGE_DIGEST,
tier_registry_generation: first.info.tier_registry_generation,
},
},
tier_registry_names,
bucket_plan_complete,
budget_elapsed,
cancelled,
)
}
#[test]
fn completed_data_usage_info_rejects_duplicate_bucket_inventory() {
let set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
let buckets = vec!["bucket".to_string(), "bucket".to_string()];
assert!(completed_data_usage_info_for_test(&[set], &buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_extra_or_detached_bucket_data() {
let buckets = vec!["bucket".to_string()];
let mut set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
set.replace(
"unlisted",
DATA_USAGE_ROOT,
DataUsageEntry {
objects: 1,
..Default::default()
},
);
assert!(completed_data_usage_info_for_test(&[set.clone()], &buckets, false, false).is_none());
set.cache
.get_mut(DATA_USAGE_ROOT)
.expect("set root")
.children
.remove(&hash_path("unlisted").key());
assert!(
completed_data_usage_info_for_test(&[set], &buckets, false, false).is_none(),
"orphaned data must not disappear from authoritative accounting"
);
}
#[test]
fn completed_data_usage_info_rejects_disconnected_expected_bucket() {
let buckets = vec!["bucket".to_string()];
let mut set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
set.cache.get_mut(DATA_USAGE_ROOT).expect("set root").children.clear();
assert!(completed_data_usage_info_for_test(&[set], &buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_root_scalar_data_and_unknown_key_format() {
let buckets = vec!["bucket".to_string()];
let set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
let mut scalar_root = set.clone();
scalar_root.cache.get_mut(DATA_USAGE_ROOT).expect("set root").size = 10;
assert!(completed_data_usage_info_for_test(&[scalar_root], &buckets, false, false).is_none());
let mut future_format = set;
future_format.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT + 1;
assert!(completed_data_usage_info_for_test(&[future_format], &buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_binds_all_results_to_requested_identity() {
let buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let sources = HashSet::from([source]);
let set = completed_root_cache("bucket", 2, 10, source);
let identity = ScannerSnapshotIdentity {
cycle: 0,
leader_epoch: 0,
plan_digest: TEST_PLAN_DIGEST,
coverage_digest: TEST_COVERAGE_DIGEST,
tier_registry_generation: None,
};
let results = [set];
for expected in [
ScannerSnapshotIdentity { cycle: 1, ..identity },
ScannerSnapshotIdentity {
leader_epoch: 1,
..identity
},
ScannerSnapshotIdentity {
plan_digest: DataUsageScanPlanDigest([9; 32]),
..identity
},
ScannerSnapshotIdentity {
tier_registry_generation: Some(1),
..identity
},
ScannerSnapshotIdentity {
coverage_digest: DataUsageScanPlanDigest([4; 32]),
..identity
},
] {
let scope = ScannerSnapshotScope {
sources: &sources,
buckets: &buckets,
identity: expected,
};
assert!(completed_data_usage_info(&results, &scope, &[], true, false, false).is_none());
}
let scope = ScannerSnapshotScope {
sources: &sources,
buckets: &buckets,
identity,
};
let (usage, _) = completed_data_usage_info(&results, &scope, &[], true, false, false)
.expect("the requested complete scope remains publishable");
assert_eq!(usage.objects_total_count, 2);
assert!(usage.is_complete_bucket_usage_snapshot());
}
fn lkg_root_cache(bucket: &str, objects: usize, source: DataUsageCacheSource) -> DataUsageCache {
@@ -136,7 +296,7 @@ fn partial_usage_is_observational_not_authoritative_for_quota() {
let expected = HashSet::from([current_source, stalled_source]);
assert!(
completed_data_usage_info(&[current.clone(), stalled.clone()], &expected, &all_buckets, &[], true, false, false)
completed_usage_for_scope(&[current.clone(), stalled.clone()], &expected, &all_buckets, &[], true, false, false)
.is_none()
);
let (observed, _) = observational_data_usage_info(&[current, stalled], &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 8, 3)
@@ -509,7 +669,7 @@ fn completed_data_usage_info_accepts_unknown_only_with_current_registry_generati
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(
completed_data_usage_info(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_some()
completed_usage_for_scope(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_some()
);
}
@@ -547,7 +707,7 @@ fn completed_data_usage_info_rejects_non_registry_tier_in_current_generation() {
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(
completed_data_usage_info(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_none()
completed_usage_for_scope(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_none()
);
}
@@ -619,6 +779,7 @@ fn current_cache_root_with_new_tier_generation_resets_old_cache() {
DataUsageCacheReuseOptions {
require_source: false,
tier_registry_generation: Some(2),
checkpoint_identity: None,
},
);
@@ -698,6 +859,8 @@ fn completed_data_usage_info_publishes_confirmed_empty_namespace() {
source: Some(DataUsageCacheSource::new(0, 0)),
snapshot_complete: true,
scan_plan_digest: Some(TEST_PLAN_DIGEST),
scan_coverage_digest: Some(TEST_COVERAGE_DIGEST),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
..Default::default()
@@ -855,7 +1018,7 @@ fn completed_data_usage_info_requires_exact_topology_sources() {
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0), DataUsageCacheSource::new(1, 0)]);
assert!(
completed_data_usage_info(&[first_set, unexpected_set], &expected_sources, &all_buckets, &[], true, false, false)
completed_usage_for_scope(&[first_set, unexpected_set], &expected_sources, &all_buckets, &[], true, false, false)
.is_none()
);
}
@@ -866,7 +1029,7 @@ fn completed_data_usage_info_rejects_incomplete_bucket_plan() {
let set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(completed_data_usage_info(&[set], &expected_sources, &all_buckets, &[], false, false, false).is_none());
assert!(completed_usage_for_scope(&[set], &expected_sources, &all_buckets, &[], false, false, false).is_none());
}
#[test]
@@ -1176,6 +1339,90 @@ 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 scoped_scan_complete_root_requires_current_coverage_from_every_set() {
let sources = HashSet::from([DataUsageCacheSource::new(0, 0), DataUsageCacheSource::new(1, 0)]);
let buckets = vec!["bucket".to_string()];
let coverage = DataUsageScanPlanDigest([4; 32]);
let scope = ScannerSnapshotScope {
sources: &sources,
buckets: &buckets,
identity: ScannerSnapshotIdentity {
cycle: 0,
leader_epoch: 0,
plan_digest: TEST_PLAN_DIGEST,
coverage_digest: coverage,
tier_registry_generation: None,
},
};
for (first_coverage, second_coverage, valid) in [
(Some(coverage), Some(coverage), true),
(None, Some(coverage), false),
(Some(coverage), None, false),
(None, None, false),
(Some(coverage), Some(DataUsageScanPlanDigest([5; 32])), false),
] {
let mut first = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
let mut second = completed_root_cache("bucket", 3, 10, DataUsageCacheSource::new(1, 0));
first.info.scan_coverage_digest = first_coverage;
second.info.scan_coverage_digest = second_coverage;
assert_eq!(
completed_data_usage_info(&[first, second], &scope, &[], true, false, false).is_some(),
valid
);
}
}
#[test]
fn scanner_cache_lock_resource_is_scoped_to_cache_source() {
let cache_name = "photos/.usage-cache.bin";
+336 -8
View File
@@ -124,6 +124,41 @@ 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, 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, 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, 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, HealScanMode::Normal)
.await
.is_err()
);
}
#[tokio::test]
#[serial]
async fn scanner_cache_locks_block_same_source_workers() {
@@ -283,6 +318,211 @@ async fn scanner_cycle_is_deferred_while_terminal_decommission_is_blocked() {
}
}
#[tokio::test]
#[serial]
async fn scoped_scan_production_entry_preserves_deep_and_full_maintenance_work() {
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 mut baseline = None;
for (index, (scan_mode, requires_full_scan, explicit_scope)) in [
(HealScanMode::Normal, true, false),
(HealScanMode::Normal, false, false),
(HealScanMode::Deep, false, false),
(HealScanMode::Normal, true, false),
(HealScanMode::Deep, false, true),
(HealScanMode::Normal, true, true),
]
.into_iter()
.enumerate()
{
if index > 0 {
let mut reader = ScannerPutObjReader::from_vec(b"maintenance".to_vec());
store.pools[0].disk_set[0]
.put_object("cold-bucket", &format!("added-{index}"), &mut reader, &ScannerObjectOptions::default())
.await
.expect("cold bucket mutation should persist");
// Only the hot bucket is in the usage hint. The cold result must
// come from this cycle's storage walk, not its previous baseline.
record_dirty_usage_bucket("hot-bucket");
}
let requested_scope = if explicit_scope {
ScannerBucketScanScope::from_dirty_buckets(
HashSet::from(["hot-bucket".to_string()]),
DataUsageScanPlanDigest([7; 32]),
)
} else {
ScannerBucketScanScope::default()
};
let ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default());
let (updates, mut receiver) = mpsc::channel(1);
let (observer, observed_scope) = tokio::sync::oneshot::channel();
let cycle = u64::try_from(index + 1).expect("test cycle should fit");
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,
scan_scope: requested_scope,
persisted_usage_baseline: baseline,
requires_full_scan,
resolved_scope_observer: Some(observer),
},
),
)
.await
.expect("cycle should finish within the test deadline")
.expect("cycle should succeed");
assert_eq!(result.status, ScannerCycleStatus::Complete, "cycle {cycle}");
let resolved = observed_scope.await.expect("production resolver should report its scope");
if index == 1 {
assert_eq!(
resolved.selected_buckets.as_deref(),
Some(&HashSet::from(["hot-bucket".to_string()])),
"ordinary dirty work must retain the existing planner"
);
} else {
assert!(resolved.is_default(), "cycle {cycle} must visit the full maintenance scope");
}
let mut snapshot = receiver.recv().await.expect("cycle should publish a snapshot");
assert!(snapshot.usage_snapshot_complete, "cycle {cycle}");
assert_eq!(
snapshot.buckets_usage["cold-bucket"].objects_count,
u64::try_from(index + 1).expect("count should fit")
);
assert_eq!(snapshot.buckets_usage["hot-bucket"].objects_count, 1);
assert_eq!(snapshot.scanner_cycle, Some(cycle));
assert_eq!(snapshot.scanner_epoch, Some(11));
snapshot.usage_snapshot_converged = Some(true);
baseline = Some(Bytes::from(serde_json::to_vec(&snapshot).expect("complete baseline should encode")));
}
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::Normal, false),
(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");
if scan_mode == HealScanMode::Normal && !requires_full_scan {
record_dirty_usage_bucket("cold-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);
@@ -893,6 +1133,7 @@ fn complete_set_usage_cache(buckets: &[(&str, usize)], scan_plan_digest: DataUsa
source: Some(DataUsageCacheSource::new(1, 2)),
snapshot_complete: true,
scan_plan_digest: Some(scan_plan_digest),
scan_coverage_digest: Some(scan_plan_digest),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
tier_registry_generation: Some(13),
..Default::default()
@@ -1010,6 +1251,8 @@ async fn set_snapshot_reuse_requires_execution_identity_and_fences_stale_writers
all_buckets: Arc::new(Vec::new()),
scope: ScannerBucketScanScope::default(),
digest: DataUsageScanPlanDigest([6; 32]),
bucket_coverage_digest: DataUsageScanPlanDigest([6; 32]),
requires_full_scan: false,
execution_digest: empty_execution,
leader_epoch: 11,
tier_registry_generation: 13,
@@ -1139,6 +1382,36 @@ fn scoped_scan_selects_only_current_dirty_buckets_after_baseline_validation() {
assert_ne!(scope.baseline_scan_plan_digest, Some(baseline_scan_plan_digest));
}
#[test]
fn scoped_scan_baseline_work_proof_requires_uniform_known_set_identity() {
let source = DataUsageCacheSource::new(1, 2);
let second_source = DataUsageCacheSource::new(1, 3);
let sources = HashSet::from([source, second_source]);
let structural = DataUsageScanPlanDigest([9; 32]);
let full = scanner_bucket_work_digest(structural, HealScanMode::Normal, true);
let deep = scanner_bucket_work_digest(structural, HealScanMode::Deep, true);
let encoded = complete_usage_baseline(source, full, 7, 11);
let baseline: DataUsageInfo = serde_json::from_slice(&encoded).expect("baseline should decode");
for (second_plan, expected) in [(full, Some(full)), (deep, None), (DataUsageScanPlanDigest([8; 32]), None)] {
let mut candidate = baseline.clone();
let mut second = candidate.usage_snapshot_set_states[0].clone();
second.set_index = 3;
second.scan_plan_digest = Some(second_plan.0);
candidate.usage_snapshot_set_states.push(second);
let data = Bytes::from(serde_json::to_vec(&candidate).expect("candidate should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
data: Some(&data),
expected_sources: &sources,
leader_epoch: 11,
want_cycle: 8,
scan_plan_digest: structural,
}),
expected
);
}
}
fn peer_dirty_usage_snapshot(
instance_id: &str,
generation: u64,
@@ -1221,8 +1494,15 @@ fn verified_remote_dirty_usage_buckets_rejects_incomplete_or_stale_peer_state()
}
}
fn bucket_info_with_created_time(name: &str) -> BucketInfo {
BucketInfo {
created: Some(time::OffsetDateTime::UNIX_EPOCH),
..bucket_info(name)
}
}
#[test]
fn scoped_set_scan_preserves_unselected_usage_and_drops_deleted_buckets() {
fn scoped_set_scan_rebuilds_selected_buckets_and_drops_deleted_buckets() {
let baseline_digest = DataUsageScanPlanDigest([1; 32]);
let current_digest = DataUsageScanPlanDigest([2; 32]);
let mut old_cache = complete_set_usage_cache(&[("stable", 10), ("dirty", 20), ("deleted", 30)], baseline_digest);
@@ -1235,8 +1515,11 @@ fn scoped_set_scan_preserves_unselected_usage_and_drops_deleted_buckets() {
..Default::default()
},
);
let all_buckets = vec![bucket_info("stable"), bucket_info("dirty")];
let selected_buckets = Arc::new(HashSet::from(["dirty".to_string(), "deleted".to_string()]));
let all_buckets = vec![
bucket_info_with_created_time("stable"),
bucket_info_with_created_time("dirty"),
];
let selected_buckets = Arc::new(HashSet::from(["stable".to_string(), "dirty".to_string(), "deleted".to_string()]));
let prepared = prepare_scoped_set_scan(
&old_cache,
@@ -1256,12 +1539,16 @@ fn scoped_set_scan_preserves_unselected_usage_and_drops_deleted_buckets() {
)
.expect("complete matching set cache should support a scoped scan");
assert_eq!(prepared.buckets.iter().map(|bucket| bucket.name.as_str()).collect::<Vec<_>>(), ["dirty"]);
assert_eq!(
prepared.buckets.iter().map(|bucket| bucket.name.as_str()).collect::<Vec<_>>(),
["stable", "dirty"]
);
let stable = prepared
.cache
.checked_flatten("stable")
.expect("unselected bucket subtree should be retained");
assert_eq!((stable.size, stable.objects), (15, 2));
.expect("selected bucket placeholder should exist");
assert_eq!((stable.size, stable.objects), (0, 0));
assert!(prepared.cache.find("stable/prefix").is_none());
assert_eq!(prepared.cache.find("dirty").map(|entry| (entry.size, entry.objects)), Some((0, 0)));
assert!(prepared.cache.find("deleted").is_none());
assert_eq!(prepared.cache.info.scan_plan_digest, Some(current_digest));
@@ -1272,11 +1559,41 @@ fn scoped_set_scan_preserves_unselected_usage_and_drops_deleted_buckets() {
assert_eq!(prepared.cache.info.lkg_scan_plan_digest, Some(baseline_digest));
}
#[test]
fn scoped_set_scan_rejects_unbound_bucket_incarnations() {
let baseline_digest = DataUsageScanPlanDigest([1; 32]);
let old_cache = complete_set_usage_cache(&[("stable", 10), ("dirty", 20)], baseline_digest);
let scope = ScannerBucketScanScope {
selected_buckets: Some(Arc::new(HashSet::from(["dirty".to_string()]))),
baseline_scan_plan_digest: Some(baseline_digest),
};
let generation = ScannerSetCacheGeneration {
want_cycle: 8,
leader_epoch: 11,
tier_registry_generation: 13,
source: DataUsageCacheSource::new(1, 2),
scan_plan_digest: DataUsageScanPlanDigest([2; 32]),
};
for created in [
None,
Some(OffsetDateTime::UNIX_EPOCH),
Some(OffsetDateTime::UNIX_EPOCH + time::Duration::days(1)),
] {
let mut stable = bucket_info("stable");
stable.created = created;
let buckets = vec![stable, bucket_info_with_created_time("dirty")];
assert!(
prepare_scoped_set_scan(&old_cache, &buckets, &buckets, &scope, generation).is_none(),
"missing identity, volume timestamps and same-name recreation must all rebuild"
);
}
}
#[test]
fn scoped_set_scan_falls_back_when_an_unselected_bucket_has_no_baseline() {
let baseline_digest = DataUsageScanPlanDigest([3; 32]);
let old_cache = complete_set_usage_cache(&[("stable", 10)], baseline_digest);
let all_buckets = vec![bucket_info("stable"), bucket_info("new")];
let all_buckets = vec![bucket_info_with_created_time("stable"), bucket_info_with_created_time("new")];
assert!(
prepare_scoped_set_scan(
@@ -1302,7 +1619,7 @@ fn scoped_set_scan_falls_back_when_an_unselected_bucket_has_no_baseline() {
#[test]
fn scoped_set_scan_requires_an_exact_complete_baseline() {
let baseline_digest = DataUsageScanPlanDigest([5; 32]);
let all_buckets = vec![bucket_info("dirty")];
let all_buckets = vec![bucket_info_with_created_time("dirty")];
let scope = ScannerBucketScanScope {
selected_buckets: Some(Arc::new(HashSet::from(["dirty".to_string()]))),
baseline_scan_plan_digest: Some(baseline_digest),
@@ -1323,6 +1640,10 @@ fn scoped_set_scan_requires_an_exact_complete_baseline() {
not_durable.info.last_update = None;
assert!(prepare_scoped_set_scan(&not_durable, &all_buckets, &all_buckets, &scope, generation).is_none());
let mut unscoped_usage = complete_set_usage_cache(&[("dirty", 10)], baseline_digest);
unscoped_usage.cache.get_mut(DATA_USAGE_ROOT).expect("set root").objects = 1;
assert!(prepare_scoped_set_scan(&unscoped_usage, &all_buckets, &all_buckets, &scope, generation).is_none());
let mut wrong_digest = complete_set_usage_cache(&[("dirty", 10)], baseline_digest);
wrong_digest.info.scan_plan_digest = Some(DataUsageScanPlanDigest([7; 32]));
assert!(prepare_scoped_set_scan(&wrong_digest, &all_buckets, &all_buckets, &scope, generation).is_none());
@@ -1333,6 +1654,13 @@ fn scoped_set_scan_requires_an_exact_complete_baseline() {
};
let complete = complete_set_usage_cache(&[("dirty", 10)], baseline_digest);
assert!(prepare_scoped_set_scan(&complete, &all_buckets, &all_buckets, &empty_scope, generation).is_none());
assert!(prepare_scoped_set_scan(&complete, &all_buckets, &all_buckets, &scope, generation).is_some());
let unidentified_buckets = vec![bucket_info("dirty")];
assert!(
prepare_scoped_set_scan(&complete, &unidentified_buckets, &unidentified_buckets, &scope, generation).is_some(),
"fully selected buckets are rebuilt without reusing an unproven incarnation"
);
let mut future_cache = complete_set_usage_cache(&[("dirty", 10)], baseline_digest);
future_cache.info.next_cycle = generation.want_cycle.saturating_add(1);
+10 -6
View File
@@ -9,14 +9,18 @@ cargo test -p rustfs-scanner --lib checkpoint_fixture -- --list
RUST_MIN_STACK=4194304 cargo test -p rustfs-scanner --lib checkpoint_fixture -- --nocapture
```
The unchanged-plan case requires durable static coverage to increase each round. The hot-plan diagnostic changes the bucket plan digest between rounds and reports where coverage is lost without asserting that a particular defect must remain present. To require progress in this diagnostic as well:
Both the unchanged-plan and hot-plan cases require durable static coverage to increase each round. `LostAtPrepare` identifies invalidation before traversal; `LostAtReload` identifies loss between the returned cache and persisted data; `WalkWithoutRetention` identifies visited growth without durable coverage growth. Missing, corrupt, empty-root, and oversized checkpoint inputs are rejected by the strict fixture reader. Save failure and publication-epoch rejection must preserve the preceding file bytes. Parent cancellation is checked separately from object-budget exhaustion. Superseded classification is tested separately from either incomplete outcome.
```sh
RUST_MIN_STACK=4194304 RUSTFS_CHECKPOINT_REQUIRE_PROGRESS=1 cargo test -p rustfs-scanner --lib checkpoint_fixture_hot_digest_diagnostic -- --nocapture
```
After three interrupted rounds, the fixture overwrites a previously visited object with two versions, deletes another visited object, and creates one more hot object. It then keeps the same four-object budget until the stable namespace is certified. Finishing a sweep that spans different mutation plans must first return partial; a subsequent verification sweep must produce exactly 25 objects, 2 versioned entries, and 34 logical bytes. There is no final unbudgeted sweep.
A nonzero exit from the strict command means that walked work did not become additional retained static coverage. `LostAtPrepare` identifies invalidation before traversal; `LostAtReload` identifies loss between the returned cache and persisted data; `WalkWithoutRetention` identifies visited growth without durable coverage growth. Missing, corrupt, empty-root, and oversized checkpoint inputs are rejected by the strict fixture reader. Save failure and publication-epoch rejection must preserve the preceding file bytes. Parent cancellation is checked separately from object-budget exhaustion. Superseded classification is tested separately from either incomplete outcome.
The new bucket checkpoint binds the persisted bucket incarnation, set layout, publication epoch, tier generation and scan mode, with the existing source/leader/key-format checks. Its forward sweep records the starting and requested mutation plans separately. Partial sweeps omit the legacy `scan_plan_digest`, so older readers rebuild instead of treating mixed observations as a current complete snapshot. Completed sweeps restore that digest only after covering one mutation plan. Unsupported or missing identities retain the legacy rebuild path. The stable-plan fast path requires a complete snapshot without unfinished checkpoint state; its interrupted result must enter forward validation on reload.
A coverage receipt binds the completed traversal frontier to its scope, starting plan and canonical covered-prefix digest. It excludes ancestor aggregates and the unvisited suffix, so unrelated suffix changes cannot invalidate completed work. Cancellation seals only the completed frontier; failed child traversal and known failed-metadata skips block further frontier advancement. A saved cursor pointing at an existing but unvisited old subtree is rejected unless it agrees with that receipt. The receipt is a consistency check for storage owned by the scanner, not authentication against a party able to forge the entire cache and recompute its digest. New metadata remains map-encoded with optional top-level fields.
The Normal-to-Deep regression holds the mutation plan fixed, changes metadata in a previously visited prefix, and verifies that the real Deep disk-scan entry point reads that prefix again. It also checks that a complete Normal cache cannot satisfy the Deep `Current` path. The fixture disables heal side effects; it proves traversal re-entry, not actual bitrot detection or repair. Additional tests cover map round-trips, stale identities, existing-but-uncovered cursors, coverage gaps, and per-instance metadata ownership.
This fixture bounds object processing after directory enumeration. It does not prove fixed-budget enumeration of arbitrarily wide directories or real process-restart convergence. Those gates require a storage-owned resumable enumeration capability, including its initial construction cost; a readdir offset, an in-memory iterator or a last-name filter is not that capability.
For every saved partial cache, the fixture also passes its progress through the production authenticated remote terminal-frame writer and stream consumer. A remote partial result must remain partial even when its progress reports visited objects. This covers the return-frame contract; it does not execute the remote RPC server, distributed locks, EC quorum persistence, mixed-version peers, process crashes, or fsync durability. The file backend models revision preconditions and persistence errors, not a concurrent object store.
The synthetic namespace contains no customer data. Temporary files are removed with their owning fixture. Production scan semantics and persistent formats are unchanged, so rollback consists of removing these tests and this guide. A passing fixture alone does not establish that the field report in [issue #7108](https://github.com/rustfs/rustfs/issues/7108) has been independently reproduced or fixed. A field diagnosis must separately identify the source capture, cycle and leader identity, and decoded bucket/set caches.
The synthetic namespace contains no customer data. Temporary files are removed with their owning fixture. Rolling back to a reader without the optional checkpoint metadata rebuilds partial coverage; it must not clear quota floors or complete authoritative snapshots. A passing fixture alone does not establish that the field report in [issue #7108](https://github.com/rustfs/rustfs/issues/7108) has been independently reproduced or fixed. A field diagnosis must separately identify the source capture, cycle and leader identity, and decoded bucket/set caches.
+1 -1
View File
@@ -38,7 +38,6 @@ use crate::on_demand_migration::{
SourceListRequest, SourceObject, SourcePage, decode_continuation_token, source_list_plan,
};
use futures::StreamExt;
use http::HeaderMap;
use rustfs_utils::http::{SUFFIX_SOURCE_PROXY_REQUEST, get_header};
use std::sync::Arc;
use std::time::Instant;
@@ -476,6 +475,7 @@ mod tests {
FilterConfig, MAX_LIST_NO_PROGRESS_PAGES, OnDemandMigrationConfig, PathStyle, PolicyConfig, Provider, SourceConfig,
SourceCredentials, TlsConfig,
};
use http::HeaderMap;
use std::time::Duration;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
@@ -133,6 +133,7 @@ impl NativeHttp {
self.send_classified(request, error_code_header, false).await
}
#[cfg(feature = "gcs")]
pub(super) async fn send_object(
&self,
request: reqwest::Request,
@@ -210,6 +211,7 @@ pub(super) async fn read_text(response: reqwest::Response, max_bytes: usize) ->
/// Base64 digest (`Content-MD5`, `md5Hash`, `x-goog-hash`) as lowercase hex.
/// `None` when the value is not a 16-byte digest, so a CRC32C never passes as
/// an MD5.
#[cfg(any(test, feature = "gcs"))]
pub(super) fn base64_md5_to_hex(value: &str) -> Option<String> {
let raw = base64_simd::STANDARD.decode_to_vec(value.trim().as_bytes()).ok()?;
(raw.len() == 16).then(|| faster_hex::hex_string(&raw))
+1 -1
View File
@@ -4034,7 +4034,7 @@ mod tests {
let enabled_before = sys.is_module_enabled();
sys.set_module_enabled(true);
let bucket = format!("odm-capture-failure-{}", uuid::Uuid::new_v4());
let mut config: crate::on_demand_migration::OnDemandMigrationConfig = serde_json::from_str(r#"{"source":{"provider":"minio","endpoint":"https://source.example.com","bucket":"source","credentials":{"access_key":"test","secret_key":"test"}}}"#).expect("source config");
let mut config: crate::on_demand_migration::OnDemandMigrationConfig = serde_json::from_str(r#"{"source":{"provider":"minio","endpoint":"https://source.example.com","region":"us-east-1","bucket":"source","credentials":{"access_key":"test","secret_key":"test"}}}"#).expect("source config");
config.policy.list_through = true;
sys.apply_for_incarnation(&bucket, uuid::Uuid::new_v4(), Some(&config)).await;
let mut get = build_request(