mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-05 19:55:37 +00:00
Compare commits
33 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 982ed2888f | |||
| e1608fbd9c | |||
| eb1b17802c | |||
| eec6d3a76d | |||
| 112f70914d | |||
| bbca907781 | |||
| 51982be687 | |||
| b8f0119320 | |||
| 20ca988e85 | |||
| 045770d09d | |||
| 278cbaefd5 | |||
| 982bebf6a5 | |||
| 1aafe477ec | |||
| c33a42c5ab | |||
| e3fd1ae19a | |||
| 6c1d401215 | |||
| c80ff445c0 | |||
| 3d4c05ddf9 | |||
| f9dd0388fa | |||
| 527860d71e | |||
| 7d2c120073 | |||
| 489c3276a4 | |||
| b09a4f3706 | |||
| a0c32e7c6b | |||
| d5a1530409 | |||
| e7b0349a4f | |||
| bdbdca07c8 | |||
| 53efaa2b8f | |||
| ec9672a397 | |||
| cee35f7e54 | |||
| ef7e7afd8c | |||
| 652ebb12c6 | |||
| 38c03d9d5d |
@@ -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(())
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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;
|
||||
};
|
||||
|
||||
@@ -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))]);
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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) => {
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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(¬_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);
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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(
|
||||
|
||||
Reference in New Issue
Block a user