From 402ac432078c56dc4f913ef3300a90af13c98648 Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 29 Aug 2026 23:35:56 +0800 Subject: [PATCH] fix(scanner): retain post-scan observations Preserve a complete scanner walk as a non-converged observation when the final activity probe is unavailable. Advance the cycle as partial without acknowledging dirty usage.\n\nCo-Authored-By: heihutu --- crates/scanner/src/scanner.rs | 49 ++++++++++++++++++++--- crates/scanner/src/scanner/tests.rs | 20 +++++++++ crates/scanner/src/scanner_io.rs | 15 +++++++ crates/scanner/src/scanner_io/io_cycle.rs | 40 ++++++++++-------- crates/scanner/src/scanner_io/tests.rs | 41 +++++++++++++++++++ 5 files changed, 143 insertions(+), 22 deletions(-) diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index d985931ad..3e4f310d2 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -1462,6 +1462,20 @@ async fn run_data_scanner_cycle_with_budget( { Some(ScannerCycleDeferReason::DataMovement) } + // A complete walk can still be retained as an observational snapshot + // when only the final activity proof was unavailable. It must not + // block the observation receiver: the authoritative publication + // fence remains enforced by the usage store and the cycle is advanced + // as partial without acknowledging dirty usage. + Ok(result) + if result.has_observational_snapshot() + && matches!( + result.status, + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable) + ) => + { + None + } Ok(result) => final_data_usage_publication_defer_reason(storeapi.as_ref(), result.status).await, Err(_) => Some(ScannerCycleDeferReason::ActivityBaselineUnavailable), }; @@ -2803,12 +2817,7 @@ fn finalize_scanner_cycle_result( scan_cycle_result: crate::scanner_io::ScannerCycleResult, usage_persist_outcome: DataUsagePersistOutcome, ) -> (ScannerCycleOutcome, bool, Vec) { - let completion_outcome = scanner_cycle_completion_outcome( - scan_cycle_result.status, - usage_persist_outcome, - scan_cycle_result.has_dirty_usage_to_acknowledge(), - scan_cycle_result.has_failed_dirty_usage(), - ); + let completion_outcome = scanner_cycle_completion_outcome_for_result(&scan_cycle_result, usage_persist_outcome); let pending_maintenance_work = scan_cycle_result.has_pending_maintenance_work(); let durable_complete_snapshot = scan_cycle_result.status == ScannerCycleStatus::Complete && matches!( @@ -2823,6 +2832,34 @@ fn finalize_scanner_cycle_result( (completion_outcome, pending_maintenance_work, remote_dirty_usage_acknowledgements) } +fn scanner_cycle_completion_outcome_for_result( + scan_cycle_result: &crate::scanner_io::ScannerCycleResult, + usage_persist_outcome: DataUsagePersistOutcome, +) -> ScannerCycleOutcome { + let has_dirty_usage = scan_cycle_result.has_dirty_usage_to_acknowledge(); + let has_failed_dirty_usage = scan_cycle_result.has_failed_dirty_usage(); + if scan_cycle_result.has_observational_snapshot() + && matches!( + scan_cycle_result.status, + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable) + ) + { + return match usage_persist_outcome { + DataUsagePersistOutcome::Saved + | DataUsagePersistOutcome::AlreadyDurable + | DataUsagePersistOutcome::PriorCycleDurable + | DataUsagePersistOutcome::Current + if !has_failed_dirty_usage => + { + ScannerCycleOutcome::Partial + } + DataUsagePersistOutcome::Deferred(reason) => ScannerCycleOutcome::Deferred(reason), + _ => ScannerCycleOutcome::Failed, + }; + } + scanner_cycle_completion_outcome(scan_cycle_result.status, usage_persist_outcome, has_dirty_usage, has_failed_dirty_usage) +} + /// Decide whether an incoming usage snapshot must be skipped as stale, given the local /// wall clock `now`. Mirrors `stale_data_usage_persist_reason` in /// `crates/ecstore/src/data_usage/mod.rs` — keep the two consistent. diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 08c882841..b4af5fbda 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -4478,6 +4478,26 @@ fn finalizing_a_deferred_usage_save_keeps_dirty_work_pending() { crate::scanner_io::clear_dirty_usage_bucket("photos"); } +#[test] +#[serial] +fn finalizing_post_scan_observation_advances_partially_without_dirty_ack() { + crate::scanner_io::clear_dirty_usage_bucket("photos"); + crate::scanner_io::record_dirty_usage_bucket("photos"); + let dirty_snapshot = crate::scanner_io::dirty_usage_buckets_for_tests(); + let observed = crate::scanner_io::ScannerCycleResult::new( + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable), + Some(dirty_snapshot), + ) + .with_observational_snapshot_published(true); + + let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(observed, DataUsagePersistOutcome::Saved); + + assert_eq!(outcome, ScannerCycleOutcome::Partial); + assert!(acknowledgements.is_empty()); + assert!(crate::scanner_io::dirty_usage_buckets_pending()); + crate::scanner_io::clear_dirty_usage_bucket("photos"); +} + #[tokio::test] async fn scanner_cycle_keeps_remote_pending_acknowledgement() { let pending = remote_dirty_usage_acknowledgement_pending(7, 1, std::future::ready(Ok::(true))).await; diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 8753b2139..be39397a6 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -269,6 +269,10 @@ fn should_publish_usage_snapshot(status: ScannerCycleStatus) -> bool { matches!(status, ScannerCycleStatus::Complete | ScannerCycleStatus::Superseded) } +fn should_publish_observational_snapshot(status: ScannerCycleStatus) -> bool { + matches!(status, ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable)) +} + fn prepare_usage_snapshot_for_publication( status: ScannerCycleStatus, mut data_usage_info: DataUsageInfo, @@ -611,6 +615,7 @@ fn scanner_activity_preflight( pub(crate) struct ScannerCycleResult { pub(crate) status: ScannerCycleStatus, publication_epoch: Option, + observational_snapshot_published: bool, dirty_usage_clear: Option, remote_dirty_usage_acknowledgements: Vec, remote_publication_lease_targets: Vec<(String, String, u64)>, @@ -624,6 +629,7 @@ impl ScannerCycleResult { Self { status, publication_epoch: None, + observational_snapshot_published: false, dirty_usage_clear, remote_dirty_usage_acknowledgements: Vec::new(), remote_publication_lease_targets: Vec::new(), @@ -642,6 +648,15 @@ impl ScannerCycleResult { self.publication_epoch } + pub(crate) fn with_observational_snapshot_published(mut self, published: bool) -> Self { + self.observational_snapshot_published = published; + self + } + + pub(crate) fn has_observational_snapshot(&self) -> bool { + self.observational_snapshot_published + } + fn with_failed_dirty_usage(mut self, failed_dirty_usage: bool) -> Self { self.failed_dirty_usage = failed_dirty_usage; self diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index a5c447be3..ce9236277 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -162,18 +162,18 @@ impl ScannerIOCycle for ECStore { dirty_usage_status, activity_status, ); - if !publish_usage_snapshot( - &updates, - status, - DataUsageInfo { - last_update: Some(SystemTime::now()), - scanner_cycle: Some(want_cycle), - usage_snapshot_complete: true, - ..Default::default() - }, - ) - .await? - { + let empty_usage = DataUsageInfo { + last_update: Some(SystemTime::now()), + scanner_cycle: Some(want_cycle), + usage_snapshot_complete: true, + ..Default::default() + }; + let observational_snapshot_published = if should_publish_observational_snapshot(status) { + publish_observational_snapshot(&updates, empty_usage).await? + } else { + publish_usage_snapshot(&updates, status, empty_usage).await? + }; + if !observational_snapshot_published { return Ok(ScannerCycleResult::new(status, None).with_publication_epoch(publication_epoch)); } if status == ScannerCycleStatus::Complete { @@ -188,6 +188,7 @@ impl ScannerIOCycle for ECStore { }; return Ok(ScannerCycleResult::new(status, dirty_usage_clear) .with_publication_epoch(publication_epoch) + .with_observational_snapshot_published(observational_snapshot_published) .with_remote_publication_lease_targets(remote_publication_lease_targets) .with_remote_dirty_usage_acknowledgements(remote_dirty_usage_acknowledgements)); } @@ -437,13 +438,19 @@ impl ScannerIOCycle for ECStore { dirty_usage_status, activity_status, ); - if let Some((data_usage_info, _)) = completed_usage { - publish_usage_snapshot(&updates, cycle_status, data_usage_info).await?; + let observational_snapshot_published = if let Some((data_usage_info, _)) = completed_usage { + if should_publish_observational_snapshot(cycle_status) { + publish_observational_snapshot(&updates, data_usage_info).await? + } else { + publish_usage_snapshot(&updates, cycle_status, data_usage_info).await? + } } else if !ctx.is_cancelled() && let Some((data_usage_info, _)) = observational_usage { - publish_observational_snapshot(&updates, data_usage_info).await?; - } + publish_observational_snapshot(&updates, data_usage_info).await? + } else { + false + }; let dirty_usage_clear = should_clear_dirty_usage_snapshot( result.is_ok(), structurally_complete_snapshot, @@ -463,6 +470,7 @@ impl ScannerIOCycle for ECStore { }; Ok(ScannerCycleResult::new(cycle_status, dirty_usage_clear) .with_publication_epoch(publication_epoch) + .with_observational_snapshot_published(observational_snapshot_published) .with_remote_publication_lease_targets(remote_publication_lease_targets) .with_remote_dirty_usage_acknowledgements(remote_dirty_usage_acknowledgements) .with_failed_dirty_usage(!failed_buckets.is_empty()) diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index ef0e98063..df86fdc1d 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -909,6 +909,47 @@ async fn structurally_complete_superseded_cycles_publish_without_claiming_conver ); } +#[tokio::test] +async fn post_scan_activity_failure_retains_complete_usage_as_observation() { + let (updates, mut receiver) = mpsc::channel(1); + let status = ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable); + + assert!(should_publish_observational_snapshot(status)); + assert!( + publish_observational_snapshot( + &updates, + DataUsageInfo { + last_update: Some(SystemTime::now()), + scanner_cycle: Some(7), + objects_total_count: 3, + objects_total_size: 12, + usage_snapshot_complete: true, + ..Default::default() + }, + ) + .await + .expect("post-scan activity failure should retain an observation") + ); + + let observed = receiver.recv().await.expect("observational update should be queued"); + assert!(!observed.usage_snapshot_complete); + assert!(observed.usage_snapshot_partial); + assert_eq!(observed.usage_snapshot_converged, Some(false)); + assert_eq!(observed.objects_total_count, 3); + assert_eq!(observed.objects_total_size, 12); +} + +#[test] +fn only_unverified_activity_allows_post_scan_observation() { + assert!(should_publish_observational_snapshot(ScannerCycleStatus::Deferred( + ScannerCycleDeferReason::ActivityBaselineUnavailable + ))); + assert!(!should_publish_observational_snapshot(ScannerCycleStatus::Deferred( + ScannerCycleDeferReason::DataMovement + ))); + assert!(!should_publish_observational_snapshot(ScannerCycleStatus::Incomplete)); +} + #[test] fn scanner_cycle_fails_closed_for_namespace_disappearance() { for activity_status in [ScannerCycleActivityStatus::Changed, ScannerCycleActivityStatus::Unchanged] {