test(scanner): add durable checkpoint diagnostics (#7175)

* chore(deps): refresh SDKs and pin clock skew regression coverage

Refresh compatible dependencies for Scanner/Heal V2 batch 1 and verify
the production S3 retry/signing path with a deterministic clock.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* test(scanner): add durable checkpoint diagnostics

Refs rustfs/backlog#2260 and rustfs/backlog#2240.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-05 15:34:29 +08:00
committed by GitHub
parent cbfd5b92f4
commit 2e4ab045b6
12 changed files with 791 additions and 109 deletions
+1
View File
@@ -244,6 +244,7 @@ windows-sys = { workspace = true, features = [
windows-sys = { workspace = true, features = ["Win32_System_Ioctl"] }
[dev-dependencies]
aws-smithy-async.workspace = true
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "test-util", "fs"] }
criterion = { workspace = true, features = ["html_reports"] }
temp-env = { workspace = true, features = ["async_closure"] }
+170 -1
View File
@@ -652,9 +652,10 @@ async fn build_aws_s3_http_client_from_tls_path() -> Option<SharedHttpClient> {
#[cfg(test)]
mod tests {
use super::*;
use aws_smithy_async::time::TimeSource;
use aws_smithy_runtime_api::http::StatusCode as SmithyStatusCode;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
fn spec(endpoint: &str, secure: bool) -> RemoteS3EndpointSpec {
RemoteS3EndpointSpec {
@@ -824,6 +825,174 @@ mod tests {
);
}
#[derive(Clone, Debug)]
struct ClockSkewTimeSource(Arc<AtomicU64>);
impl TimeSource for ClockSkewTimeSource {
fn now(&self) -> SystemTime {
SystemTime::UNIX_EPOCH + Duration::from_secs(self.0.load(Ordering::SeqCst))
}
}
#[derive(Clone, Debug)]
struct ClockSkewConnector {
request_headers: RecordedHeaders,
error_code: &'static str,
skew_seconds: i64,
clock: ClockSkewTimeSource,
}
fn recorded_header<'a>(headers: &'a [(String, String)], name: &str) -> &'a str {
headers
.iter()
.find(|(key, _)| key.eq_ignore_ascii_case(name))
.map(|(_, value)| value.as_str())
.unwrap_or_else(|| panic!("signed request must contain {name}"))
}
fn signing_time(headers: &[(String, String)]) -> chrono::NaiveDateTime {
chrono::NaiveDateTime::parse_from_str(recorded_header(headers, "x-amz-date"), "%Y%m%dT%H%M%SZ")
.expect("SDK signing timestamp must use the SigV4 format")
}
impl SmithyHttpConnector for ClockSkewConnector {
fn call(&self, request: HttpRequest) -> HttpConnectorFuture {
let mut headers = self.request_headers.lock().expect("clock skew request capture lock");
assert!(headers.len() < 3, "clock skew fixture must not exceed two GET attempts and one HEAD");
headers.push(
request
.headers()
.iter()
.map(|(key, value)| (key.to_string(), value.to_string()))
.collect(),
);
let server_time = chrono::DateTime::<chrono::Utc>::from(self.clock.now()).naive_utc()
+ chrono::Duration::seconds(self.skew_seconds);
let (status, body) = if headers.len() == 1 {
(
403,
format!("<Error><Code>{}</Code><Message>Clock skew fixture</Message></Error>", self.error_code),
)
} else {
(200, String::new())
};
let response = http::Response::builder()
.status(status)
.header("date", server_time.format("%a, %d %b %Y %H:%M:%S GMT").to_string())
.header("content-type", "application/xml")
.header("content-length", body.len())
.body(SdkBody::from(body))
.expect("clock skew fixture response");
HttpConnectorFuture::ready(Ok(HttpResponse::try_from(response).expect("Smithy fixture response")))
}
}
async fn clock_skew_client(
error_code: &'static str,
skew_seconds: i64,
retry: RemoteS3RetryPolicy,
) -> (S3Client, RecordedHeaders, ClockSkewTimeSource) {
let headers: RecordedHeaders = Arc::new(Mutex::new(Vec::new()));
let clock = ClockSkewTimeSource(Arc::new(AtomicU64::new(1_700_000_000)));
let connector = SharedHttpConnector::new(ClockSkewConnector {
request_headers: Arc::clone(&headers),
error_code,
skew_seconds,
clock: clock.clone(),
});
let mut spec = spec("s3.example.com", true);
spec.retry = retry;
let config = build_remote_s3_config(&spec)
.await
.expect("clock skew fixture uses the production outbound configuration")
.http_client(http_client_fn(move |_settings, _components| connector.clone()))
.time_source(clock.clone())
.build();
(S3Client::from_conf(config), headers, clock)
}
#[tokio::test(start_paused = true)]
async fn remote_s3_clock_skew_retries_resign_and_seed_next_operation() {
for error_code in ["RequestTimeTooSkewed", "SignatureDoesNotMatch"] {
for skew_seconds in [-600, 600] {
let (client, headers, clock) = clock_skew_client(error_code, skew_seconds, REPLICATION_TARGET_RETRY_POLICY).await;
let initial = chrono::DateTime::<chrono::Utc>::from(clock.now()).naive_utc();
client
.get_object()
.bucket("bucket")
.key("object")
.send()
.await
.expect("clock skew GET must retry successfully");
assert_eq!(
headers.lock().expect("captured requests").len(),
2,
"{error_code}: GET needs exactly one retry"
);
clock.0.fetch_add(17, Ordering::SeqCst);
// SDK signing time is independent of Tokio's retry/scheduler clock.
tokio::time::advance(Duration::from_secs(61)).await;
client
.head_bucket()
.bucket("bucket")
.send()
.await
.expect("subsequent HEAD must use the client's cached skew");
let headers = headers.lock().expect("captured signed requests");
assert_eq!(headers.len(), 3, "subsequent operation must succeed on its first attempt");
assert_eq!(signing_time(&headers[0]), initial, "the first attempt must use the injected clock");
assert_eq!(
signing_time(&headers[1]),
initial + chrono::Duration::seconds(skew_seconds),
"{error_code}: retry must apply the measured offset exactly"
);
assert_eq!(
signing_time(&headers[2]),
initial + chrono::Duration::seconds(skew_seconds + 17),
"{error_code}: the next operation must apply cached skew to the advanced signing clock"
);
let signature = |index: usize| {
recorded_header(&headers[index], "authorization")
.rsplit_once("Signature=")
.expect("SigV4 authorization contains a signature")
.1
};
assert_ne!(
signature(0),
signature(1),
"{error_code}: retry must be signed again after adjusting its date"
);
}
}
}
#[tokio::test(start_paused = true)]
async fn remote_s3_clock_skew_respects_one_attempt_policy() {
use aws_smithy_types::error::metadata::ProvideErrorMetadata;
for error_code in ["RequestTimeTooSkewed", "SignatureDoesNotMatch"] {
for retry in [
RemoteS3RetryPolicy::Disabled,
RemoteS3RetryPolicy::Standard { max_attempts: 1 },
] {
let (client, headers, _clock) = clock_skew_client(error_code, 600, retry).await;
let error = client
.get_object()
.bucket("bucket")
.key("object")
.send()
.await
.expect_err("clock skew must not override the caller's one-attempt budget");
assert_eq!(error.as_service_error().and_then(ProvideErrorMetadata::code), Some(error_code));
assert_eq!(
headers.lock().expect("captured requests").len(),
1,
"{error_code}: {retry:?} must send exactly one request"
);
}
}
}
#[test]
fn path_style_auto_and_path_force_path_style() {
assert!(PathStyle::Auto.force_path_style());
+3
View File
@@ -52,6 +52,9 @@ static REMOTE_SCANNER_CYCLE_REFRESH: LazyLock<AsyncMutex<()>> = LazyLock::new(||
mod stream;
#[cfg(test)]
pub(crate) use stream::checkpoint_fixture_partial_return;
pub use stream::{RemoteScannerAdmission, RemoteScannerRequest, serve_remote_scanner_request};
pub(crate) use stream::{RemoteScannerOutcome, RemoteScannerScanSpec, scan_remote_bucket};
use stream::{RemoteScannerReplayCache, RemoteScannerRequestWire, RemoteScannerValidatedCycle};
@@ -1017,6 +1017,48 @@ fn finish_remote_scanner_stream(
#[cfg(test)]
const TEST_NEXT_CYCLE: u64 = 11;
#[cfg(test)]
pub(crate) async fn checkpoint_fixture_partial_return(progress: (u64, u64), entries_visited: u64) {
let request_id = Uuid::new_v4();
let writer_auth = FrameAuthenticator::for_test(request_id);
let reader_auth = FrameAuthenticator::for_test(request_id);
let mut bytes = Vec::new();
write_frame(
&mut bytes,
&writer_auth,
&mut 0,
&RemoteScannerFrame::terminal(
RemoteScannerProgress {
objects_scanned: progress.0,
directories_started: progress.1,
entries_visited,
},
RemoteScannerFrameResult::Partial,
),
)
.await
.expect("checkpoint partial frame must encode");
let frame = read_frame(&mut std::io::Cursor::new(bytes.as_slice()), &reader_auth, &mut 0)
.await
.expect("checkpoint progress frame must authenticate");
assert_eq!(frame.progress.entries_visited, entries_visited);
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(&parent, Default::default());
let result = consume_remote_scanner_stream(
std::io::Cursor::new(bytes),
parent,
budget.clone(),
"bucket",
DataUsageCacheSource::new(0, 0),
DataUsageScanPlanDigest([17; 32]),
reader_auth,
)
.await
.expect("checkpoint partial frame must decode");
assert!(matches!(result, RemoteScannerOutcome::Partial));
assert_eq!(budget.progress(), progress);
}
#[cfg(test)]
async fn consume_remote_scanner_stream<R>(
reader: R,
@@ -24,6 +24,8 @@ use std::io::Write;
use std::os::unix::fs::{PermissionsExt, symlink};
use std::sync::Mutex;
mod checkpoint_fixture;
/// Reset the process-global alert cooldown map; test-only.
fn reset_alert_cooldowns() {
*SCANNER_ALERT_EMISSION_COOLDOWN
@@ -0,0 +1,410 @@
// Copyright 2026 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use super::*;
use crate::scanner_budget::ScannerCycleBudgetConfig;
use crate::scanner_io::{ScannerDiskScanOutcome, ScannerIODisk};
use crate::storage_api::scanner_io::ObjectIO;
use crate::{DataUsageCacheSource, DataUsageScanPlanDigest};
use std::io::Cursor;
use tokio::io::AsyncReadExt;
const CACHE_NAME: &str = "bucket/checkpoint-fixture.bin";
const STATIC_OBJECTS: u64 = 24;
const MAX_CACHE_BYTES: u64 = 1024 * 1024;
const SOURCE: DataUsageCacheSource = DataUsageCacheSource::new(0, 0);
const PLAN: DataUsageScanPlanDigest = DataUsageScanPlanDigest([17; 32]);
/// Real cache persistence codec and CAS calls, backed by two bounded local files.
#[derive(Debug)]
struct FixtureStore {
root: tempfile::TempDir,
reject_save: AtomicBool,
}
impl FixtureStore {
fn new() -> Arc<Self> {
Arc::new(Self {
root: tempfile::tempdir().expect("checkpoint fixture storage directory"),
reject_save: AtomicBool::new(false),
})
}
fn path(&self, object: &str) -> std::path::PathBuf {
assert!(object.ends_with(CACHE_NAME) || object.ends_with(&format!("{CACHE_NAME}.bkp")));
self.root
.path()
.join(if object.ends_with(".bkp") { "backup" } else { "main" })
}
async fn strict_load(&self) -> DataUsageCache {
let bytes = tokio::fs::read(self.root.path().join("main"))
.await
.expect("saved checkpoint fixture must exist");
decode_fixture(&bytes).expect("saved checkpoint fixture must contain a valid bucket root")
}
}
#[async_trait::async_trait]
impl ObjectIO for FixtureStore {
type Error = crate::EcstoreError;
type RangeSpec = crate::storage_api::scanner_io::HTTPRangeSpec;
type HeaderMap = http::HeaderMap;
type ObjectOptions = crate::ScannerObjectOptions;
type ObjectInfo = crate::ScannerObjectInfo;
type GetObjectReader = crate::ScannerGetObjectReader;
type PutObjectReader = crate::ScannerPutObjReader;
async fn get_object_reader(
&self,
_bucket: &str,
object: &str,
_range: Option<Self::RangeSpec>,
_headers: Self::HeaderMap,
_options: &Self::ObjectOptions,
) -> crate::EcstoreResult<Self::GetObjectReader> {
let bytes = tokio::fs::read(self.path(object)).await.map_err(|error| {
if error.kind() == std::io::ErrorKind::NotFound {
crate::EcstoreError::FileNotFound
} else {
crate::EcstoreError::from(error)
}
})?;
assert!(u64::try_from(bytes.len()).expect("cache length") <= MAX_CACHE_BYTES);
Ok(crate::ScannerGetObjectReader {
stream: Box::new(Cursor::new(bytes)),
object_info: crate::ScannerObjectInfo {
etag: Some("fixture".into()),
..Default::default()
},
buffered_body: None,
body_source: Default::default(),
})
}
async fn put_object(
&self,
_bucket: &str,
object: &str,
data: &mut Self::PutObjectReader,
options: &Self::ObjectOptions,
) -> crate::EcstoreResult<Self::ObjectInfo> {
if self.reject_save.load(Ordering::SeqCst) {
return Err(crate::EcstoreError::PreconditionFailed);
}
let path = self.path(object);
let exists = tokio::fs::try_exists(&path).await?;
let preconditions = options.http_preconditions.as_ref().expect("checkpoint writes must use CAS");
if (exists && preconditions.if_none_match_value() == Some("*"))
|| (!exists && preconditions.if_match_value().is_some())
|| (exists && preconditions.if_match_value() != Some("fixture"))
{
return Err(crate::EcstoreError::PreconditionFailed);
}
let mut bytes = Vec::new();
(&mut data.stream).take(MAX_CACHE_BYTES + 1).read_to_end(&mut bytes).await?;
assert!(u64::try_from(bytes.len()).expect("cache length") <= MAX_CACHE_BYTES);
tokio::fs::write(path, bytes).await?;
Ok(crate::ScannerObjectInfo {
etag: Some("fixture".into()),
..Default::default()
})
}
}
#[async_trait::async_trait]
impl crate::ScannerConfigObjectDelete for FixtureStore {
async fn delete_config_object(
&self,
_bucket: &str,
_object: &str,
_options: crate::ScannerObjectOptions,
) -> crate::EcstoreResult<crate::ScannerObjectInfo> {
Err(crate::EcstoreError::NotImplemented)
}
async fn scanner_data_usage_publication_admission(&self) -> Option<crate::ScannerDataUsagePublicationAdmission> {
Some(crate::ScannerDataUsagePublicationAdmission::unfenced())
}
}
fn decode_fixture(bytes: &[u8]) -> Result<DataUsageCache, &'static str> {
if bytes.is_empty() || bytes.len() > usize::try_from(MAX_CACHE_BYTES).expect("fixture bound") {
return Err("missing or oversized checkpoint fixture");
}
let cache = DataUsageCache::unmarshal(bytes).map_err(|_| "corrupt checkpoint fixture")?;
if cache.info.name != "bucket" || cache.checked_flatten("bucket").is_none() {
return Err("checkpoint fixture has no valid bucket root");
}
Ok(cache)
}
fn retained(cache: &DataUsageCache) -> u64 {
assert!(
!cache.root().is_some_and(|root| root.compacted),
"a compacted bucket root cannot prove static-prefix coverage"
);
cache
.checked_flatten("bucket/static")
.map_or(0, |entry| u64::try_from(entry.objects).expect("fixture object count fits u64"))
}
#[derive(Debug, PartialEq, Eq)]
enum CoverageDiagnosis {
Progress,
NoNewWork,
LostAtPrepare,
LostAtReload,
WalkWithoutRetention,
}
fn diagnose(previous: u64, prepared: u64, walked: u64, scanned: u64, reloaded: u64) -> CoverageDiagnosis {
if reloaded < scanned {
CoverageDiagnosis::LostAtReload
} else if prepared < previous {
CoverageDiagnosis::LostAtPrepare
} else if walked > 0 && reloaded <= previous {
CoverageDiagnosis::WalkWithoutRetention
} else if reloaded > previous {
CoverageDiagnosis::Progress
} else {
CoverageDiagnosis::NoNewWork
}
}
#[test]
fn checkpoint_fixture_diagnosis_rejects_walk_without_retention() {
assert_eq!(diagnose(4, 4, 9, 8, 8), CoverageDiagnosis::Progress);
assert_eq!(diagnose(4, 4, 9, 4, 4), CoverageDiagnosis::WalkWithoutRetention);
assert_eq!(diagnose(4, 0, 9, 4, 4), CoverageDiagnosis::LostAtPrepare);
assert_eq!(diagnose(4, 4, 9, 8, 4), CoverageDiagnosis::LostAtReload);
assert_eq!(diagnose(4, 4, 0, 4, 4), CoverageDiagnosis::NoNewWork);
}
#[test]
fn checkpoint_fixture_missing_and_corrupt_inputs_fail() {
for bytes in [
vec![],
vec![0xc1],
DataUsageCache::default().marshal_msg().expect("empty cache encoding"),
vec![0; usize::try_from(MAX_CACHE_BYTES + 1).expect("oversized fixture")],
] {
assert!(decode_fixture(&bytes).is_err(), "invalid fixture must not become an empty complete root");
}
}
#[test]
fn checkpoint_fixture_compaction_preserves_aggregate_not_child_enumeration() {
let mut cache = DataUsageCache::default();
cache.info.name = "bucket".to_string();
cache.replace("bucket", "", DataUsageEntry::default());
cache.replace("bucket/static", "bucket", DataUsageEntry::default());
for index in 0..4 {
cache.replace(
&format!("bucket/static/{index}"),
"bucket/static",
DataUsageEntry {
objects: 1,
..Default::default()
},
);
}
cache.reduce_children_of(&hash_path("bucket/static"), 1, true);
let decoded = decode_fixture(&cache.marshal_msg().expect("encode compacted cache")).expect("decode compacted fixture");
let entry = decoded
.find("bucket/static")
.expect("compaction must retain the static subtree root");
assert!(entry.compacted);
assert!(entry.children.is_empty());
assert_eq!(
retained(&decoded),
4,
"compaction retains aggregate coverage even when leaf keys are absent"
);
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_save_reload_resume() {
run_checkpoint_fixture(false).await;
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_hot_digest_diagnostic() {
run_checkpoint_fixture(true).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;
}
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;
let mut cache = DataUsageCache::default();
let revisions = cache
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("load checkpoint revisions");
if round > 0 {
assert_eq!(retained(&store.strict_load().await), previous);
}
let plan = crate::scanner_io::checkpoint_fixture_bucket_digest(PLAN, change_digest.then_some(u64::from(round)));
crate::scanner_io::current_cache_root_or_prepare_with_generation(
&mut cache,
"bucket",
SOURCE,
11,
7,
plan,
crate::scanner_io::DataUsageCacheReuseOptions {
require_source: true,
tier_registry_generation: None,
},
);
let prepared = retained(&cache);
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("budgeted local disk scan returns partial cache");
let ScannerDiskScanOutcome::Partial(cache) = outcome else {
panic!("budgeted fixture must remain partial")
};
assert!(!cache.info.snapshot_complete, "partial must never publish a complete root");
assert_eq!(budget.reason(), Some(crate::scanner_budget::ScannerCycleBudgetReason::Objects));
let scanned = retained(&cache);
cache
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect("persist partial checkpoint");
let mut loaded = DataUsageCache::default();
loaded
.load(store.clone(), CACHE_NAME)
.await
.expect("reload persisted partial checkpoint");
let reloaded = retained(&loaded);
assert_eq!(reloaded, retained(&store.strict_load().await));
assert_eq!(scanned, reloaded, "save/load must retain static subtree coverage");
assert!(!loaded.info.snapshot_complete);
visited += budget.entries_visited();
let diagnosis = diagnose(previous, prepared, budget.entries_visited(), scanned, reloaded);
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"
);
}
crate::remote_scanner::checkpoint_fixture_partial_return(budget.progress(), budget.entries_visited()).await;
previous = reloaded;
}
assert!(visited > 0, "fixture must exercise the directory walk");
assert!(previous > 0, "fixture must retain and enumerate static subtree entries");
let mut loaded = DataUsageCache::default();
let revisions = loaded
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("load final checkpoint");
let before = tokio::fs::read(store.root.path().join("main"))
.await
.expect("read durable checkpoint bytes");
let epoch_error = loaded
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 1)
.await
.expect_err("stale publication epoch must reject persistence");
assert!(epoch_error.to_string().contains(crate::SCANNER_PUBLICATION_EPOCH_CHANGED));
store.reject_save.store(true, Ordering::SeqCst);
loaded.info.next_cycle += 1;
loaded
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect_err("injected save failure must not report durable progress");
assert_eq!(
tokio::fs::read(store.root.path().join("main"))
.await
.expect("read unchanged checkpoint bytes"),
before
);
let parent = CancellationToken::new();
parent.cancel();
let budget = ScannerCycleBudget::new(&parent, Default::default());
let result = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget.clone(),
vec![scanner.local_disk.clone()],
loaded.clone(),
None,
HealScanMode::Normal,
)
.await;
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,
)
.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")
);
}
+8
View File
@@ -209,6 +209,14 @@ fn scanner_bucket_cache_digest(
DataUsageScanPlanDigest(hasher.finalize().into())
}
#[cfg(test)]
pub(crate) fn checkpoint_fixture_bucket_digest(
scan_plan_digest: DataUsageScanPlanDigest,
dirty_generation: Option<u64>,
) -> DataUsageScanPlanDigest {
scanner_bucket_cache_digest(scan_plan_digest, dirty_generation)
}
fn finalize_nsscanner_result(results: &[DataUsageCache], first_err: Option<Error>) -> Result<()> {
if results.iter().any(|result| result.info.last_update.is_some()) {
return Ok(());
+21
View File
@@ -1048,6 +1048,27 @@ fn scanner_cycle_status_requires_a_clean_complete_snapshot() {
}
}
#[test]
fn checkpoint_fixture_superseded_is_distinct_from_partial_and_cancel() {
for (budget, cancelled, bucket, expected) in [
(false, false, ScannerBucketScanStatus::Complete, ScannerCycleStatus::Superseded),
(true, false, ScannerBucketScanStatus::Partial, ScannerCycleStatus::Incomplete),
(false, true, ScannerBucketScanStatus::Partial, ScannerCycleStatus::Incomplete),
] {
assert_eq!(
classify_nsscanner_cycle(
true,
budget,
cancelled,
bucket,
DirtyUsageSnapshotStatus::Changed,
ScannerCycleActivityStatus::Unchanged
),
expected,
);
}
}
#[test]
fn unverified_activity_defers_partial_and_floor_cycles() {
let expected = ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable);