From 2d159635eddb9f17ad504410be91e20557eee4fe Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Sat, 5 Sep 2026 16:32:02 +0800 Subject: [PATCH 1/5] feat(scanner): plan dirty bucket cache refreshes (#7146) * feat(scanner): plan dirty bucket cache refreshes * fix(scanner): route peer snapshot through storage boundary --------- Co-authored-by: Henry Guo Co-authored-by: cxymds Co-authored-by: Zhengchao An --- crates/scanner/src/scanner.rs | 23 ++- crates/scanner/src/scanner/activity.rs | 41 ++++ crates/scanner/src/scanner/tests.rs | 38 ++++ crates/scanner/src/scanner_io.rs | 119 ++++++++++- crates/scanner/src/scanner_io/cache.rs | 18 ++ crates/scanner/src/scanner_io/io_cycle.rs | 95 ++++++++- .../src/scanner_io/publish_gate_tests.rs | 22 ++ crates/scanner/src/scanner_io/tests.rs | 190 ++++++++++++++++++ crates/scanner/src/storage_api.rs | 4 +- 9 files changed, 539 insertions(+), 11 deletions(-) diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index fcf67eed0..0f97924a4 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -1703,14 +1703,18 @@ where let (sender, receiver) = mpsc::channel::(1); let done_cycle = Metrics::time(Metric::ScanCycle); - let scan_result = crate::scanner_io::nsscanner_with_storage_status( + let scan_result = crate::scanner_io::nsscanner_with_storage_status_scoped( storeapi.as_ref(), - cycle_budget.token(), - cycle_budget.clone(), - sender, - cycle_info.current, - leader_epoch, - scan_mode, + crate::scanner_io::ScannerCycleRequest { + ctx: cycle_budget.token(), + budget: cycle_budget.clone(), + updates: sender, + want_cycle: cycle_info.current, + leader_epoch, + scan_mode, + scan_scope: crate::scanner_io::ScannerBucketScanScope::default(), + persisted_usage_baseline: usage_persist_baseline.data.clone(), + }, ) .await; let publication_defer_reason = match &scan_result { @@ -3424,10 +3428,13 @@ use cycle_state::*; use leadership::*; use usage_store::*; +#[cfg(test)] +pub(crate) use activity::scanner_activity_snapshot_digest; pub use activity::scanner_topology_digest; pub(crate) use activity::{ ScannerActivitySnapshot, ScannerDirtyUsageAcknowledgement, probe_scanner_activity, scanner_activity_allows_usage_publication, - scanner_activity_publication_lease_targets, scanner_activity_snapshot_digest, scanner_dirty_usage_acknowledgements, + scanner_activity_dirty_usage_state_for_host, scanner_activity_publication_lease_targets, scanner_activity_structural_digest, + scanner_dirty_usage_acknowledgements, }; pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance}; pub use backlog::{ diff --git a/crates/scanner/src/scanner/activity.rs b/crates/scanner/src/scanner/activity.rs index 2f8aed1b5..ffcbc5313 100644 --- a/crates/scanner/src/scanner/activity.rs +++ b/crates/scanner/src/scanner/activity.rs @@ -902,6 +902,7 @@ where observation } +#[cfg(test)] pub(crate) fn scanner_activity_snapshot_digest(snapshot: &ScannerActivitySnapshot) -> [u8; 32] { let mut hasher = Sha256::new(); hasher.update(u64::try_from(snapshot.len()).unwrap_or(u64::MAX).to_be_bytes()); @@ -925,6 +926,30 @@ pub(crate) fn scanner_activity_snapshot_digest(snapshot: &ScannerActivitySnapsho hasher.finalize().into() } +/// Hash the activity inputs that make an existing scanner cache unsafe to +/// reuse. Regular namespace writes and dirty-usage generations are omitted: +/// their affected buckets are tracked separately and may be refreshed from a +/// complete authoritative cache baseline. +pub(crate) fn scanner_activity_structural_digest(snapshot: &ScannerActivitySnapshot) -> [u8; 32] { + let mut hasher = Sha256::new(); + hasher.update(u64::try_from(snapshot.len()).unwrap_or(u64::MAX).to_be_bytes()); + for (host, activity) in snapshot { + let host = host.as_bytes(); + let instance_id = activity.instance_id.as_bytes(); + hasher.update(u64::try_from(host.len()).unwrap_or(u64::MAX).to_be_bytes()); + hasher.update(host); + hasher.update(u64::try_from(instance_id.len()).unwrap_or(u64::MAX).to_be_bytes()); + hasher.update(instance_id); + hasher.update(activity.maintenance_generation.to_be_bytes()); + hasher.update(activity.protocol_version.to_be_bytes()); + hasher.update(activity.topology_digest); + hasher.update([u8::from(activity.data_movement_active)]); + hasher.update(activity.movement_generation.to_be_bytes()); + hasher.update([u8::from(activity.publication_blocked)]); + } + hasher.finalize().into() +} + pub(crate) fn scanner_activity_allows_usage_publication(snapshot: &ScannerActivitySnapshot) -> bool { !snapshot.is_empty() && snapshot.values().all(|activity| { @@ -955,6 +980,22 @@ pub(crate) fn scanner_dirty_usage_acknowledgements(snapshot: &ScannerActivitySna .collect() } +pub(crate) fn scanner_activity_dirty_usage_state_for_host<'a>( + snapshot: &'a ScannerActivitySnapshot, + host: &str, +) -> Option<(&'a str, u64, bool)> { + snapshot + .get(host) + .filter(|_| host != LOCAL_SCANNER_ACTIVITY_NODE) + .map(|activity| { + ( + activity.instance_id.as_str(), + activity.dirty_usage_generation, + activity.dirty_usage_pending, + ) + }) +} + pub fn scanner_topology_digest(storeapi: &ECStore) -> [u8; 32] { let endpoint_pools = storeapi.endpoints(); let mut hasher = Sha256::new(); diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index efe29f612..244e9088f 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -8169,6 +8169,44 @@ fn scanner_activity_snapshot_digest_fences_dirty_usage_state() { assert_ne!(scanner_activity_snapshot_digest(&clean), scanner_activity_snapshot_digest(&pending)); } +#[test] +fn scanner_activity_structural_digest_ignores_regular_bucket_writes() { + let baseline = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]); + let mut written = baseline.clone(); + let activity = written.get_mut("node-2").expect("node should exist"); + activity.namespace_generation = 8; + activity.dirty_usage_generation = 6; + activity.dirty_usage_pending = true; + + assert_ne!(scanner_activity_snapshot_digest(&baseline), scanner_activity_snapshot_digest(&written)); + assert_eq!( + scanner_activity_structural_digest(&baseline), + scanner_activity_structural_digest(&written), + "bucket writes are refreshed through the dirty-bucket scope rather than invalidating every cache" + ); +} + +#[test] +fn scanner_activity_structural_digest_fences_restart_and_maintenance() { + let baseline = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]); + let mut restarted = baseline.clone(); + restarted.get_mut("node-2").expect("node should exist").instance_id = "epoch-b".to_string(); + let mut maintained = baseline.clone(); + maintained + .get_mut("node-2") + .expect("node should exist") + .maintenance_generation = 4; + + assert_ne!( + scanner_activity_structural_digest(&baseline), + scanner_activity_structural_digest(&restarted) + ); + assert_ne!( + scanner_activity_structural_digest(&baseline), + scanner_activity_structural_digest(&maintained) + ); +} + #[test] fn scanner_dirty_usage_acknowledgements_exclude_local_and_clean_nodes() { let snapshot = BTreeMap::from([ diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 4807e11b8..c8353bae9 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -21,6 +21,7 @@ use crate::{ DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, DataUsageSnapshotSetState, ScannerError, SizeSummary, TierStats, }; +use bytes::Bytes; use futures::future::join_all; use metrics::counter; use rand::seq::SliceRandom as _; @@ -54,6 +55,7 @@ use tokio_util::task::AbortOnDropHandle; use tracing::{debug, error, warn}; use crate::ScannerObjectInfo as ObjectInfo; +use crate::storage_api::EcstoreScannerPeerDirtyUsageSnapshot; use crate::storage_api::ScannerStorage; use crate::storage_api::scan::NamespaceLocking as _; use crate::storage_api::scanner_io::{BucketInfo, BucketOptions}; @@ -111,6 +113,121 @@ pub(crate) struct ScannerBucketScanScope { baseline_scan_plan_digest: Option, } +impl ScannerBucketScanScope { + fn is_default(&self) -> bool { + self.selected_buckets.is_none() && self.baseline_scan_plan_digest.is_none() + } + + fn from_dirty_buckets(selected_buckets: HashSet, baseline_scan_plan_digest: DataUsageScanPlanDigest) -> Self { + Self { + selected_buckets: Some(Arc::new(selected_buckets)), + baseline_scan_plan_digest: Some(baseline_scan_plan_digest), + } + } +} + +#[derive(Clone, Copy)] +pub(super) struct ScannerCacheBaselineProof<'a> { + pub(super) data: Option<&'a Bytes>, + pub(super) expected_sources: &'a HashSet, + pub(super) leader_epoch: u64, + pub(super) want_cycle: u64, + pub(super) scan_plan_digest: DataUsageScanPlanDigest, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +struct ScannerPeerDirtyUsageExpectation { + instance_id: String, + generation: u64, + pending: bool, +} + +fn verified_remote_dirty_usage_buckets( + expected_peers: &HashMap, + peer_snapshots: Vec<(String, EcstoreScannerPeerDirtyUsageSnapshot)>, +) -> Option> { + if expected_peers.is_empty() || peer_snapshots.len() != expected_peers.len() { + return None; + } + + let mut received_peers = HashSet::with_capacity(peer_snapshots.len()); + let mut dirty_buckets = HashSet::new(); + for (host, snapshot) in peer_snapshots { + let expected = expected_peers.get(&host)?; + if !received_peers.insert(host) + || snapshot.instance_id != expected.instance_id + || snapshot.generation != expected.generation + || snapshot.generation == u64::MAX + || snapshot.protocol_version != crate::SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION + || !snapshot.complete + || snapshot.pending_bucket_count != u64::try_from(snapshot.buckets.len()).unwrap_or(u64::MAX) + || (expected.pending && snapshot.pending_bucket_count == 0) + { + return None; + } + dirty_buckets.extend(snapshot.buckets.into_keys()); + } + + (received_peers.len() == expected_peers.len()).then_some(dirty_buckets) +} + +fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<'_>) -> Option { + let data = proof.data?; + let baseline = serde_json::from_slice::(data).ok()?; + if !baseline.is_complete_bucket_usage_snapshot() + || baseline.usage_snapshot_partial + || baseline.usage_snapshot_converged != Some(true) + || baseline.scanner_epoch != Some(proof.leader_epoch) + || baseline.usage_snapshot_set_states.len() != proof.expected_sources.len() + { + 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()?); + if !proof.expected_sources.contains(&source) + || !states.insert(source) + || !state.complete + || 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) + { + return None; + } + } + + (states == *proof.expected_sources).then_some(proof.scan_plan_digest) +} + +fn scoped_scan_scope_from_dirty_buckets( + requested_scope: ScannerBucketScanScope, + dirty_buckets: HashSet, + dirty_snapshot_complete: bool, + all_buckets: &[BucketInfo], + baseline_proof: ScannerCacheBaselineProof<'_>, +) -> ScannerBucketScanScope { + if !requested_scope.is_default() || !dirty_snapshot_complete { + return requested_scope; + } + + let current_buckets = all_buckets.iter().map(|bucket| bucket.name.as_str()).collect::>(); + let selected_buckets = dirty_buckets + .into_iter() + .filter(|bucket| current_buckets.contains(bucket.as_str())) + .collect::>(); + if selected_buckets.is_empty() { + return requested_scope; + } + + let Some(baseline_scan_plan_digest) = complete_scanner_cache_baseline_plan_digest(baseline_proof) else { + return requested_scope; + }; + + ScannerBucketScanScope::from_dirty_buckets(selected_buckets, baseline_scan_plan_digest) +} + pub(crate) fn is_scanner_metadata_corrupt_error(err: &StorageError) -> bool { matches!(err, StorageError::Io(io) if io.to_string().starts_with(SCANNER_METADATA_CORRUPT_ERROR)) } @@ -749,7 +866,7 @@ mod io_cache; mod io_cycle; #[cfg(test)] use io_cache::{ScannerSetCacheGeneration, prepare_scoped_set_scan}; -pub(crate) use io_cycle::nsscanner_with_storage_status; +pub(crate) use io_cycle::{ScannerCycleRequest, nsscanner_with_storage_status_scoped}; mod io_disk; #[cfg(test)] mod publish_gate_tests; diff --git a/crates/scanner/src/scanner_io/cache.rs b/crates/scanner/src/scanner_io/cache.rs index 9522ffed2..7ff684c99 100644 --- a/crates/scanner/src/scanner_io/cache.rs +++ b/crates/scanner/src/scanner_io/cache.rs @@ -282,9 +282,26 @@ pub(super) fn completed_data_usage_info( .iter() .map(|(bucket, usage)| (bucket.clone(), usage.size)) .collect(); + let mut usage_snapshot_set_states = results + .iter() + .map(|result| { + let source = result.info.source?; + Some(DataUsageSnapshotSetState { + pool_index: u64::try_from(source.pool_index).ok()?, + set_index: u64::try_from(source.set_index).ok()?, + scanner_cycle: Some(result.info.next_cycle), + scanner_epoch: Some(result.info.leader_epoch), + scan_plan_digest: Some(result.info.scan_plan_digest?.0), + complete: true, + tombstone: false, + }) + }) + .collect::>>()?; + 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), 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()?, @@ -295,6 +312,7 @@ pub(super) fn completed_data_usage_info( bucket_sizes, buckets_usage, usage_snapshot_complete: true, + usage_snapshot_set_states, ..Default::default() }; Some((data_usage_info, merged_last_update)) diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index 3b14e9ff2..2947d638d 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -71,6 +71,7 @@ where leader_epoch, scan_mode, scan_scope: ScannerBucketScanScope::default(), + persisted_usage_baseline: None, }; nsscanner_with_storage_status_scoped(store, request).await } @@ -83,6 +84,79 @@ pub(crate) struct ScannerCycleRequest { pub(crate) leader_epoch: u64, pub(crate) scan_mode: HealScanMode, pub(crate) scan_scope: ScannerBucketScanScope, + pub(crate) persisted_usage_baseline: Option, +} + +struct ScannerBucketScopeResolution<'a> { + requested_scope: ScannerBucketScanScope, + baseline_proof: ScannerCacheBaselineProof<'a>, + activity_before: &'a crate::scanner::ScannerActivitySnapshot, + dirty_usage_snapshot: &'a DirtyUsageSnapshot, + all_buckets: &'a [BucketInfo], +} + +async fn resolve_scanner_bucket_scan_scope( + store: &S, + distributed: bool, + resolution: ScannerBucketScopeResolution<'_>, +) -> ScannerBucketScanScope +where + S: ScannerStorage, +{ + if !resolution.requested_scope.is_default() + || !resolution.dirty_usage_snapshot.covers_all_pending + || resolution.dirty_usage_snapshot.generation == u64::MAX + || resolution.dirty_usage_snapshot.buckets.len() > crate::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES + { + return resolution.requested_scope; + } + + let mut dirty_buckets = resolution + .dirty_usage_snapshot + .buckets + .keys() + .cloned() + .collect::>(); + if distributed { + let Some(notification_system) = store.scanner_notification_system() else { + return resolution.requested_scope; + }; + let Ok(peer_snapshots) = notification_system.scanner_dirty_usage_snapshots().await else { + return resolution.requested_scope; + }; + let mut expected_peers = HashMap::new(); + for (host, lease_instance_id, _) in crate::scanner::scanner_activity_publication_lease_targets(resolution.activity_before) + { + let Some((activity_instance_id, generation, pending)) = + crate::scanner::scanner_activity_dirty_usage_state_for_host(resolution.activity_before, &host) + else { + return resolution.requested_scope; + }; + if activity_instance_id != lease_instance_id || expected_peers.contains_key(&host) { + return resolution.requested_scope; + } + expected_peers.insert( + host, + ScannerPeerDirtyUsageExpectation { + instance_id: activity_instance_id.to_string(), + generation, + pending, + }, + ); + } + let Some(remote_dirty_buckets) = verified_remote_dirty_usage_buckets(&expected_peers, peer_snapshots) else { + return resolution.requested_scope; + }; + dirty_buckets.extend(remote_dirty_buckets); + } + + scoped_scan_scope_from_dirty_buckets( + resolution.requested_scope, + dirty_buckets, + true, + resolution.all_buckets, + resolution.baseline_proof, + ) } pub(crate) async fn nsscanner_with_storage_status_scoped(store: &S, request: ScannerCycleRequest) -> Result @@ -97,6 +171,7 @@ where leader_epoch, scan_mode, scan_scope, + persisted_usage_baseline, } = request; let child_token = ctx.child_token(); let _tier_cycle_guard = begin_tier_registry_cycle(want_cycle, leader_epoch); @@ -186,8 +261,26 @@ where } bucket_plan_complete &= buckets_by_source.keys().copied().collect::>() == *expected_sources; let scan_plan_digest = - scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_snapshot_digest(&activity_before)); + scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before)); 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, + distributed, + ScannerBucketScopeResolution { + requested_scope: scan_scope, + baseline_proof: ScannerCacheBaselineProof { + data: persisted_usage_baseline.as_ref(), + expected_sources: &expected_sources, + leader_epoch, + want_cycle, + scan_plan_digest, + }, + activity_before: &activity_before, + dirty_usage_snapshot: &dirty_usage_snapshot, + all_buckets: &all_buckets, + }, + ) + .await; 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; diff --git a/crates/scanner/src/scanner_io/publish_gate_tests.rs b/crates/scanner/src/scanner_io/publish_gate_tests.rs index d740cb2eb..f131cf437 100644 --- a/crates/scanner/src/scanner_io/publish_gate_tests.rs +++ b/crates/scanner/src/scanner_io/publish_gate_tests.rs @@ -655,9 +655,31 @@ fn completed_data_usage_info_requires_every_set_before_publish() { .expect("all completed sets should produce a publishable data usage snapshot"); assert_eq!(last_update, SystemTime::UNIX_EPOCH + Duration::from_secs(20)); assert_eq!(data_usage_info.scanner_cycle, Some(0)); + assert_eq!(data_usage_info.scanner_epoch, Some(0)); assert_eq!(data_usage_info.objects_total_count, 3); assert_eq!(data_usage_info.buckets_usage.len(), 3); assert!(data_usage_info.usage_snapshot_complete); + assert_eq!( + data_usage_info + .usage_snapshot_set_states + .iter() + .map(|state| { + ( + state.pool_index, + state.set_index, + state.scanner_cycle, + state.scanner_epoch, + state.scan_plan_digest, + state.complete, + state.tombstone, + ) + }) + .collect::>(), + vec![ + (0, 0, Some(0), Some(0), Some(TEST_PLAN_DIGEST.0), true, false), + (1, 0, Some(0), Some(0), Some(TEST_PLAN_DIGEST.0), true, false), + ] + ); assert_eq!( data_usage_info .buckets_usage diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index 6fddafdfe..ec2c1ad65 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -17,6 +17,7 @@ use super::io_disk::tier_stats_template; use super::*; use crate::scanner_budget::ScannerCycleBudgetConfig; use crate::scanner_folder::ScannerItem; +use crate::storage_api::EcstoreScannerPeerDirtyUsageSnapshot; use crate::storage_api::owner::{ EcstorePoolDecommissionInfo, EcstoreRebalStatus, EcstoreRebalanceInfo, EcstoreRebalanceMeta, EcstoreRebalanceStats, }; @@ -796,6 +797,195 @@ fn complete_set_usage_cache(buckets: &[(&str, usize)], scan_plan_digest: DataUsa cache } +fn complete_usage_baseline( + source: DataUsageCacheSource, + scan_plan_digest: DataUsageScanPlanDigest, + scanner_cycle: u64, + scanner_epoch: u64, +) -> bytes::Bytes { + let baseline = DataUsageInfo { + last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(10)), + scanner_cycle: Some(scanner_cycle), + scanner_epoch: Some(scanner_epoch), + buckets_count: 1, + buckets_usage: HashMap::from([("photos".to_string(), Default::default())]), + usage_snapshot_complete: true, + usage_snapshot_converged: Some(true), + usage_snapshot_set_states: vec![DataUsageSnapshotSetState { + pool_index: u64::try_from(source.pool_index).expect("test pool index should fit"), + set_index: u64::try_from(source.set_index).expect("test set index should fit"), + scanner_cycle: Some(scanner_cycle), + scanner_epoch: Some(scanner_epoch), + scan_plan_digest: Some(scan_plan_digest.0), + complete: true, + tombstone: false, + }], + ..Default::default() + }; + bytes::Bytes::from(serde_json::to_vec(&baseline).expect("test baseline should encode")) +} + +#[test] +fn scoped_scan_requires_a_converged_complete_baseline_with_exact_set_provenance() { + let source = DataUsageCacheSource::new(1, 2); + let expected_sources = HashSet::from([source]); + let scan_plan_digest = DataUsageScanPlanDigest([9; 32]); + let baseline = complete_usage_baseline(source, scan_plan_digest, 7, 11); + + assert_eq!( + complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof { + data: Some(&baseline), + expected_sources: &expected_sources, + leader_epoch: 11, + want_cycle: 8, + scan_plan_digest, + }), + Some(scan_plan_digest) + ); + + let mut incomplete = serde_json::from_slice::(&baseline).expect("test baseline should decode"); + incomplete.usage_snapshot_converged = Some(false); + let incomplete = bytes::Bytes::from(serde_json::to_vec(&incomplete).expect("test baseline should encode")); + assert_eq!( + complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof { + data: Some(&incomplete), + expected_sources: &expected_sources, + leader_epoch: 11, + want_cycle: 8, + scan_plan_digest, + }), + None + ); + + let mut wrong_provenance = serde_json::from_slice::(&baseline).expect("test baseline should decode"); + wrong_provenance.usage_snapshot_set_states[0].scan_plan_digest = Some([8; 32]); + let wrong_provenance = bytes::Bytes::from(serde_json::to_vec(&wrong_provenance).expect("test baseline should encode")); + assert_eq!( + complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof { + data: Some(&wrong_provenance), + expected_sources: &expected_sources, + leader_epoch: 11, + want_cycle: 8, + scan_plan_digest, + }), + None + ); +} + +#[test] +fn scoped_scan_selects_only_current_dirty_buckets_after_baseline_validation() { + let source = DataUsageCacheSource::new(1, 2); + let expected_sources = HashSet::from([source]); + let baseline_scan_plan_digest = DataUsageScanPlanDigest([4; 32]); + let current_scan_plan_digest = DataUsageScanPlanDigest([5; 32]); + let baseline = complete_usage_baseline(source, current_scan_plan_digest, 7, 11); + let scope = scoped_scan_scope_from_dirty_buckets( + ScannerBucketScanScope::default(), + HashSet::from(["photos".to_string(), "deleted".to_string()]), + true, + &[bucket_info("photos")], + ScannerCacheBaselineProof { + data: Some(&baseline), + expected_sources: &expected_sources, + leader_epoch: 11, + want_cycle: 8, + scan_plan_digest: current_scan_plan_digest, + }, + ); + + assert_eq!(scope.baseline_scan_plan_digest, Some(current_scan_plan_digest)); + assert_eq!( + scope + .selected_buckets + .as_deref() + .expect("validated scope should select a bucket"), + &HashSet::from(["photos".to_string()]) + ); + assert_ne!(scope.baseline_scan_plan_digest, Some(baseline_scan_plan_digest)); +} + +fn peer_dirty_usage_snapshot( + instance_id: &str, + generation: u64, + complete: bool, + buckets: &[(&str, u64)], +) -> EcstoreScannerPeerDirtyUsageSnapshot { + EcstoreScannerPeerDirtyUsageSnapshot { + instance_id: instance_id.to_string(), + generation, + pending_bucket_count: u64::try_from(buckets.len()).expect("test bucket count should fit"), + protocol_version: crate::SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION, + complete, + buckets: buckets + .iter() + .map(|(bucket, generation)| ((*bucket).to_string(), *generation)) + .collect(), + } +} + +#[test] +fn verified_remote_dirty_usage_buckets_merges_only_complete_current_snapshots() { + let expected_peers = HashMap::from([ + ( + "node-a:9000".to_string(), + ScannerPeerDirtyUsageExpectation { + instance_id: "instance-a".to_string(), + generation: 7, + pending: true, + }, + ), + ( + "node-b:9000".to_string(), + ScannerPeerDirtyUsageExpectation { + instance_id: "instance-b".to_string(), + generation: 3, + pending: false, + }, + ), + ]); + + assert_eq!( + verified_remote_dirty_usage_buckets( + &expected_peers, + vec![ + ( + "node-a:9000".to_string(), + peer_dirty_usage_snapshot("instance-a", 7, true, &[("photos", 7)]), + ), + ( + "node-b:9000".to_string(), + peer_dirty_usage_snapshot("instance-b", 3, true, &[("archive", 3)]), + ), + ], + ), + Some(HashSet::from(["photos".to_string(), "archive".to_string()])) + ); +} + +#[test] +fn verified_remote_dirty_usage_buckets_rejects_incomplete_or_stale_peer_state() { + let expected_peers = HashMap::from([( + "node-a:9000".to_string(), + ScannerPeerDirtyUsageExpectation { + instance_id: "instance-a".to_string(), + generation: 7, + pending: true, + }, + )]); + + for snapshot in [ + peer_dirty_usage_snapshot("instance-a", 7, false, &[("photos", 7)]), + peer_dirty_usage_snapshot("instance-a", 6, true, &[("photos", 6)]), + peer_dirty_usage_snapshot("instance-b", 7, true, &[("photos", 7)]), + peer_dirty_usage_snapshot("instance-a", 7, true, &[]), + ] { + assert!( + verified_remote_dirty_usage_buckets(&expected_peers, vec![("node-a:9000".to_string(), snapshot)]).is_none(), + "incomplete, stale, mismatched, or empty pending peer state must fall back to a full scan" + ); + } +} + #[test] fn scoped_set_scan_preserves_unselected_usage_and_drops_deleted_buckets() { let baseline_digest = DataUsageScanPlanDigest([1; 32]); diff --git a/crates/scanner/src/storage_api.rs b/crates/scanner/src/storage_api.rs index a0aa874b0..f81065d59 100644 --- a/crates/scanner/src/storage_api.rs +++ b/crates/scanner/src/storage_api.rs @@ -103,7 +103,9 @@ pub(crate) use rustfs_ecstore::api::rebalance::{ RebalStatus as EcstoreRebalStatus, RebalanceInfo as EcstoreRebalanceInfo, RebalanceMeta as EcstoreRebalanceMeta, RebalanceStats as EcstoreRebalanceStats, }; -pub(crate) use rustfs_ecstore::api::rpc::ScannerBucketListing as EcstoreScannerBucketListing; +pub(crate) use rustfs_ecstore::api::rpc::{ + ScannerBucketListing as EcstoreScannerBucketListing, ScannerPeerDirtyUsageSnapshot as EcstoreScannerPeerDirtyUsageSnapshot, +}; #[cfg(test)] pub(crate) use rustfs_ecstore::api::runtime::InstanceContext as EcstoreInstanceContext; pub(crate) use rustfs_ecstore::api::runtime::{ From cf9688898d9fc9ae11d6b16793f355435618d122 Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 5 Sep 2026 16:39:26 +0800 Subject: [PATCH 2/5] fix(heal): retain completed task progress (#7177) * 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 Co-Authored-By: zhi22915 * fix(heal): retain completed task progress Refs rustfs/backlog#2262 and rustfs/backlog#2240. Co-Authored-By: heihutu Co-Authored-By: zhi22915 --------- Co-authored-by: heihutu Co-authored-by: zhi22915 --- crates/heal/src/heal/manager.rs | 72 ++++- crates/heal/src/heal/manager/queue.rs | 129 ++++++++ crates/heal/src/heal/manager/scheduler.rs | 114 ++++--- crates/heal/src/heal/manager/tests.rs | 351 +++++++++++++++++++++- crates/heal/src/heal/task.rs | 2 +- 5 files changed, 613 insertions(+), 55 deletions(-) diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index d36c9f9a7..971590fdb 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -45,6 +45,11 @@ use tracing::{debug, error, info, warn}; use super::{DiskError, Endpoint, HealDiskExt as _, local_disk_map_read}; const KEEP_HEAL_TASK_STATUS_DURATION: Duration = Duration::from_secs(10 * 60); +// Each cache includes alias tokens in its count and byte budget. Eviction +// removes every token sharing a snapshot; neither cache retains repair state. +const MAX_COMPLETED_HEAL_TOKENS: usize = 1024; +const MAX_COMPLETED_HEAL_BYTES: usize = 64 * 1024 * 1024; +const MAX_COMPLETED_HEAL_RESULT_BYTES: usize = 1024 * 1024; const DISPLACED_HEAL_REASON: &str = "reason=displaced; retry_hint=submit_again"; const LOG_COMPONENT_HEAL: &str = "heal"; const LOG_SUBSYSTEM_DISK_SCANNER: &str = "disk_scanner"; @@ -180,6 +185,8 @@ fn record_displaced_terminal( request: &HealRequest, ) -> Arc { let terminal = Arc::new(CompletedHealStatus { + progress: None, + retained_bytes: std::sync::OnceLock::new(), heal_type: request.heal_type.clone(), status: HealTaskStatus::Failed { error: format!("heal task displaced by a higher-priority request ({DISPLACED_HEAL_REASON})"), @@ -193,6 +200,7 @@ fn record_displaced_terminal( let mut terminals = lock_displaced_terminals(registry); prune_completed_heal_statuses(&mut terminals); terminals.insert(request.id.clone(), Arc::clone(&terminal)); + prune_completed_heal_statuses(&mut terminals); terminal } @@ -209,9 +217,15 @@ async fn remove_displaced_task_aliases( .collect::>(); let mut displaced_terminals = lock_displaced_terminals(terminals); prune_completed_heal_statuses(&mut displaced_terminals); - for alias_id in alias_ids { - displaced_terminals.insert(alias_id, Arc::clone(terminal)); + if displaced_terminals + .get(task_id) + .is_some_and(|current| Arc::ptr_eq(current, terminal)) + { + for alias_id in alias_ids { + displaced_terminals.insert(alias_id, Arc::clone(terminal)); + } } + prune_completed_heal_statuses(&mut displaced_terminals); aliases.retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id); } @@ -222,6 +236,36 @@ async fn remove_task_aliases_for_task(registry: &Arc +// retrying (when needed) -> aliases -> completed; queries release aliases +// before looking up active state. Publishing aliases before removing their +// mapping keeps both an already-resolved token and a new lookup valid. +async fn publish_completed_heal( + completed_heals: &Mutex>>, + task_aliases: &Mutex>, + task_id: &str, + completed: CompletedHealStatus, + terminal: bool, +) { + let completed = Arc::new(completed); + completed.retained_bytes(); + let mut aliases = task_aliases.lock().await; + let mut retained = completed_heals.lock().await; + if let Some(previous) = retained.get(task_id).cloned() { + for entry in retained.values_mut().filter(|entry| Arc::ptr_eq(entry, &previous)) { + *entry = Arc::clone(&completed); + } + } + retained.insert(task_id.to_owned(), Arc::clone(&completed)); + if terminal { + for (alias_id, _) in aliases.iter().filter(|(_, alias)| alias.task_id == task_id) { + retained.insert(alias_id.clone(), Arc::clone(&completed)); + } + aliases.retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id); + } + prune_completed_heal_statuses(&mut retained); +} + #[derive(Debug, Clone)] pub struct HealTaskReport { pub status: HealTaskStatus, @@ -268,7 +312,7 @@ fn completed_task_report(completed: &CompletedHealStatus, since: Option) -> let result_items = match since { None => completed.seqed_items.iter().map(|(_, item)| item.clone()).collect(), Some(cursor) => { - if cursor + 1 < completed.min_seq { + if cursor.saturating_add(1) < completed.min_seq { lagged = true; } completed @@ -283,7 +327,7 @@ fn completed_task_report(completed: &CompletedHealStatus, since: Option) -> status: completed.status.clone(), result_items, result_items_truncated: completed.result_items_truncated || lagged, - progress: None, + progress: completed.progress.clone(), next_seq: completed.next_seq, min_seq: completed.min_seq, } @@ -1847,14 +1891,14 @@ impl HealManager { pub async fn get_task_progress(&self, task_id: &str) -> Result { let canonical_task_id = self.canonical_task_id(task_id).await; - let active_heals = self.active_heals.lock().await; - if let Some(task) = active_heals.get(&canonical_task_id) { - Ok(task.get_progress().await) - } else { - Err(Error::TaskNotFound { - task_id: task_id.to_string(), - }) - } + let progress = match self.lookup_task_state(&canonical_task_id, None).await { + TaskStateLookup::Active(task) => Some(task.get_progress().await), + TaskStateLookup::Completed(completed) => completed.progress.clone(), + _ => None, + }; + progress.ok_or_else(|| Error::TaskNotFound { + task_id: task_id.to_string(), + }) } /// Cancel task @@ -1864,6 +1908,8 @@ impl HealManager { let mut active_heals = self.active_heals.lock().await; if let Some(task) = active_heals.get(&canonical_task_id) { task.cancel().await?; + let completed = CompletedHealStatus::snapshot(task, HealTaskStatus::Cancelled).await; + publish_completed_heal(&self.completed_heals, &self.task_aliases, &canonical_task_id, completed, true).await; active_heals.remove(&canonical_task_id); publish_active_heal_count(&active_heals); info!( @@ -1940,6 +1986,8 @@ impl HealManager { for task_id in &task_ids { if let Some(task) = active_heals.get(task_id) { task.cancel().await?; + let completed = CompletedHealStatus::snapshot(task, HealTaskStatus::Cancelled).await; + publish_completed_heal(&self.completed_heals, &self.task_aliases, task_id, completed, true).await; } active_heals.remove(task_id); cancelled += 1; diff --git a/crates/heal/src/heal/manager/queue.rs b/crates/heal/src/heal/manager/queue.rs index b4608f27b..aceabe42a 100644 --- a/crates/heal/src/heal/manager/queue.rs +++ b/crates/heal/src/heal/manager/queue.rs @@ -82,6 +82,8 @@ pub(super) enum QueuePushOutcome { pub(super) struct CompletedHealStatus { pub(super) heal_type: HealType, pub(super) status: HealTaskStatus, + pub(super) progress: Option, + pub(super) retained_bytes: std::sync::OnceLock, pub(super) result_items_truncated: bool, pub(super) completed_at: SystemTime, /// Sequence-stamped retained window, archived with the completion so @@ -92,6 +94,133 @@ pub(super) struct CompletedHealStatus { pub(super) min_seq: u64, } +impl CompletedHealStatus { + // Account for owned capacities, including nested drive arrays. Aliases + // conservatively charge the shared allocation again, keeping both token + // count and retained payload bounded without a second ownership index. + pub(super) fn retained_bytes(&self) -> usize { + *self.retained_bytes.get_or_init(|| self.measure_retained_bytes()) + } + + fn measure_retained_bytes(&self) -> usize { + let mut bytes = size_of::(); + let mut add = |amount: usize| bytes = bytes.saturating_add(amount); + match &self.heal_type { + HealType::Cluster => {} + HealType::Bucket { bucket } => add(bucket.capacity()), + HealType::Object { + bucket, + object, + version_id, + } + | HealType::ECDecode { + bucket, + object, + version_id, + } => { + add(bucket.capacity()); + add(object.capacity()); + add(version_id.as_ref().map_or(0, String::capacity)); + } + HealType::Prefix { bucket, prefix } => { + add(bucket.capacity()); + add(prefix.capacity()); + } + HealType::Metadata { bucket, object } => { + add(bucket.capacity()); + add(object.capacity()); + } + HealType::ErasureSet { buckets, set_disk_id } => { + add(buckets.capacity().saturating_mul(size_of::())); + for bucket in buckets { + add(bucket.capacity()); + } + add(set_disk_id.capacity()); + } + } + if let HealTaskStatus::Failed { error } | HealTaskStatus::Retrying { error, .. } = &self.status { + add(error.capacity()); + } + add(self + .progress + .as_ref() + .and_then(|progress| progress.current_object.as_ref()) + .map_or(0, String::capacity)); + add(self.seqed_items.capacity().saturating_mul(size_of::<(u64, HealResultItem)>())); + for (_, item) in &self.seqed_items { + add(Self::result_item_heap_bytes(item)); + } + bytes + } + + fn result_item_heap_bytes(item: &HealResultItem) -> usize { + let mut bytes = 0usize; + let mut add = |amount: usize| bytes = bytes.saturating_add(amount); + for value in [ + &item.heal_item_type, + &item.bucket, + &item.object, + &item.version_id, + &item.detail, + ] { + add(value.capacity()); + } + for infos in [&item.before, &item.after] { + add(infos + .drives + .capacity() + .saturating_mul(size_of::())); + for drive in &infos.drives { + add(drive.uuid.capacity()); + add(drive.endpoint.capacity()); + add(drive.state.capacity()); + } + } + bytes + } + + pub(super) fn bound_result_window(&mut self) { + let mut bytes = 0usize; + let retained = self + .seqed_items + .iter() + .rev() + .take_while(|(_, item)| { + bytes = bytes + .saturating_add(size_of::<(u64, HealResultItem)>()) + .saturating_add(Self::result_item_heap_bytes(item)); + bytes <= MAX_COMPLETED_HEAL_RESULT_BYTES + }) + .count(); + let truncated = retained < self.seqed_items.len(); + if truncated { + self.seqed_items.drain(..self.seqed_items.len() - retained); + self.seqed_items.shrink_to_fit(); + self.min_seq = self.seqed_items.first().map_or(self.next_seq, |(seq, _)| *seq); + self.result_items_truncated = true; + self.retained_bytes.take(); + } + } + + pub(super) async fn snapshot(task: &HealTask, status: HealTaskStatus) -> Self { + let seqed_items = task.get_seqed_result_items().await; + let (next_seq, min_seq) = task.result_seq_cursors(); + let mut snapshot = Self { + heal_type: task.heal_type.clone(), + status, + progress: Some(task.get_progress().await), + retained_bytes: std::sync::OnceLock::new(), + result_items_truncated: task.result_items_truncated(), + completed_at: SystemTime::now(), + seqed_items, + next_seq, + min_seq, + }; + snapshot.bound_result_window(); + snapshot + } +} + #[derive(Debug, Clone)] pub(super) struct HealTaskAlias { pub(super) task_id: String, diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index c0deae524..feb0c1231 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -264,7 +264,7 @@ impl HealManager { error: error.clone(), retry_attempt: request.retry_attempts, }); - let retry_request_for_queue = retry_request; + let mut retry_request_for_queue = retry_request; let retry_cancel_token = retry_request_for_queue.as_ref().map(|_| CancellationToken::new()); if retry_request_for_queue.is_none() { replacement_recovery_anchors_clone @@ -272,7 +272,35 @@ impl HealManager { .unwrap_or_else(|poisoned| poisoned.into_inner()) .remove(&task_id); } + let mut completed_status = match retry_request_for_status { + Some(status) => status, + None => task.get_status().await, + }; + let mut completed_status_entry = CompletedHealStatus::snapshot(&task, completed_status.clone()).await; + let completed_progress = task.get_progress().await; + #[cfg(test)] + tests::pause_completed_retention_before_publish(&task_id, &completed_status).await; let mut active_heals_guard = active_heals_clone.lock().await; + let owns_completion = active_heals_guard.contains_key(&task_id); + let cancelled_completion = if owns_completion { + false + } else { + // Cancellation can win while a finished worker waits + // for active ownership. It must not resurrect a retry + // or replace an acknowledged cancellation with success. + retry_request_for_queue = None; + completed_heals_clone + .lock() + .await + .get(&task_id) + .is_some_and(|completed| completed.status == HealTaskStatus::Cancelled) + }; + if cancelled_completion { + completed_status = HealTaskStatus::Cancelled; + completed_status_entry.status = HealTaskStatus::Cancelled; + } + let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. }); + let successful_completion = matches!(completed_status, HealTaskStatus::Completed); // Keep retry ownership continuous: status snapshots acquire // these locks in the same active -> retrying order. let mut retrying_heals_guard = if let (Some((request, _, error)), Some(cancel_token)) = @@ -295,6 +323,16 @@ impl HealManager { } else { None }; + if owns_completion || cancelled_completion { + publish_completed_heal( + &completed_heals_clone, + &task_aliases_clone, + &task_id, + completed_status_entry, + terminal_completion, + ) + .await; + } let completed_task = active_heals_guard.remove(&task_id); if let Some(completed_task) = completed_task.as_ref() { publish_active_heal_count(&active_heals_guard); @@ -304,33 +342,10 @@ impl HealManager { drop(retrying_heals_guard.take()); drop(active_heals_guard); - if let Some(completed_task) = completed_task { - let completed_status = if let Some(status) = retry_request_for_status { - status - } else { - completed_task.get_status().await - }; - let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. }); - let successful_completion = matches!(completed_status, HealTaskStatus::Completed); - let completed_progress = completed_task.get_progress().await; - // Single snapshot of the retained window: the task is - // finished and already off the active map, so there is - // no concurrent writer to race with. - let seqed_items = completed_task.get_seqed_result_items().await; - let (next_seq, min_seq) = completed_task.result_seq_cursors(); - let completed_status_entry = CompletedHealStatus { - heal_type: completed_task.heal_type.clone(), - status: completed_status.clone(), - result_items_truncated: completed_task.result_items_truncated(), - completed_at: SystemTime::now(), - seqed_items, - next_seq, - min_seq, - }; - let mut completed_heals_guard = completed_heals_clone.lock().await; - prune_completed_heal_statuses(&mut completed_heals_guard); - completed_heals_guard.insert(task_id.clone(), Arc::new(completed_status_entry)); - drop(completed_heals_guard); + #[cfg(test)] + tests::pause_completed_retention_handoff(&task_id).await; + + if completed_task.is_some() { // update statistics let mut stats = statistics_clone.write().await; match completed_status { @@ -352,10 +367,6 @@ impl HealManager { } else { release_mrf_repair_notice_targets(notice_targets); } - task_aliases_clone - .lock() - .await - .retain(|alias_id, alias| alias_id != &task_id && alias.task_id != task_id); } } @@ -718,17 +729,42 @@ pub(super) fn heal_request_set_key_for_task(task: &HealTask) -> Option { } pub(super) fn prune_completed_heal_statuses(completed_heals: &mut HashMap>) { - let Ok(now) = SystemTime::now().duration_since(SystemTime::UNIX_EPOCH) else { - return; - }; + prune_completed_heal_statuses_at(completed_heals, SystemTime::now()); +} +pub(super) fn prune_completed_heal_statuses_at(completed_heals: &mut HashMap>, now: SystemTime) { completed_heals.retain(|_, completed| { - completed - .completed_at - .duration_since(SystemTime::UNIX_EPOCH) - .map(|completed_at| now.saturating_sub(completed_at) <= KEEP_HEAL_TASK_STATUS_DURATION) + now.duration_since(completed.completed_at) + .map(|age| age <= KEEP_HEAL_TASK_STATUS_DURATION) .unwrap_or(false) }); + let entry_bytes = |key: &String, value: &Arc| { + key.capacity() + .saturating_add(size_of::<(String, Arc)>()) + .saturating_add(value.retained_bytes()) + }; + let mut bytes = completed_heals + .iter() + .fold(0usize, |total, (key, value)| total.saturating_add(entry_bytes(key, value))); + while completed_heals.len() > MAX_COMPLETED_HEAL_TOKENS || bytes > MAX_COMPLETED_HEAL_BYTES { + let Some(oldest) = completed_heals + .iter() + .min_by(|(left_id, left), (right_id, right)| { + left.completed_at.cmp(&right.completed_at).then_with(|| left_id.cmp(right_id)) + }) + .map(|(_, value)| Arc::clone(value)) + else { + break; + }; + completed_heals.retain(|key, value| { + if Arc::ptr_eq(value, &oldest) { + bytes = bytes.saturating_sub(entry_bytes(key, value)); + false + } else { + true + } + }); + } } pub(super) fn can_schedule_request( diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index 37a100f4b..aa1c1293f 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -101,6 +101,326 @@ async fn process_manager_queue_once(manager: &HealManager) { struct MockStorage; +fn completed_retention_fixture(completed_at: SystemTime) -> CompletedHealStatus { + CompletedHealStatus { + heal_type: HealType::Cluster, + status: HealTaskStatus::Completed, + progress: Some(HealProgress { + objects_scanned: 9, + objects_healed: 8, + objects_failed: 1, + ..Default::default() + }), + retained_bytes: std::sync::OnceLock::new(), + result_items_truncated: false, + completed_at, + seqed_items: vec![(3, HealResultItem::default()), (4, HealResultItem::default())], + next_seq: 5, + min_seq: 3, + } +} + +#[test] +fn completed_retention_cursor_boundaries_preserve_progress() { + let completed = completed_retention_fixture(SystemTime::now()); + for (cursor, count, lagged) in [ + (0, 2, true), + (1, 2, true), + (2, 2, false), + (3, 1, false), + (4, 0, false), + (5, 0, false), + (u64::MAX, 0, false), + ] { + let report = completed_task_report(&completed, Some(cursor)); + assert_eq!(report.result_items.len(), count, "cursor={cursor}"); + assert_eq!(report.result_items_truncated, lagged, "cursor={cursor}"); + assert_eq!(report.progress, completed.progress); + assert_eq!((report.next_seq, report.min_seq), (5, 3)); + } + assert_eq!(completed_task_report(&completed, None).result_items.len(), 2); +} + +#[tokio::test] +async fn completed_retention_displaced_alias_does_not_resurrect_evicted_snapshot() { + let manager = HealManager::new(Arc::new(MockStorage), None); + let request = HealRequest::bucket("bucket".to_string()); + manager.insert_task_alias("alias", &request.id).await; + let terminal = record_displaced_terminal(&manager.displaced_terminals, &request); + lock_displaced_terminals(&manager.displaced_terminals).remove(&request.id); + remove_displaced_task_aliases(&manager.task_aliases, &manager.displaced_terminals, &request.id, &terminal).await; + for token in [&request.id, &"alias".to_string()] { + assert!(matches!(manager.get_task_report(token).await, Err(Error::TaskNotFound { .. }))); + } + assert!(manager.task_aliases.lock().await.is_empty()); + assert!(lock_displaced_terminals(&manager.displaced_terminals).is_empty()); +} + +#[test] +fn completed_retention_count_ttl_and_alias_eviction_are_bounded() { + let now = SystemTime::now(); + let mut entries = HashMap::new(); + let oldest = Arc::new(completed_retention_fixture(now - KEEP_HEAL_TASK_STATUS_DURATION)); + entries.insert("oldest".to_string(), Arc::clone(&oldest)); + entries.insert("oldest-alias".to_string(), Arc::clone(&oldest)); + for index in 2..MAX_COMPLETED_HEAL_TOKENS { + entries.insert(format!("task-{index}"), Arc::new(completed_retention_fixture(now))); + } + prune_completed_heal_statuses_at(&mut entries, now); + assert_eq!(entries.len(), MAX_COMPLETED_HEAL_TOKENS); + entries.insert("cap-plus-one".to_string(), Arc::new(completed_retention_fixture(now))); + prune_completed_heal_statuses_at(&mut entries, now); + assert_eq!(entries.len(), MAX_COMPLETED_HEAL_TOKENS - 1); + assert!(!entries.contains_key("oldest")); + assert!(!entries.contains_key("oldest-alias")); + entries.clear(); + entries.insert("ttl-boundary".to_string(), oldest); + entries.insert( + "expired".to_string(), + Arc::new(completed_retention_fixture( + now - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_nanos(1), + )), + ); + entries.insert("future".to_string(), Arc::new(completed_retention_fixture(now + Duration::from_nanos(1)))); + prune_completed_heal_statuses_at(&mut entries, now); + assert_eq!(entries.len(), 1); + assert!(entries.contains_key("ttl-boundary")); + prune_completed_heal_statuses_at(&mut entries, now + Duration::from_nanos(1)); + assert!(entries.is_empty()); +} + +#[test] +fn completed_retention_total_byte_cap_and_cap_plus_one() { + let now = SystemTime::now(); + let key = "large".to_string(); + let mut entry = completed_retention_fixture(now); + let base_bytes = entry.retained_bytes() + key.capacity() + size_of::<(String, Arc)>(); + entry.retained_bytes.take(); + entry.status = HealTaskStatus::Failed { + error: "x".repeat(MAX_COMPLETED_HEAL_BYTES - base_bytes), + }; + assert_eq!( + entry.retained_bytes() + key.capacity() + size_of::<(String, Arc)>(), + MAX_COMPLETED_HEAL_BYTES + ); + let mut entries = HashMap::from([(key, Arc::new(entry))]); + prune_completed_heal_statuses_at(&mut entries, now); + assert_eq!(entries.len(), 1, "exact byte cap remains retained"); + let mut over = Arc::try_unwrap(entries.remove("large").expect("entry retained")).expect("entry not shared"); + over.retained_bytes.take(); + if let HealTaskStatus::Failed { error } = &mut over.status { + *error = "x".repeat(error.len() + 1); + } + entries.insert("large".to_string(), Arc::new(over)); + prune_completed_heal_statuses_at(&mut entries, now); + assert!(entries.is_empty(), "oversized metadata cannot escape total byte bound"); +} + +#[tokio::test] +async fn completed_retention_large_window_keeps_cursors_and_progress() { + let task = HealTask::from_request(HealRequest::bucket("bucket".to_string()), Arc::new(MockStorage)); + let mut snapshot = completed_retention_fixture(SystemTime::now()); + snapshot.seqed_items[0].1.detail = "x".repeat(MAX_COMPLETED_HEAL_RESULT_BYTES); + snapshot.bound_result_window(); + assert_eq!(snapshot.seqed_items.len(), 1); + assert_eq!((snapshot.min_seq, snapshot.next_seq), (4, 5)); + assert!(snapshot.result_items_truncated); + assert!(snapshot.retained_bytes() < MAX_COMPLETED_HEAL_RESULT_BYTES); + let report = completed_task_report(&snapshot, Some(0)); + assert_eq!(report.progress.expect("progress retained").objects_scanned, 9); + assert!(report.result_items_truncated); + let active_max = task.get_result_items_since(Some(u64::MAX)).await; + assert!(active_max.items.is_empty()); + assert!(!active_max.lagged); +} + +#[test] +fn completed_retention_result_byte_cap_and_cap_plus_one() { + for extra in [0, 1] { + let mut snapshot = completed_retention_fixture(SystemTime::now()); + snapshot.seqed_items = vec![( + 4, + HealResultItem { + detail: "x".repeat(MAX_COMPLETED_HEAL_RESULT_BYTES - size_of::<(u64, HealResultItem)>() + extra), + ..Default::default() + }, + )]; + snapshot.min_seq = 4; + snapshot.bound_result_window(); + assert_eq!(snapshot.seqed_items.len(), 1 - extra); + assert_eq!(snapshot.result_items_truncated, extra == 1); + assert_eq!(snapshot.min_seq, if extra == 0 { 4 } else { 5 }); + assert_eq!(snapshot.next_seq, 5); + assert_eq!(snapshot.progress.as_ref().expect("progress retained").objects_scanned, 9); + } +} + +#[derive(Default)] +struct CompletedRetentionHook { + started: Notify, + execute: Notify, + handoff: Notify, + finish: Notify, + pause_before_publish: bool, + before_publish: Notify, + publish: Notify, + prepared_status: Mutex>, +} + +static COMPLETED_RETENTION_HOOKS: LazyLock>>> = + LazyLock::new(|| Mutex::new(HashMap::new())); + +pub(super) async fn pause_completed_retention_handoff(task_id: &str) { + let hook = COMPLETED_RETENTION_HOOKS.lock().await.get(task_id).cloned(); + if let Some(hook) = hook { + hook.handoff.notify_one(); + hook.finish.notified().await; + } +} + +pub(super) async fn pause_completed_retention_before_publish(task_id: &str, status: &HealTaskStatus) { + let hook = COMPLETED_RETENTION_HOOKS.lock().await.get(task_id).cloned(); + if let Some(hook) = hook.filter(|hook| hook.pause_before_publish) { + *hook.prepared_status.lock().await = Some(status.clone()); + hook.before_publish.notify_one(); + hook.publish.notified().await; + } +} + +#[tokio::test] +async fn completed_retention_cancel_wins_over_a_prepared_retry_snapshot() { + let bucket = "completed-retention-retry-cancel"; + let manager = HealManager::new(Arc::new(MockStorage), None); + let request = HealRequest::object(bucket.to_string(), "object".to_string(), None); + let task_id = request.id.clone(); + let duplicate = HealRequest::object(bucket.to_string(), "object".to_string(), None); + let alias = duplicate.id.clone(); + let hook = Arc::new(CompletedRetentionHook { + pause_before_publish: true, + ..Default::default() + }); + { + let mut hooks = COMPLETED_RETENTION_HOOKS.lock().await; + hooks.insert(bucket.to_string(), Arc::clone(&hook)); + hooks.insert(task_id.clone(), Arc::clone(&hook)); + } + manager.submit_heal_request(request).await.expect("admit original"); + manager.submit_heal_request(duplicate).await.expect("admit alias"); + process_manager_queue_once(&manager).await; + tokio::time::timeout(Duration::from_secs(5), hook.started.notified()) + .await + .expect("scheduler starts"); + let task = manager.active_heals.lock().await.get(&task_id).cloned().expect("active task"); + task.progress.write().await.update_object_progress(1, 1, 0, 0, 4096); + hook.execute.notify_one(); + tokio::time::timeout(Duration::from_secs(5), hook.before_publish.notified()) + .await + .expect("retry snapshot prepared"); + manager.cancel_task(&alias).await.expect("cancel wins active ownership"); + assert!(matches!(*hook.prepared_status.lock().await, Some(HealTaskStatus::Retrying { .. }))); + hook.publish.notify_one(); + tokio::time::timeout(Duration::from_secs(5), hook.handoff.notified()) + .await + .expect("scheduler finishes handoff"); + for token in [&task_id, &alias] { + let report = manager.get_task_report(token).await.expect("cancelled token retained"); + assert_eq!(report.status, HealTaskStatus::Cancelled); + assert_eq!(report.progress.expect("frozen progress").objects_scanned, 1); + } + assert!(!manager.retrying_heals.lock().await.contains_key(&task_id)); + assert!(!manager.heal_queue.lock().await.contains_request_id(&task_id)); + hook.finish.notify_one(); + COMPLETED_RETENTION_HOOKS + .lock() + .await + .retain(|key, _| key != bucket && key != &task_id); +} + +#[tokio::test] +async fn completed_retention_scheduler_preserves_progress_aliases_and_atomic_handoff() { + for outcome in ["success", "failed", "cancelled"] { + let bucket = format!("completed-retention-{outcome}"); + let hook = Arc::new(CompletedRetentionHook::default()); + let manager = Arc::new(HealManager::new(Arc::new(MockStorage), None)); + let request = HealRequest::object(bucket.clone(), "object".to_string(), None); + let task_id = request.id.clone(); + let duplicate = HealRequest::object(bucket.clone(), "object".to_string(), None); + let alias = duplicate.id.clone(); + { + let mut hooks = COMPLETED_RETENTION_HOOKS.lock().await; + hooks.insert(bucket.clone(), Arc::clone(&hook)); + hooks.insert(task_id.clone(), Arc::clone(&hook)); + } + manager.submit_heal_request(request).await.expect("admit original"); + manager.submit_heal_request(duplicate).await.expect("admit alias"); + process_manager_queue_once(&manager).await; + tokio::time::timeout(Duration::from_secs(5), hook.started.notified()) + .await + .expect("scheduler reaches storage"); + let task = manager + .active_heals + .lock() + .await + .get(&task_id) + .cloned() + .expect("task is active"); + task.progress.write().await.update_object_progress(1, 1, 0, 0, 4096); + let before = manager.get_task_report(&alias).await.expect("alias resolves active progress"); + assert_eq!(before.progress.as_ref().expect("active progress").objects_scanned, 1); + let poll_manager = Arc::clone(&manager); + let poll_alias = alias.clone(); + let stop = CancellationToken::new(); + let poll_stop = stop.clone(); + let polling = tokio::spawn(async move { + while !poll_stop.is_cancelled() { + let report = poll_manager + .get_task_report(&poll_alias) + .await + .expect("handoff must never return NotFound"); + assert!(report.progress.expect("progress never disappears").objects_scanned >= 1); + tokio::task::yield_now().await; + } + }); + if outcome == "cancelled" { + manager.cancel_task(&alias).await.expect("cancel active task by alias"); + } else { + hook.execute.notify_one(); + } + tokio::time::timeout(Duration::from_secs(5), hook.handoff.notified()) + .await + .expect("scheduler archives terminal"); + assert!(!manager.active_heals.lock().await.contains_key(&task_id)); + let expected = task.get_progress().await; + for token in [&task_id, &alias] { + assert_eq!(manager.get_task_progress(token).await.expect("terminal progress query"), expected); + let report = manager + .get_task_report_for_path_since(&format!("{bucket}/object"), token, Some(u64::MAX)) + .await + .expect("terminal token remains queryable at handoff"); + assert_eq!(report.progress.as_ref(), Some(&expected)); + assert!(report.result_items.is_empty()); + match outcome { + "success" => assert_eq!(report.status, HealTaskStatus::Completed), + "failed" => assert!(matches!(report.status, HealTaskStatus::Failed { .. })), + _ => assert_eq!(report.status, HealTaskStatus::Cancelled), + } + } + let retained = manager.completed_heals.lock().await; + assert!(Arc::ptr_eq(&retained[&task_id], &retained[&alias])); + drop(retained); + stop.cancel(); + polling.await.expect("concurrent polling succeeds"); + // Archived progress must not alias a mutable live progress object. + task.progress.write().await.objects_scanned = 999; + assert_eq!(manager.get_task_report(&alias).await.expect("frozen report").progress, Some(expected)); + hook.finish.notify_one(); + COMPLETED_RETENTION_HOOKS + .lock() + .await + .retain(|key, _| key != &bucket && key != &task_id); + } +} + #[async_trait::async_trait] impl HealStorageAPI for MockStorage { async fn get_object_meta(&self, _bucket: &str, _object: &str) -> Result> { @@ -123,6 +443,12 @@ impl HealStorageAPI for MockStorage { } async fn object_exists(&self, bucket: &str, _object: &str) -> Result { + let hook = COMPLETED_RETENTION_HOOKS.lock().await.get(bucket).cloned(); + if let Some(hook) = hook { + hook.started.notify_one(); + hook.execute.notified().await; + return Ok(true); + } Ok(bucket == "retry-transition") } @@ -133,13 +459,18 @@ impl HealStorageAPI for MockStorage { _version_id: Option<&str>, _opts: &HealOpts, ) -> Result<(HealResultItem, Option)> { + if bucket == "completed-retention-failed" { + return Err(Error::TaskExecutionFailed { + message: "retention fixture failure".to_string(), + }); + } if let Some(hook) = manager_recovery_test_hook() { *hook .heal_object_calls .lock() .expect("manager recovery object call lock should not poison") += 1; } - if bucket == "retry-transition" { + if matches!(bucket, "retry-transition" | "completed-retention-retry-cancel") { return Ok(( HealResultItem::default(), Some(Error::Storage(EcstoreError::InsufficientReadQuorum( @@ -1145,7 +1476,13 @@ async fn test_active_duplicate_token_can_query_and_cancel_original_task() { .expect("duplicate token should cancel merged active task"); assert!(manager.active_heals.lock().await.get(&active_task_id).is_none()); - assert!(matches!(manager.get_task_status(&active_task_id).await, Err(Error::TaskNotFound { .. }))); + assert_eq!( + manager + .get_task_status(&active_task_id) + .await + .expect("cancelled task remains queryable"), + HealTaskStatus::Cancelled + ); } #[tokio::test] @@ -1638,6 +1975,8 @@ async fn insert_retrying_request(manager: &HealManager, request: HealRequest) -> manager.completed_heals.lock().await.insert( task_id, Arc::new(CompletedHealStatus { + progress: None, + retained_bytes: std::sync::OnceLock::new(), heal_type: request.heal_type, status: HealTaskStatus::Retrying { error: "Lock acquisition timeout".to_string(), @@ -2053,7 +2392,7 @@ async fn admin_force_start_cancels_overlapping_active_task_first() { "the overlapping admin task must be cancelled (removed from the active table) before the new one starts" ); assert!( - matches!(manager.get_task_status(&old_id).await, Err(Error::TaskNotFound { .. })), + matches!(manager.get_task_status(&old_id).await, Ok(HealTaskStatus::Cancelled)), "a cancelled task must no longer resolve as an active heal" ); } @@ -2360,6 +2699,8 @@ async fn test_retrying_completion_outranks_the_queue_for_the_same_id() { manager.completed_heals.lock().await.insert( task_id.clone(), Arc::new(CompletedHealStatus { + progress: None, + retained_bytes: std::sync::OnceLock::new(), heal_type: request.heal_type.clone(), status: HealTaskStatus::Retrying { error: "transient disk failure".to_string(), @@ -2395,6 +2736,8 @@ async fn test_get_task_status_reads_recent_completed_status() { manager.completed_heals.lock().await.insert( "completed-token".to_string(), Arc::new(CompletedHealStatus { + progress: None, + retained_bytes: std::sync::OnceLock::new(), heal_type: HealType::Bucket { bucket: "bucket".to_string(), }, @@ -2424,6 +2767,8 @@ async fn test_get_task_report_for_path_reads_completed_items() { manager.completed_heals.lock().await.insert( "completed-token".to_string(), Arc::new(CompletedHealStatus { + progress: None, + retained_bytes: std::sync::OnceLock::new(), heal_type: HealType::Object { bucket: "bucket".to_string(), object: "object".to_string(), diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index c95770834..c4c1d902b 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -999,7 +999,7 @@ impl HealTask { let items = match since { None => result_items.iter().map(|(_, item)| item.clone()).collect::>(), Some(cursor) => { - if cursor + 1 < min_seq { + if cursor.saturating_add(1) < min_seq { lagged = true; } result_items From e6bf2a464660ac5cd605fda090f6e73f8035efd0 Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 5 Sep 2026 16:45:45 +0800 Subject: [PATCH 3/5] fix(admin): report partial background heal coverage (#7178) * 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 Co-Authored-By: zhi22915 * fix(admin): report partial background heal coverage Refs rustfs/backlog#2035 and rustfs/backlog#2240. Co-Authored-By: heihutu Co-Authored-By: zhi22915 --------- Co-authored-by: heihutu Co-authored-by: zhi22915 --- crates/madmin/src/client.rs | 72 +++++++++++ rustfs/src/admin/handlers/heal.rs | 197 +++++++++++++++++++++++++++--- 2 files changed, 249 insertions(+), 20 deletions(-) diff --git a/crates/madmin/src/client.rs b/crates/madmin/src/client.rs index 83f6c33d7..b49fde280 100644 --- a/crates/madmin/src/client.rs +++ b/crates/madmin/src/client.rs @@ -182,6 +182,9 @@ pub struct BackgroundHealStatus { pub heal_active_tasks: u64, #[serde(default)] pub cluster_status_complete: bool, + /// Missing on older servers; absent coverage or counts mean unknown. + #[serde(default)] + pub coverage: Option, #[serde(default)] pub progress: Option, /// Remaining wire fields (flattened `BackgroundHealInfo` plus the @@ -190,6 +193,22 @@ pub struct BackgroundHealStatus { pub extra: serde_json::Map, } +/// Node coverage of a background heal status snapshot. Counters describe only +/// nodes with usable snapshots; unknown peers may still be running heal work. +#[derive(Debug, Clone, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct BackgroundHealCoverage { + #[serde(default)] + pub expected: Option, + #[serde(default)] + pub responded: Option, + #[serde(default)] + pub unknown: Option, + /// Stable reason codes; unknown future codes are preserved verbatim. + #[serde(default)] + pub reasons: Vec, +} + /// `GET /v3/scanner/status` response, typed at the fields operators branch /// on; everything else passes through verbatim. #[derive(Debug, Clone, Deserialize)] @@ -630,9 +649,39 @@ mod tests { assert_eq!(status.state, "active"); assert_eq!(status.heal_queue_length, 3); assert!(status.cluster_status_complete); + assert!(status.coverage.is_none(), "legacy payloads have unknown coverage"); assert!(status.extra.contains_key("healOperations"), "unknown nested payloads must pass through"); } + #[test] + fn background_heal_status_missing_coverage_fields_remain_unknown() { + for raw in [json!({"state": "degraded"}), json!({"state": "degraded", "coverage": {}})] { + let status: BackgroundHealStatus = serde_json::from_value(raw).expect("partial legacy payload decodes"); + assert!(!status.cluster_status_complete); + if let Some(coverage) = status.coverage { + assert_eq!(coverage.expected, None); + assert_eq!(coverage.responded, None); + assert_eq!(coverage.unknown, None); + } + } + } + + #[test] + fn background_heal_status_preserves_future_fields_and_reasons() { + let raw = json!({ + "state": "degraded", "clusterStatusComplete": false, + "coverage": {"expected": 3, "responded": 1, "unknown": 2, "reasons": ["future_reason"], "futureCoverage": true}, + "futureStatus": {"value": 7} + }); + let status: BackgroundHealStatus = serde_json::from_value(raw).expect("future additive fields decode"); + assert_eq!(status.extra["futureStatus"]["value"], 7); + let coverage = status.coverage.expect("coverage supplied"); + assert_eq!(coverage.expected, Some(3)); + assert_eq!(coverage.responded, Some(1)); + assert_eq!(coverage.unknown, Some(2)); + assert_eq!(coverage.reasons, ["future_reason"]); + } + #[test] fn scanner_status_defaults_freshness_to_unknown() { let raw = json!({"enabled": true, "freshness": {"state": "stale"}, "metrics": {}}); @@ -721,6 +770,7 @@ mod tests { let status = client.background_heal_status().await.expect("status decodes"); assert_eq!(status.state, "idle"); + assert!(status.coverage.is_none(), "older HTTP responses retain unknown coverage"); let request = server.recorded(); // The server registers this route POST-only; a GET here answers 405. assert_eq!(request.method, "POST"); @@ -728,6 +778,28 @@ mod tests { assert_eq!(request.query, ""); } + #[tokio::test] + async fn background_heal_status_decodes_partial_coverage_over_http() { + let body = r#"{"state":"degraded","healQueueLength":0,"healActiveTasks":0,"clusterStatusComplete":false,"coverage":{"expected":3,"responded":1,"unknown":2,"reasons":["notification_system_unavailable"]},"futureStatus":true}"#; + let server = TestServer::spawn(body, 200).await; + let client = AdminClient::new(&format!("http://{}", server.addr), "ak", "sk").expect("client builds"); + let status = client + .background_heal_status() + .await + .expect("partial status is a successful response"); + assert_eq!(status.state, "degraded"); + assert!(!status.cluster_status_complete); + assert_eq!(status.extra["futureStatus"], true); + let coverage = status.coverage.expect("partial coverage supplied"); + assert_eq!(coverage.expected, Some(3)); + assert_eq!(coverage.responded, Some(1)); + assert_eq!(coverage.unknown, Some(2)); + assert_eq!(coverage.reasons, ["notification_system_unavailable"]); + let request = server.recorded(); + assert_eq!(request.method, "POST"); + assert_eq!(request.query, "", "reading status must not send heal control parameters"); + } + #[tokio::test] async fn http_error_status_maps_to_a_typed_error_with_body() { let server = TestServer::spawn(r#"{"code":"AccessDenied","message":"denied"}"#, 403).await; diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index f10ad9123..74a2da980 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -39,7 +39,7 @@ use rustfs_utils::path::path_join; use s3s::header::{CONTENT_LENGTH, CONTENT_TYPE}; use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; use serde::{Deserialize, Serialize}; -use std::collections::{BTreeMap, HashSet}; +use std::collections::{BTreeMap, BTreeSet, HashSet}; use std::future::Future; use std::path::PathBuf; use std::sync::Arc; @@ -261,6 +261,7 @@ struct BackgroundHealStatus<'a> { heal_active_tasks: u64, heal_operations: rustfs_heal::HealOperationsSnapshot, cluster_status_complete: bool, + coverage: &'a BackgroundHealCoverage, #[serde(skip_serializing_if = "Option::is_none")] progress: Option, } @@ -300,6 +301,23 @@ fn background_heal_runtime_state( type BackgroundHealProgress = rustfs_heal::HealProgress; +#[derive(Debug, Serialize)] +struct BackgroundHealCoverage { + expected: usize, + responded: usize, + unknown: usize, + reasons: BTreeSet, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize)] +#[serde(rename_all = "snake_case")] +enum BackgroundHealCoverageReason { + NotificationSystemUnavailable, + PeerTopologyIncomplete, + PeerStatusUnsupported, + PeerStatusUnavailable, +} + #[derive(Debug)] struct ClusterHealStatusSnapshot { info: BackgroundHealInfo, @@ -307,6 +325,7 @@ struct ClusterHealStatusSnapshot { operations: rustfs_heal::HealOperationsSnapshot, progress: Option, complete: bool, + coverage: BackgroundHealCoverage, } fn add_priority_counts(total: &mut rustfs_heal::HealPriorityCounts, next: rustfs_heal::HealPriorityCounts) { @@ -338,6 +357,7 @@ fn add_operations(total: &mut rustfs_heal::HealOperationsSnapshot, next: rustfs_ } fn aggregate_cluster_heal_status(snapshots: Vec) -> ClusterHealStatusSnapshot { + let responded = snapshots.len(); let mut info = BackgroundHealInfo::default(); let mut operations = rustfs_heal::HealOperationsSnapshot::default(); let mut progress = Vec::new(); @@ -379,6 +399,12 @@ fn aggregate_cluster_heal_status(snapshots: Vec) -> Clus operations, progress, complete: true, + coverage: BackgroundHealCoverage { + expected: responded, + responded, + unknown: 0, + reasons: BTreeSet::new(), + }, } } @@ -413,12 +439,14 @@ fn merge_peer_heal_statuses( mut snapshots: Vec, peer_statuses: Vec, String>>, expected_nodes: usize, - topology_complete: bool, + coverage_reason: Option, ) -> S3Result { + let mut reasons: BTreeSet<_> = coverage_reason.into_iter().collect(); for peer_status in peer_statuses { match peer_status { Ok(Some(snapshot)) => snapshots.push(snapshot), Ok(None) => { + reasons.insert(BackgroundHealCoverageReason::PeerStatusUnsupported); warn!( event = EVENT_ADMIN_REQUEST_FAILED, component = LOG_COMPONENT_ADMIN_API, @@ -430,6 +458,7 @@ fn merge_peer_heal_statuses( ); } Err(err) => { + reasons.insert(BackgroundHealCoverageReason::PeerStatusUnavailable); warn!( event = EVENT_ADMIN_REQUEST_FAILED, component = LOG_COMPONENT_ADMIN_API, @@ -452,9 +481,12 @@ fn merge_peer_heal_statuses( // so during a reconfiguration the count can equal `expected_nodes` while // the topology is known-incomplete. Counting alone would report a // definitive answer precisely when the membership itself is in doubt. - let complete = topology_complete && snapshots.len() == expected_nodes; + let complete = reasons.is_empty() && snapshots.len() == expected_nodes; let mut status = aggregate_cluster_heal_status(snapshots); status.complete = complete; + status.coverage.expected = expected_nodes; + status.coverage.unknown = expected_nodes.saturating_sub(status.coverage.responded); + status.coverage.reasons = reasons; // A partial answer must never be mistakable for a definitive verdict: an // unreachable peer might be mid-heal, so reporting the reachable nodes' // "idle" (or disabled/uninitialized) as the cluster state would falsely @@ -497,7 +529,12 @@ async fn read_cluster_heal_status( return Ok(aggregate_cluster_heal_status(snapshots)); } let Some(notification_system) = notification_system else { - return Err(cluster_heal_status_unavailable("notification_system_unavailable")); + return merge_peer_heal_statuses( + snapshots, + Vec::new(), + expected_nodes, + Some(BackgroundHealCoverageReason::NotificationSystemUnavailable), + ); }; // An incomplete peer topology (a down member's client slot, a rolling // upgrade) previously failed the whole endpoint here, before any peer was @@ -540,7 +577,12 @@ async fn read_cluster_heal_status( })) .await; - merge_peer_heal_statuses(snapshots, peer_statuses, expected_nodes, topology_complete) + merge_peer_heal_statuses( + snapshots, + peer_statuses, + expected_nodes, + (!topology_complete).then_some(BackgroundHealCoverageReason::PeerTopologyIncomplete), + ) } async fn query_peer_replacement_recovery_status( @@ -1164,6 +1206,7 @@ fn encode_background_heal_status( heal_operations: rustfs_heal::HealOperationsSnapshot, progress: Option, cluster_status_complete: bool, + coverage: &BackgroundHealCoverage, ) -> S3Result> { let status = BackgroundHealStatus { info, @@ -1172,6 +1215,7 @@ fn encode_background_heal_status( heal_active_tasks: heal_operations.active_tasks, heal_operations, cluster_status_complete, + coverage, progress, }; serde_json::to_vec(&status).map_err(|e| { @@ -1461,6 +1505,7 @@ impl Operation for BackgroundHealStatusHandler { cluster_status.operations, cluster_status.progress, cluster_status.complete, + &cluster_status.coverage, )?; info!( event = EVENT_ADMIN_RESPONSE_EMITTED, @@ -1515,13 +1560,14 @@ impl Operation for ReplacementRecoveryStatusHandler { mod tests { use super::extract_heal_init_params; use super::{ - BackgroundHealProgress, HealInitParams, HealResp, HealRuntimeState, aggregate_cluster_heal_status, - aggregate_replacement_recovery_cluster_status, background_heal_runtime_state, build_heal_channel_request, - build_replacement_recovery_status_response, encode_background_heal_status, encode_heal_control_path, - encode_heal_start_success, encode_heal_task_status, execute_after_heal_control_capability, heal_channel_response_items, - heal_channel_response_progress, heal_channel_response_summary, heal_control_response_id, json_response, - map_heal_response, merge_peer_heal_statuses, peer_topology_complete, query_peer_heal_status, - query_peer_replacement_recovery_status, reject_heal_admission, validate_heal_request_mode, validate_heal_target, + BackgroundHealCoverage, BackgroundHealCoverageReason, BackgroundHealProgress, HealInitParams, HealResp, HealRuntimeState, + aggregate_cluster_heal_status, aggregate_replacement_recovery_cluster_status, background_heal_runtime_state, + build_heal_channel_request, build_replacement_recovery_status_response, encode_background_heal_status, + encode_heal_control_path, encode_heal_start_success, encode_heal_task_status, execute_after_heal_control_capability, + heal_channel_response_items, heal_channel_response_progress, heal_channel_response_summary, heal_control_response_id, + json_response, map_heal_response, merge_peer_heal_statuses, peer_topology_complete, query_peer_heal_status, + query_peer_replacement_recovery_status, read_cluster_heal_status, reject_heal_admission, validate_heal_request_mode, + validate_heal_target, }; use crate::storage::rpc::node_service::heal::{ NodeHealProgress, NodeHealStatusSnapshot, NodeReplacementRecoveryStatusSnapshot, encode_node_replacement_recovery_status, @@ -2175,7 +2221,13 @@ mod tests { ..Default::default() }; - let encoded = encode_background_heal_status(&info, HealRuntimeState::Active, operations, None, true) + let coverage = BackgroundHealCoverage { + expected: 1, + responded: 1, + unknown: 0, + reasons: Default::default(), + }; + let encoded = encode_background_heal_status(&info, HealRuntimeState::Active, operations, None, true, &coverage) .expect("background heal info should serialize"); let json: serde_json::Value = serde_json::from_slice(&encoded).expect("json should deserialize"); @@ -2225,6 +2277,12 @@ mod tests { rustfs_heal::HealOperationsSnapshot::default(), Some(progress), true, + &BackgroundHealCoverage { + expected: 1, + responded: 1, + unknown: 0, + reasons: Default::default(), + }, ) .expect("background heal info should serialize"); let json: serde_json::Value = serde_json::from_slice(&encoded).expect("json should deserialize"); @@ -2398,6 +2456,84 @@ mod tests { assert!(peer_topology_complete(1, 0, 0, 1, 0)); } + #[tokio::test] + async fn test_background_heal_status_without_notification_preserves_local_snapshot() { + let initialized = rustfs_heal::heal_runtime_initialized(); + let info = BackgroundHealInfo { + bitrot_start_cycle: 37, + current_scan_mode: HealScanMode::Deep, + ..Default::default() + }; + for expected in [1, 3] { + let status = tokio::time::timeout(Duration::from_secs(1), read_cluster_heal_status(info.clone(), None, expected)) + .await + .expect("local status must not wait for remote peers") + .expect("missing notification must retain the local snapshot"); + assert_eq!(status.info.bitrot_start_cycle, 37); + assert_eq!(status.info.current_scan_mode, HealScanMode::Deep); + assert_eq!(status.complete, expected == 1); + assert_eq!(status.coverage.expected, expected); + assert_eq!(status.coverage.responded, 1); + assert_eq!(status.coverage.unknown, expected - 1); + if expected == 1 { + assert!(status.coverage.reasons.is_empty()); + } else { + assert!(matches!(status.state, HealRuntimeState::Degraded | HealRuntimeState::Active)); + assert_eq!( + status.coverage.reasons, + [BackgroundHealCoverageReason::NotificationSystemUnavailable].into() + ); + } + let encoded = encode_background_heal_status( + &status.info, + status.state, + status.operations, + status.progress, + status.complete, + &status.coverage, + ) + .expect("fallback status must encode"); + let decoded: rustfs_madmin::client::BackgroundHealStatus = + serde_json::from_slice(&encoded).expect("the actual madmin client must decode the server response"); + assert_eq!(decoded.cluster_status_complete, expected == 1); + let coverage = decoded.coverage.expect("new server supplies coverage"); + assert_eq!(coverage.expected, Some(expected)); + assert_eq!(coverage.responded, Some(1)); + assert_eq!(coverage.unknown, Some(expected - 1)); + if expected > 1 { + assert_eq!(coverage.reasons, ["notification_system_unavailable"]); + } + } + assert_eq!( + rustfs_heal::heal_runtime_initialized(), + initialized, + "reading status must not initialize heal" + ); + } + + #[test] + fn test_background_heal_status_coverage_reasons_are_bounded() { + let local = NodeHealStatusSnapshot::for_test(true, true, BackgroundHealInfo::default(), Default::default(), None); + let peers = (0..100) + .map(|index| { + if index % 2 == 0 { + Ok(None) + } else { + Err("peer unavailable".to_owned()) + } + }) + .collect(); + let status = merge_peer_heal_statuses(vec![local], peers, 101, None).expect("local status remains available"); + assert_eq!(status.coverage.responded, 1); + assert_eq!(status.coverage.unknown, 100); + assert_eq!(status.coverage.reasons.len(), 2); + let encoded = serde_json::to_vec(&status.coverage).expect("coverage encodes"); + assert!(encoded.len() < 256, "coverage must not grow with peer failures"); + let decoded: rustfs_madmin::client::BackgroundHealCoverage = + serde_json::from_slice(&encoded).expect("client coverage decodes"); + assert_eq!(decoded.reasons, ["peer_status_unsupported", "peer_status_unavailable"]); + } + #[test] fn test_peer_status_merge_degrades_explicitly_and_never_claims_idle() { let local = || { @@ -2413,15 +2549,20 @@ mod tests { // but the safety property of the previous fail-closed behaviour is // preserved: the partial answer is labelled Degraded, never Idle, so // unknown peer work cannot be mistaken for "nothing is running". - let partial = merge_peer_heal_statuses(vec![local()], vec![Err("peer timeout".to_string())], 2, true) + let partial = merge_peer_heal_statuses(vec![local()], vec![Err("peer timeout".to_string())], 2, None) .expect("an unreachable peer degrades the answer instead of destroying it"); assert!(!partial.complete); assert_eq!(partial.state, HealRuntimeState::Degraded); + assert_eq!(partial.coverage.expected, 2); + assert_eq!(partial.coverage.responded, 1); + assert_eq!(partial.coverage.unknown, 1); + assert_eq!(partial.coverage.reasons, [BackgroundHealCoverageReason::PeerStatusUnavailable].into()); - let older_peer = merge_peer_heal_statuses(vec![local()], vec![Ok(None)], 2, true) + let older_peer = merge_peer_heal_statuses(vec![local()], vec![Ok(None)], 2, None) .expect("an older peer degrades the answer instead of destroying it"); assert!(!older_peer.complete); assert_eq!(older_peer.state, HealRuntimeState::Degraded); + assert_eq!(older_peer.coverage.reasons, [BackgroundHealCoverageReason::PeerStatusUnsupported].into()); let known_active = NodeHealStatusSnapshot::for_test( true, @@ -2433,12 +2574,12 @@ mod tests { }, None, ); - let partial_active = merge_peer_heal_statuses(vec![known_active], vec![Ok(None)], 2, true) + let partial_active = merge_peer_heal_statuses(vec![known_active], vec![Ok(None)], 2, None) .expect("known active work may be reported as an explicit partial status"); assert!(!partial_active.complete); assert_eq!(partial_active.state, HealRuntimeState::Active); - merge_peer_heal_statuses(Vec::new(), vec![Err("peer timeout".to_string())], 2, true) + merge_peer_heal_statuses(Vec::new(), vec![Err("peer timeout".to_string())], 2, None) .expect_err("no snapshot at all still fails closed"); } @@ -2459,12 +2600,22 @@ mod tests { None, ) }; - let full_count_incomplete_topology = merge_peer_heal_statuses(vec![snapshot()], vec![Ok(Some(snapshot()))], 2, false) - .expect("incomplete topology degrades the answer instead of destroying it"); + let full_count_incomplete_topology = merge_peer_heal_statuses( + vec![snapshot()], + vec![Ok(Some(snapshot()))], + 2, + Some(BackgroundHealCoverageReason::PeerTopologyIncomplete), + ) + .expect("incomplete topology degrades the answer instead of destroying it"); assert!(!full_count_incomplete_topology.complete); assert_eq!(full_count_incomplete_topology.state, HealRuntimeState::Degraded); + assert_eq!(full_count_incomplete_topology.coverage.unknown, 0); + assert_eq!( + full_count_incomplete_topology.coverage.reasons, + [BackgroundHealCoverageReason::PeerTopologyIncomplete].into() + ); - let full_count_complete_topology = merge_peer_heal_statuses(vec![snapshot()], vec![Ok(Some(snapshot()))], 2, true) + let full_count_complete_topology = merge_peer_heal_statuses(vec![snapshot()], vec![Ok(Some(snapshot()))], 2, None) .expect("complete topology and full count is a definitive answer"); assert!(full_count_complete_topology.complete); assert_eq!(full_count_complete_topology.state, HealRuntimeState::Idle); @@ -2478,6 +2629,12 @@ mod tests { rustfs_heal::HealOperationsSnapshot::default(), None, false, + &BackgroundHealCoverage { + expected: 2, + responded: 1, + unknown: 1, + reasons: [BackgroundHealCoverageReason::PeerStatusUnavailable].into(), + }, ) .expect("degraded status must serialize"); let json: serde_json::Value = serde_json::from_slice(&encoded).expect("valid json"); From 42c32381b60d09a46e38554a37664352e6c66ec8 Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 5 Sep 2026 16:45:53 +0800 Subject: [PATCH 4/5] fix(scanner): make reset cleanup safely reentrant (#7180) * 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 Co-Authored-By: zhi22915 * fix(scanner): make reset cleanup safely reentrant Refs rustfs/backlog#2264 and rustfs/backlog#2240. Co-Authored-By: heihutu Co-Authored-By: zhi22915 --------- Co-authored-by: heihutu Co-authored-by: zhi22915 --- crates/scanner/src/scanner/cycle_state.rs | 331 ++++++++++++++++--- crates/scanner/src/scanner/leadership.rs | 8 +- crates/scanner/src/scanner/tests.rs | 380 +++++++++++++++++++++- 3 files changed, 664 insertions(+), 55 deletions(-) diff --git a/crates/scanner/src/scanner/cycle_state.rs b/crates/scanner/src/scanner/cycle_state.rs index dc29e3331..6913ed836 100644 --- a/crates/scanner/src/scanner/cycle_state.rs +++ b/crates/scanner/src/scanner/cycle_state.rs @@ -379,6 +379,12 @@ pub(super) fn decode_recovery_marker_for_reset( if !matches!(marker_revision, DataUsageCacheRevision::Etag(_)) { return Err(ScannerError::Other("cycle recovery marker has no object revision".to_string())); } + if let Ok(value) = serde_json::from_slice::(data) + && let Some(state) = value.get("state") + && !matches!(state.as_str(), Some("blocked" | "cleanup-pending")) + { + return Err(ScannerError::Other("cycle recovery marker state is unsupported".to_string())); + } let compat = serde_json::from_slice::(data).ok(); let _schema_version = compat.as_ref().and_then(|marker| marker.schema_version); let primary_revision = compat @@ -406,7 +412,10 @@ pub(super) fn decode_recovery_marker_for_reset( }; let state = match compat.as_ref().and_then(|marker| marker.state.as_deref()) { Some("cleanup-pending") => "cleanup-pending", - _ => "blocked", + Some("blocked") | None => "blocked", + Some(_) => { + return Err(ScannerError::Other("cycle recovery marker state is unsupported".to_string())); + } }; let now = unix_now_secs(); Ok(ScannerCycleRecoveryMarker { @@ -721,17 +730,19 @@ async fn mark_cycle_recovery_cleanup_pending( mut marker: ScannerCycleRecoveryMarker, marker_revision: &DataUsageCacheRevision, expected_epoch: u64, + owns_reset: &(impl Fn() -> bool + Sync), ) -> Result<(ScannerCycleRecoveryMarker, DataUsageCacheRevision), ScannerError> { marker.state = "cleanup-pending".to_string(); marker.last_attempt_at_unix_secs = unix_now_secs(); let bytes = serde_json::to_vec(&marker) .map_err(|err| ScannerError::Other(format!("failed to encode cycle recovery marker: {err}")))?; - let info = save_config_with_publication_admission_for_epoch( + let info = save_reset_config( storeapi.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), bytes, marker_revision.preconditions(), expected_epoch, + owns_reset, ) .await .map_err(|err| ScannerError::Other(format!("failed to mark cycle recovery cleanup pending: {err}")))?; @@ -933,6 +944,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< .get_write_lock_quiet(Duration::from_secs(5)) .await .map_err(|err| ScannerError::Other(format!("scanner leader lock is busy: {err}")))?; + let owns_reset = || !guard.is_lock_lost() && !ctx.is_cancelled(); if guard.is_lock_lost() { return Err(ScannerError::Other("scanner leader lock was lost before recovery reset".to_string())); @@ -952,7 +964,27 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< } Err(err) => return Err(ScannerError::Other(format!("failed to read cycle recovery marker: {err}"))), }; - let marker_data = marker_data.ok_or_else(|| ScannerError::Other("scanner cycle recovery marker is absent".to_string()))?; + let Some(marker_data) = marker_data else { + // A delete may commit before its reply is lost. Confirm both durable + // fences before treating a retry without its marker as completed. + let (cycle, epoch, revision) = read_cycle_state_for_usage_reset(storeapi.clone()).await?; + let floor = persisted_usage_floor(storeapi.clone()).await?; + if !matches!(revision, DataUsageCacheRevision::Etag(_)) + || epoch < floor.leader_epoch + || cycle.next < floor.next_cycle + || !owns_reset() + || scanner_publication_admission_for_epoch(storeapi.clone(), reset_epoch) + .await + .is_none() + { + return Err(ScannerError::Other( + "scanner cycle recovery marker is absent without a completed reset fence".to_string(), + )); + } + set_scanner_cycle_recovery_status(recovery_status("healthy", None, false)); + super::notify_scanner_cycle_recovery_wake(); + return Ok(()); + }; let (marker, force_full_rescan) = match serde_json::from_slice::(&marker_data) { Ok(marker) if validate_recovery_marker(&marker).is_ok() => (marker, false), _ => (decode_recovery_marker_for_reset(&marker_data, &marker_revision)?, true), @@ -1026,8 +1058,10 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< } }; if let Some((primary_cycle, primary_epoch)) = primary_state { + verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?; let (cleanup_marker, cleanup_marker_revision) = - mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision, reset_epoch).await?; + mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision, reset_epoch, &owns_reset) + .await?; set_scanner_cycle_recovery_status(recovery_status_from_marker(&cleanup_marker, "cleanup-pending")); let usage_floor = persisted_usage_floor(storeapi.clone()).await?; let fence_epoch = primary_epoch @@ -1047,12 +1081,14 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< "preserved scanner cycle state exceeds the bounded object size".to_string(), )); } - let preserved_info = save_config_with_publication_admission_for_epoch( + verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?; + let preserved_info = save_reset_config( storeapi.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), preserved_data, primary_revision.preconditions(), reset_epoch, + &owns_reset, ) .await .map_err(|err| { @@ -1072,9 +1108,17 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< "scanner leader lock was lost after fencing newer cycle state".to_string(), )); } - fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), fence_epoch, Some(reset_epoch), false) - .await - .map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?; + verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?; + fence_scanner_usage_epoch_with_expected_epoch( + &ctx, + storeapi.clone(), + fence_epoch, + Some(reset_epoch), + false, + &owns_reset, + ) + .await + .map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?; if guard.is_lock_lost() { return Err(ScannerError::Other( "scanner leader lock was lost after fencing newer cycle state".to_string(), @@ -1088,7 +1132,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< "scanner cycle state changed before recovery marker cleanup".to_string(), )); } - delete_config_with_publication_admission_for_epoch( + verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?; + delete_reset_config( storeapi.clone(), RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), @@ -1100,6 +1145,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< ..Default::default() }, reset_epoch, + &owns_reset, ) .await .map_err(|err| { @@ -1149,17 +1195,20 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< // Persist the cleanup-pending phase before rewriting the primary. If the // process dies after the rewrite, startup still sees a durable fence and // cannot mistake the partially completed reset for a healthy state. + verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?; let (marker, marker_revision) = if marker.state == "cleanup-pending" { (marker, marker_revision) } else { - mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision, reset_epoch).await? + mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision, reset_epoch, &owns_reset).await? }; - let rebuilt_info = save_config_with_publication_admission_for_epoch( + verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?; + let rebuilt_info = save_reset_config( storeapi.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), data, primary_revision.preconditions(), reset_epoch, + &owns_reset, ) .await .map_err(|err| { @@ -1178,8 +1227,10 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< "scanner leader lock was lost after rebuilding cycle state".to_string(), )); } + verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?; if let Err(err) = - fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false).await + fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false, &owns_reset) + .await { set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { path: DATA_USAGE_BLOOM_NAME_PATH.clone(), @@ -1249,7 +1300,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< )); } - if let Err(err) = delete_config_with_publication_admission_for_epoch( + verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?; + if let Err(err) = delete_reset_config( storeapi.clone(), RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), @@ -1261,6 +1313,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< ..Default::default() }, reset_epoch, + &owns_reset, ) .await { @@ -1310,6 +1363,57 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< Ok(()) } +async fn verify_cycle_reset_intent( + storeapi: Arc, + expected_revision: &DataUsageCacheRevision, + owns_reset: &(impl Fn() -> bool + Sync), +) -> Result<(), ScannerError> { + let revision = read_config_revision(storeapi, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to verify scanner cycle reset intent: {err}")))?; + if &revision != expected_revision { + return Err(ScannerError::Other("scanner cycle reset intent changed".to_string())); + } + if !owns_reset() { + return Err(ScannerError::Other("scanner cycle reset ownership was lost".to_string())); + } + Ok(()) +} + +async fn save_reset_config( + storeapi: Arc, + path: &str, + data: Vec, + preconditions: crate::HTTPPreconditions, + expected_epoch: u64, + owns_reset: &(impl Fn() -> bool + Sync), +) -> Result { + let Some(_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_epoch).await else { + return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED)); + }; + if !owns_reset() { + return Err(EcstoreError::other("scanner reset ownership was lost before write")); + } + save_config_with_preconditions(storeapi, path, data, preconditions).await +} + +async fn delete_reset_config( + storeapi: Arc, + bucket: &str, + path: &str, + options: ScannerObjectOptions, + expected_epoch: u64, + owns_reset: &(impl Fn() -> bool + Sync), +) -> Result { + let Some(_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_epoch).await else { + return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED)); + }; + if !owns_reset() { + return Err(EcstoreError::other("scanner reset ownership was lost before delete")); + } + storeapi.delete_config_object(bucket, path, options).await +} + fn scanner_usage_state_reset_paths() -> Vec { vec![ DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(), @@ -1333,8 +1437,14 @@ pub(super) async fn read_usage_state_reset_slots( Ok(slots) } -fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result { - let mut floor = PersistedUsageFloor::default(); +enum ScannerUsageResetFloor { + Missing, + Trusted(PersistedUsageFloor), + Corrupt, +} + +fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result { + let mut floor = None; for slot in slots { let Some(data) = slot.data.as_deref() else { continue; @@ -1342,9 +1452,21 @@ fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result(data) else { continue; }; - update_persisted_usage_floor(&mut floor, &usage, &slot.path)?; + if !data_usage_info_has_persisted_baseline_identity(&usage) + && !(slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str() && data_usage_info_is_bootstrap_pending(&usage)) + && legacy_incomplete_usage_fence(data, &usage) + .and_then(|fence| fence.claimable_epoch()) + .is_none() + { + continue; + } + update_persisted_usage_floor(floor.get_or_insert_with(PersistedUsageFloor::default), &usage, &slot.path)?; } - Ok(floor) + Ok(match floor { + Some(floor) => ScannerUsageResetFloor::Trusted(floor), + None if slots.iter().any(|slot| slot.data.is_some()) => ScannerUsageResetFloor::Corrupt, + None => ScannerUsageResetFloor::Missing, + }) } async fn read_cycle_state_for_usage_reset( @@ -1401,11 +1523,12 @@ async fn delete_usage_state_reset_slot( storeapi: Arc, slot: &ScannerUsageStateResetSlot, expected_epoch: u64, + owns_reset: &(impl Fn() -> bool + Sync), ) -> Result { if matches!(slot.revision, DataUsageCacheRevision::Missing) { return Ok(false); } - let delete_result = delete_config_with_publication_admission_for_epoch( + let delete_result = delete_reset_config( storeapi.clone(), RUSTFS_META_BUCKET, &slot.path, @@ -1415,6 +1538,7 @@ async fn delete_usage_state_reset_slot( ..Default::default() }, expected_epoch, + owns_reset, ) .await; match delete_result { @@ -1486,21 +1610,24 @@ pub(super) async fn publish_scanner_usage_bootstrap_primary( expected_publication_epoch: u64, leader_epoch: Option, context: ScannerUsageBootstrapPublishContext, + owns_publication: impl Fn() -> bool + Sync, ) -> Result<(), ScannerError> { async fn inner( storeapi: Arc, expected_revision: &DataUsageCacheRevision, expected_publication_epoch: u64, leader_epoch: Option, + owns_publication: &(impl Fn() -> bool + Sync), ) -> Result<(), ScannerUsageBootstrapPublishError> { let marker = scanner_usage_bootstrap_marker(std::time::SystemTime::now(), leader_epoch); let data = serde_json::to_vec(&marker).map_err(ScannerUsageBootstrapPublishError::Encode)?; - let save_result = save_config_with_publication_admission_for_epoch( + let save_result = save_reset_config( storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), data.clone(), expected_revision.preconditions(), expected_publication_epoch, + owns_publication, ) .await; if save_result @@ -1524,7 +1651,7 @@ pub(super) async fn publish_scanner_usage_bootstrap_primary( }) } - inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch) + inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch, &owns_publication) .await .map_err(|err| err.into_scanner_error(context)) } @@ -1534,32 +1661,108 @@ pub(super) async fn reset_scanner_usage_state_slots_for_full_rebuild( slots: &[ScannerUsageStateResetSlot], expected_epoch: u64, leader_epoch: u64, + owns_reset: impl Fn() -> bool + Sync, ) -> Result, ScannerError> { let mut reset_paths = Vec::new(); let primary = slots .iter() .find(|slot| slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str()) .ok_or_else(|| ScannerError::Other("scanner usage reset primary slot was not inspected".to_string()))?; - publish_scanner_usage_bootstrap_primary( - storeapi.clone(), - &primary.revision, - expected_epoch, - Some(leader_epoch), - ScannerUsageBootstrapPublishContext::Reset, - ) - .await?; + if !owns_reset() { + return Err(ScannerError::Other("scanner usage reset ownership was lost".to_string())); + } + let resume_epoch = usage_state_reset_resume_epoch(slots)?; + match resume_epoch { + Some(epoch) if epoch == leader_epoch => {} + Some(_) => return Err(ScannerError::Other("scanner usage reset bootstrap epoch changed".to_string())), + None => { + publish_scanner_usage_bootstrap_primary( + storeapi.clone(), + &primary.revision, + expected_epoch, + Some(leader_epoch), + ScannerUsageBootstrapPublishContext::Reset, + &owns_reset, + ) + .await?; + } + } + let (data, intent_revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to inspect scanner usage reset intent: {err}")))?; + data.as_deref() + .and_then(|data| serde_json::from_slice::(data).ok()) + .filter(|usage| data_usage_info_is_bootstrap_pending(usage) && usage.scanner_epoch == Some(leader_epoch)) + .ok_or_else(|| ScannerError::Other("scanner usage reset intent changed before cleanup".to_string()))?; + if !matches!(intent_revision, DataUsageCacheRevision::Etag(_)) + || (resume_epoch.is_some() && intent_revision != primary.revision) + { + return Err(ScannerError::Other("scanner usage reset intent revision changed".to_string())); + } reset_paths.push(DATA_USAGE_OBJ_NAME_PATH.as_str().to_string()); for slot in slots.iter().filter(|slot| slot.path != DATA_USAGE_OBJ_NAME_PATH.as_str()) { - if delete_usage_state_reset_slot(storeapi.clone(), slot, expected_epoch).await? { + if let Some(usage) = slot + .data + .as_deref() + .and_then(|data| serde_json::from_slice::(data).ok()) + && usage_epoch(&usage) >= leader_epoch + { + return Err(ScannerError::Other(format!( + "scanner usage reset slot is not older than its intent: {}", + slot.path + ))); + } + let revision = read_config_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to verify scanner usage reset intent: {err}")))?; + if revision != intent_revision { + return Err(ScannerError::Other("scanner usage reset intent changed during cleanup".to_string())); + } + if !owns_reset() { + return Err(ScannerError::Other("scanner usage reset ownership was lost".to_string())); + } + if delete_usage_state_reset_slot(storeapi.clone(), slot, expected_epoch, &owns_reset).await? { reset_paths.push(slot.path.clone()); } } + let revision = read_config_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to confirm scanner usage reset intent: {err}")))?; + if revision != intent_revision || !owns_reset() { + return Err(ScannerError::Other( + "scanner usage reset intent or ownership changed before completion".to_string(), + )); + } invalidate_admin_data_usage_snapshot_cache().await; invalidate_data_usage_snapshot_cache().await; Ok(reset_paths) } +fn usage_state_reset_resume_epoch(slots: &[ScannerUsageStateResetSlot]) -> Result, ScannerError> { + let primary = slots.iter().find(|slot| slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str()); + let usage = primary + .and_then(|slot| slot.data.as_deref()) + .and_then(|data| serde_json::from_slice::(data).ok()); + match usage { + Some(usage) if usage.usage_snapshot_bootstrap_pending => { + if !data_usage_info_is_bootstrap_pending(&usage) { + return Err(ScannerError::Other("scanner usage reset bootstrap is invalid".to_string())); + } + if usage.scanner_epoch.is_none() { + // Initial bootstrap has no reset owner yet. + return Ok(None); + } + usage + .scanner_epoch + .filter(|epoch| *epoch > 0 && *epoch < u64::MAX) + .map(Some) + .ok_or_else(|| ScannerError::Other("scanner usage reset bootstrap has no valid epoch".to_string())) + } + _ => Ok(None), + } +} + pub async fn reset_scanner_usage_state_for_full_rebuild( ctx: CancellationToken, storeapi: Arc, @@ -1584,12 +1787,31 @@ pub async fn reset_scanner_usage_state_for_full_rebuild( }; let (cycle, cycle_epoch, cycle_revision) = read_cycle_state_for_usage_reset(storeapi.clone()).await?; let slots = read_usage_state_reset_slots(storeapi.clone()).await?; - let usage_floor = usage_state_reset_floor(&slots)?; - let leader_epoch = cycle_epoch - .max(usage_floor.leader_epoch) - .checked_add(1) - .filter(|epoch| *epoch < u64::MAX) - .ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?; + let usage_floor = match usage_state_reset_floor(&slots)? { + ScannerUsageResetFloor::Trusted(floor) => floor, + ScannerUsageResetFloor::Corrupt if matches!(cycle_revision, DataUsageCacheRevision::Missing) => { + return Err(ScannerError::Other("scanner usage reset has no trusted cycle or usage floor".to_string())); + } + ScannerUsageResetFloor::Missing | ScannerUsageResetFloor::Corrupt => PersistedUsageFloor { + next_cycle: cycle.next, + leader_epoch: cycle_epoch, + }, + }; + let resume_epoch = usage_state_reset_resume_epoch(&slots)?; + let leader_epoch = if let Some(epoch) = resume_epoch { + if epoch != cycle_epoch || usage_floor.leader_epoch > epoch || usage_floor.next_cycle > cycle.next { + return Err(ScannerError::Other( + "scanner usage reset bootstrap conflicts with the persisted cycle fence".to_string(), + )); + } + epoch + } else { + cycle_epoch + .max(usage_floor.leader_epoch) + .checked_add(1) + .filter(|epoch| *epoch < u64::MAX) + .ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))? + }; let rebuilt_cycle = CurrentCycle { next: cycle.next.max(usage_floor.next_cycle), ..Default::default() @@ -1602,21 +1824,24 @@ pub async fn reset_scanner_usage_state_for_full_rebuild( "scanner leader lock was lost before fencing usage reset cycle state".to_string(), )); } - save_config_with_publication_admission_for_epoch( - storeapi.clone(), - DATA_USAGE_BLOOM_NAME_PATH.as_str(), - cycle_data, - cycle_revision.preconditions(), - reset_epoch, - ) - .await - .map_err(|err| { - if scanner_publication_epoch_changed(&err) { - ScannerError::Other("scanner usage reset deferred by a movement epoch change".to_string()) - } else { - ScannerError::Other(format!("failed to fence scanner cycle state for usage reset: {err}")) - } - })?; + if resume_epoch.is_none() { + save_reset_config( + storeapi.clone(), + DATA_USAGE_BLOOM_NAME_PATH.as_str(), + cycle_data, + cycle_revision.preconditions(), + reset_epoch, + &|| !guard.is_lock_lost() && !ctx.is_cancelled(), + ) + .await + .map_err(|err| { + if scanner_publication_epoch_changed(&err) { + ScannerError::Other("scanner usage reset deferred by a movement epoch change".to_string()) + } else { + ScannerError::Other(format!("failed to fence scanner cycle state for usage reset: {err}")) + } + })?; + } if guard.is_lock_lost() { return Err(ScannerError::Other( @@ -1624,7 +1849,10 @@ pub async fn reset_scanner_usage_state_for_full_rebuild( )); } let reset_paths = - reset_scanner_usage_state_slots_for_full_rebuild(storeapi.clone(), &slots, reset_epoch, leader_epoch).await?; + reset_scanner_usage_state_slots_for_full_rebuild(storeapi.clone(), &slots, reset_epoch, leader_epoch, || { + !guard.is_lock_lost() && !ctx.is_cancelled() + }) + .await?; if guard.is_lock_lost() { return Err(ScannerError::Other( "scanner leader lock was lost after publishing usage reset marker".to_string(), @@ -2135,6 +2363,7 @@ async fn recover_legacy_incomplete_usage_floor( expected_publication_epoch, Some(primary.epoch), ScannerUsageBootstrapPublishContext::Recovery, + || true, ) .await?; warn!( diff --git a/crates/scanner/src/scanner/leadership.rs b/crates/scanner/src/scanner/leadership.rs index 8283a987a..260476e7a 100644 --- a/crates/scanner/src/scanner/leadership.rs +++ b/crates/scanner/src/scanner/leadership.rs @@ -191,6 +191,7 @@ pub(super) async fn initialize_usage_baseline_bootstrap( expected_epoch, None, ScannerUsageBootstrapPublishContext::Initial, + || true, ) .await } @@ -201,9 +202,10 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch( claimed_epoch: u64, expected_publication_epoch: Option, allow_bootstrap_pending: bool, + owns_fence: impl Fn() -> bool, ) -> Result<(), ScannerError> { for retry in 0..=SCANNER_PERSIST_CAS_RETRIES { - if ctx.is_cancelled() { + if ctx.is_cancelled() || !owns_fence() { return Err(ScannerError::Other("scanner leadership was cancelled before usage fencing".to_string())); } @@ -264,6 +266,9 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch( "scanner usage epoch fence changed while preparing its conditional write".to_string(), )); }; + if ctx.is_cancelled() || !owns_fence() { + return Err(ScannerError::Other("scanner leadership was lost before usage fencing".to_string())); + } save_config_with_preconditions(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), data, revision.preconditions()) .await }; @@ -319,6 +324,7 @@ pub(super) async fn complete_scanner_leadership_claim( claimed_epoch, expected_publication_epoch, allow_bootstrap_pending, + || true, ) .await { diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 244e9088f..b9b4cd299 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -634,6 +634,8 @@ struct MemoryConfigStore { cancel_after_successful_puts: Mutex>, replace_after_successful_puts: Mutex)>>, error_after_commit_deletes: Mutex>, + cancel_after_deletes: Mutex>, + pause_next_publication_admission: Mutex, Arc)>>, put_counts: Mutex>, publication_admission_blocked: AtomicBool, block_publication_after_admissions: AtomicUsize, @@ -4081,6 +4083,9 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore { revisions.remove(&key); drop(revisions); drop(objects); + if let Some(token) = self.cancel_after_deletes.lock().await.remove(&key) { + token.cancel(); + } if self.error_after_commit_deletes.lock().await.remove(&key) { return Err(EcstoreError::other("injected delete error after commit")); } @@ -4088,6 +4093,11 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore { } async fn scanner_data_usage_publication_admission(&self) -> Option { + let pause = self.pause_next_publication_admission.lock().await.take(); + if let Some((entered, resume)) = pause { + entered.notify_one(); + resume.notified().await; + } if self.publication_admission_blocked.load(Ordering::Acquire) { return None; } @@ -4589,7 +4599,7 @@ async fn scanner_legacy_usage_backup_survives_fencing_and_restart_after_real_met .expect("publication must also read the intact backup"); assert_eq!(baseline.data.as_deref(), Some(data.as_slice())); assert_eq!(baseline.revision, DataUsageCacheRevision::Missing); - fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false) + fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false, || true) .await .expect("legacy backup must be fenced into v2"); let fenced = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) @@ -4818,6 +4828,28 @@ async fn scanner_usage_state_reset_publishes_fenced_bootstrap_marker() { ); } + let cycle_before_retry = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("cycle should remain before retry"); + let marker_before_retry = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("bootstrap should remain before retry"); + let retry = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone()) + .await + .expect("completed cleanup should be reentrant"); + assert_eq!(retry.leader_epoch, result.leader_epoch); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("cycle should remain"), + cycle_before_retry + ); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("bootstrap should remain"), + marker_before_retry + ); let (floor, state) = persisted_usage_floor_for_startup(store, false) .await .expect("reset marker should be resumable"); @@ -4973,7 +5005,7 @@ async fn scanner_usage_state_reset_slots_reject_primary_aba() { store.objects.lock().await.insert(key.clone(), b"newer-json".to_vec()); store.revisions.lock().await.insert(key, 2); - let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3) + let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3, || true) .await .expect_err("stale primary revision must not be overwritten"); assert!( @@ -4983,6 +5015,348 @@ async fn scanner_usage_state_reset_slots_reject_primary_aba() { ); } +#[tokio::test] +async fn scanner_usage_state_reset_resumes_every_cleanup_boundary_without_rewriting_intent() { + for completed in 0..=4 { + let store = Arc::new(MemoryConfigStore::default()); + let primary_path = DATA_USAGE_OBJ_NAME_PATH.as_str(); + let cleanup_paths = [ + format!("{primary_path}.bkp"), + LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(), + format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()), + DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str().to_string(), + ]; + for path in std::iter::once(primary_path).chain(cleanup_paths.iter().map(String::as_str)) { + let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + usage.scanner_epoch = Some(1); + save_config(store.clone(), path, serde_json::to_vec(&usage).expect("fixture should encode")) + .await + .expect("fixture should persist"); + } + // These objects belong to other owners, even when reset cleanup resumes. + for path in ["buckets/quota-reservations/ledger", "buckets/example/incarnation"] { + save_config(store.clone(), path, b"retain".to_vec()) + .await + .expect("unrelated state should persist"); + } + let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load"); + let cancelled = CancellationToken::new(); + if completed == 0 { + store + .cancel_after_successful_puts + .lock() + .await + .insert(memory_config_key(RUSTFS_META_BUCKET, primary_path), (2, cancelled.clone())); + } else { + store + .cancel_after_deletes + .lock() + .await + .insert(memory_config_key(RUSTFS_META_BUCKET, &cleanup_paths[completed - 1]), cancelled.clone()); + } + let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || !cancelled.is_cancelled()) + .await + .expect_err("interruption should stop cleanup"); + assert!(err.to_string().contains("ownership"), "boundary {completed}: {err}"); + for (index, path) in cleanup_paths.iter().enumerate() { + assert_eq!( + store + .objects + .lock() + .await + .contains_key(&memory_config_key(RUSTFS_META_BUCKET, path)), + index >= completed, + "boundary {completed}, slot {index}" + ); + } + let intent = read_config_with_revision(store.clone(), primary_path) + .await + .expect("intent should persist"); + let slots = read_usage_state_reset_slots(store.clone()) + .await + .expect("restart should reload slots"); + reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true) + .await + .expect("restart should complete the same intent"); + assert_eq!( + read_config_with_revision(store.clone(), primary_path) + .await + .expect("intent should remain"), + intent + ); + assert_eq!(store.put_counts.lock().await[&memory_config_key(RUSTFS_META_BUCKET, primary_path)], 2); + for path in cleanup_paths { + assert!( + !store + .objects + .lock() + .await + .contains_key(&memory_config_key(RUSTFS_META_BUCKET, &path)) + ); + } + for path in ["buckets/quota-reservations/ledger", "buckets/example/incarnation"] { + assert_eq!(read_config(store.clone(), path).await.expect("unrelated state should remain"), b"retain"); + } + } +} + +#[tokio::test] +async fn scanner_usage_state_reset_stops_usage_fence_after_owner_loss() { + let store = Arc::new(MemoryConfigStore::default()); + let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + usage.scanner_epoch = Some(1); + let bytes = serde_json::to_vec(&usage).expect("baseline should encode"); + save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes.clone()) + .await + .expect("baseline should persist"); + let checks = AtomicUsize::new(0); + let err = fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 3, Some(0), false, || { + checks.fetch_add(1, Ordering::SeqCst) == 0 + }) + .await + .expect_err("ownership lost during reads must prevent the write"); + assert!(err.to_string().contains("leadership was lost"), "{err}"); + assert_eq!( + read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("baseline should remain"), + bytes + ); +} + +#[tokio::test] +async fn scanner_usage_state_reset_cancels_during_publication_admission() { + for resuming in [false, true] { + let store = Arc::new(MemoryConfigStore::default()); + let usage = if resuming { + scanner_usage_bootstrap_marker(std::time::SystemTime::UNIX_EPOCH, Some(3)) + } else { + complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0) + }; + save_config( + store.clone(), + DATA_USAGE_OBJ_NAME_PATH.as_str(), + serde_json::to_vec(&usage).expect("primary should encode"), + ) + .await + .expect("primary should persist"); + save_config(store.clone(), LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str(), b"corrupt".to_vec()) + .await + .expect("cleanup target should persist"); + let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load"); + let before = store.objects.lock().await.clone(); + let revisions_before = store.revisions.lock().await.clone(); + let entered = Arc::new(tokio::sync::Notify::new()); + let resume = Arc::new(tokio::sync::Notify::new()); + *store.pause_next_publication_admission.lock().await = Some((entered.clone(), resume.clone())); + let cancelled = CancellationToken::new(); + let (result, ()) = tokio::join!( + reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || !cancelled.is_cancelled()), + async { + entered.notified().await; + cancelled.cancel(); + resume.notify_one(); + } + ); + let err = result.expect_err("losing ownership during admission must prevent mutation"); + assert!(err.to_string().contains("ownership was lost"), "resuming={resuming}: {err}"); + assert_eq!(*store.objects.lock().await, before); + assert_eq!(*store.revisions.lock().await, revisions_before); + } +} + +#[tokio::test] +#[serial] +async fn scanner_usage_state_reset_rejects_corruption_without_a_trusted_floor() { + let (_temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await; + save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), b"{corrupt".to_vec()) + .await + .expect("corrupt primary should persist"); + let before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("evidence should load"); + let err = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone()) + .await + .expect_err("corruption must not become a zero floor"); + assert!(err.to_string().contains("no trusted cycle or usage floor"), "{err}"); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("evidence should remain"), + before + ); + assert!(matches!( + read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); + let mut backup = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + backup.scanner_epoch = Some(7); + backup.scanner_cycle = Some(40); + save_config( + store.clone(), + &format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()), + serde_json::to_vec(&backup).expect("backup should encode"), + ) + .await + .expect("valid backup should persist"); + let result = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store) + .await + .expect("valid backup should supply the recovery floor"); + assert_eq!(result.leader_epoch, 8); + assert_eq!(result.next_cycle, 41); +} + +#[tokio::test] +async fn scanner_usage_state_reset_rejects_replaced_intent_and_newer_cleanup_slot() { + let store = Arc::new(MemoryConfigStore::default()); + let marker = scanner_usage_bootstrap_marker(std::time::SystemTime::UNIX_EPOCH, Some(3)); + let bytes = serde_json::to_vec(&marker).expect("marker should encode"); + save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes.clone()) + .await + .expect("intent should persist"); + let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load"); + save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes) + .await + .expect("another intent should persist"); + let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true) + .await + .expect_err("same epoch cannot replace an intent revision"); + assert!(err.to_string().contains("intent revision changed"), "{err}"); + + let mut newer = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + newer.scanner_epoch = Some(3); + let path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()); + let bytes = serde_json::to_vec(&newer).expect("newer snapshot should encode"); + save_config(store.clone(), &path, bytes.clone()) + .await + .expect("newer snapshot should persist"); + let slots = read_usage_state_reset_slots(store.clone()) + .await + .expect("slots should reload"); + let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true) + .await + .expect_err("cleanup cannot delete same-epoch progress"); + assert!(err.to_string().contains("not older than its intent"), "{err}"); + assert_eq!(read_config(store, &path).await.expect("newer snapshot should remain"), bytes); +} + +#[tokio::test] +#[serial] +async fn scanner_usage_state_reset_rejects_decodable_untrusted_floor() { + let (_temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await; + let invalid_identity = DataUsageInfo { + usage_snapshot_complete: true, + buckets_count: 1, + last_update: Some(std::time::SystemTime::UNIX_EPOCH), + ..Default::default() + }; + for usage in [DataUsageInfo::default(), invalid_identity] { + save_config( + store.clone(), + DATA_USAGE_OBJ_NAME_PATH.as_str(), + serde_json::to_vec(&usage).expect("fixture should encode"), + ) + .await + .expect("untrusted primary should persist"); + let before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("primary should load"); + let err = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone()) + .await + .expect_err("valid JSON alone cannot prove a usage floor"); + assert!(err.to_string().contains("no trusted cycle or usage floor"), "{err}"); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("evidence should remain"), + before + ); + assert!(matches!( + read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); + } +} + +#[test] +fn full_rescan_reset_rejects_unknown_marker_phase_even_with_invalid_compat_fields() { + for state in [serde_json::json!("rewrite-v2"), serde_json::json!(7), serde_json::Value::Null] { + let marker = serde_json::json!({"state": state, "retry_count": "future-type", "schema_version": 99}); + let err = super::cycle_state::decode_recovery_marker_for_reset( + &serde_json::to_vec(&marker).expect("future marker should encode"), + &DataUsageCacheRevision::Etag("intent-1".to_string()), + ) + .expect_err("unknown persistent phases must remain fenced"); + assert!(err.to_string().contains("state is unsupported"), "{err}"); + } +} + +#[tokio::test] +#[serial] +async fn full_rescan_reset_preserves_unknown_phase_and_retries_completed_cleanup() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), b"corrupt".to_vec()) + .await + .expect("corrupt primary should persist"); + save_config( + store.clone(), + DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), + br#"{"state":"future-rewrite"}"#.to_vec(), + ) + .await + .expect("future marker should persist"); + let primary_before = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("primary should load"); + let marker_before = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .expect("marker should load"); + let err = reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect_err("unknown phase must block explicit reset"); + assert!(err.to_string().contains("state is unsupported"), "{err}"); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("primary should remain"), + primary_before + ); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()) + .await + .expect("marker should remain"), + marker_before + ); + + save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{malformed".to_vec()) + .await + .expect("recoverable marker should persist"); + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("reset should complete"); + let primary = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt primary should load"); + let usage = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("fenced usage should load"); + reset_scanner_cycle_recovery(CancellationToken::new(), store.clone()) + .await + .expect("retry after marker deletion should complete"); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()) + .await + .expect("rebuilt primary should remain"), + primary + ); + assert_eq!( + read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("fenced usage should remain"), + usage + ); +} + #[tokio::test] async fn scanner_usage_state_reset_slots_defer_when_publication_epoch_moves() { let store = Arc::new(MemoryConfigStore::default()); @@ -4993,7 +5367,7 @@ async fn scanner_usage_state_reset_slots_defer_when_publication_epoch_moves() { .expect("usage reset slots should be inspected"); store.publication_admission_blocked.store(true, Ordering::Release); - let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3) + let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3, || true) .await .expect_err("movement admission loss must defer reset"); assert!( From 0d1b312673d062e7ca236d2f444728a6afcfed5e Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 5 Sep 2026 16:49:58 +0800 Subject: [PATCH 5/5] fix(ci): require fresh successful scheduled validations (#7192) --- .../scheduled-validation-freshness.yml | 1 + .../check_scheduled_validation_freshness.py | 398 +++++++++++++----- 2 files changed, 293 insertions(+), 106 deletions(-) diff --git a/.github/workflows/scheduled-validation-freshness.yml b/.github/workflows/scheduled-validation-freshness.yml index eef340869..69ddc844c 100644 --- a/.github/workflows/scheduled-validation-freshness.yml +++ b/.github/workflows/scheduled-validation-freshness.yml @@ -42,6 +42,7 @@ jobs: - name: Check latest scheduled runs env: GH_TOKEN: ${{ secrets.GITHUB_TOKEN }} + RUSTFS_DEFAULT_BRANCH: ${{ github.event.repository.default_branch }} run: | set +e python3 scripts/check_scheduled_validation_freshness.py \ diff --git a/scripts/check_scheduled_validation_freshness.py b/scripts/check_scheduled_validation_freshness.py index a5812c7a3..c694e3c89 100644 --- a/scripts/check_scheduled_validation_freshness.py +++ b/scripts/check_scheduled_validation_freshness.py @@ -1,10 +1,11 @@ #!/usr/bin/env python3 -"""Fail when a critical scheduled validation has not started recently.""" +"""Require recent scheduled attempts and completed successes on the default branch.""" from __future__ import annotations import argparse from datetime import datetime, timedelta, timezone +import io import json import os from pathlib import Path @@ -13,7 +14,7 @@ import sys import tempfile import unittest from unittest import mock -from urllib.parse import quote, urlencode +from urllib.parse import parse_qs, quote, urlencode, urlsplit from urllib.request import Request, urlopen @@ -75,15 +76,8 @@ def stale_reason( run: dict[str, object] | None, now: datetime, max_age_hours: int, - never_ran_grace_until: datetime | None = None, ) -> str | None: if run is None: - # The grace deadline only covers a workflow whose first scheduled slot - # has not arrived yet (for example a monthly cron enabled mid-month). - # A recorded-but-old run proves the schedule used to fire and stopped, - # so the grace never masks that case. - if never_ran_grace_until is not None and now <= never_ran_grace_until: - return None return "no scheduled run has been recorded" created_at = parse_timestamp(run.get("created_at")) age = now - created_at @@ -93,14 +87,23 @@ def stale_reason( def fetch_latest_scheduled_run( - repository: str, workflow: str, token: str, api_url: str + repository: str, + workflow: str, + token: str, + api_url: str, + default_branch: str, + successful: bool = False, ) -> dict[str, object] | None: owner, repo = repository.split("/", 1) workflow_name = Path(workflow).name + query = {"event": "schedule", "branch": default_branch, "per_page": 1} + if successful: + # Filter on the server: the last success may be beyond a page of failures. + query["status"] = "success" endpoint = ( f"{api_url.rstrip('/')}/repos/{quote(owner, safe='')}/{quote(repo, safe='')}" f"/actions/workflows/{quote(workflow_name, safe='')}/runs?" - + urlencode({"event": "schedule", "per_page": 1}) + + urlencode(query) ) request = Request( endpoint, @@ -110,57 +113,104 @@ def fetch_latest_scheduled_run( "X-GitHub-Api-Version": "2022-11-28", }, ) - with urlopen(request, timeout=30) as response: + # Two requests per manifest entry must fit the watchdog's ten-minute job. + with urlopen(request, timeout=15) as response: payload = json.load(response) - runs = payload.get("workflow_runs") + runs = payload.get("workflow_runs") if isinstance(payload, dict) else None if not isinstance(runs, list): raise ValueError(f"GitHub returned no workflow_runs list for {workflow}") + total_count = payload.get("total_count") + if not isinstance(total_count, int) or isinstance(total_count, bool) or total_count < len(runs): + raise ValueError(f"GitHub returned an invalid run count for {workflow}") if not runs: + if total_count: + raise ValueError(f"GitHub returned an empty first page with recorded runs for {workflow}") return None - if not isinstance(runs[0], dict): + run = runs[0] + if not isinstance(run, dict): raise ValueError(f"GitHub returned an invalid workflow run for {workflow}") - return runs[0] + if run.get("event") != "schedule" or run.get("head_branch") != default_branch: + raise ValueError(f"GitHub returned a run outside the scheduled default-branch query for {workflow}") + if not isinstance(run.get("status"), str) or not run["status"]: + raise ValueError(f"GitHub returned no run status for {workflow}") + conclusion = run.get("conclusion") + if (conclusion is not None and not isinstance(conclusion, str)) or ( + run["status"] == "completed" and not conclusion + ): + raise ValueError(f"GitHub returned an invalid run conclusion for {workflow}") + if successful and (run["status"] != "completed" or conclusion != "success"): + raise ValueError(f"GitHub returned a run without a completed success for {workflow}") + parse_timestamp(run.get("created_at")) + if not isinstance(run.get("html_url"), str) or not run["html_url"]: + raise ValueError(f"GitHub returned no run URL for {workflow}") + return run -def write_report(path: Path, failures: list[tuple[str, int, str, str]]) -> None: - lines = ["## Scheduled validation freshness"] - if not failures: - lines.append("") - lines.append("All critical scheduled validations have a recent scheduled run.") - else: - lines.extend( - [ - "", - "The following critical validations are stale or could not be inspected:", - "", - "| Workflow | Limit | Result | Last run |", - "| --- | ---: | --- | --- |", - ] - ) - for workflow, max_age_hours, reason, run_url in failures: - link = f"[open]({run_url})" if run_url else "—" - lines.append(f"| `{workflow}` | {max_age_hours}h | {reason} | {link} |") +def describe_run(run: dict[str, object] | None) -> str: + if run is None: + return "No recorded run" + outcome = run["status"] + if run.get("conclusion"): + outcome = f"{outcome}/{run['conclusion']}" + return f"[{outcome}]({run['html_url']}) — created {run['created_at']}" + + +def write_report(path: Path, rows: list[tuple[str, int, str, str, str]], default_branch: str) -> None: + lines = [ + "## Scheduled validation freshness", + "", + f"Default branch: `{default_branch}`. Ages use scheduled-run creation time; rerunning an old commit does not refresh its evidence.", + "Attempt outcomes are shown independently of successful-run freshness.", + "Success is the GitHub workflow run conclusion; suite completeness remains the responsibility of each workflow.", + "", + "| Workflow | Limit | Freshness | Last attempt | Last completed success |", + "| --- | ---: | --- | --- | --- |", + ] + for workflow, max_age_hours, result, attempt, success in rows: + cells = [f"`{workflow}`", f"{max_age_hours}h", result, attempt, success] + lines.append("| " + " | ".join(cell.replace("|", "\\|").replace("\n", " ") for cell in cells) + " |") path.write_text("\n".join(lines) + "\n") def check_freshness( - config: Path, report: Path, repository: str, token: str, api_url: str + config: Path, report: Path, repository: str, token: str, api_url: str, default_branch: str ) -> int: now = datetime.now(timezone.utc) - failures: list[tuple[str, int, str, str]] = [] + rows: list[tuple[str, int, str, str, str]] = [] + failed = False for workflow, max_age_hours, never_ran_grace_until in load_validations(config): - try: - run = fetch_latest_scheduled_run(repository, workflow, token, api_url) - reason = stale_reason(run, now, max_age_hours, never_ran_grace_until) - if reason is not None: - run_url = str(run.get("html_url", "")) if run else "" - failures.append((workflow, max_age_hours, reason, run_url)) - except Exception as error: - failures.append( - (workflow, max_age_hours, f"inspection failed: {error}", "") - ) - write_report(report, failures) - return 1 if failures else 0 + runs: dict[str, dict[str, object] | None] = {} + reasons: list[str] = [] + for label, successful in (("Last attempt", False), ("Last completed success", True)): + try: + runs[label] = fetch_latest_scheduled_run( + repository, workflow, token, api_url, default_branch, successful + ) + except Exception as error: + reasons.append(f"{label}: inspection failed: {error}") + # A failed inspection or any recorded attempt ends first-run grace. + initial_grace = ( + len(runs) == 2 + and all(run is None for run in runs.values()) + and never_ran_grace_until is not None + and now <= never_ran_grace_until + ) + if not initial_grace: + for label, run in runs.items(): + reason = stale_reason(run, now, max_age_hours) + if reason is not None: + reasons.append(f"{label}: {reason}") + failed |= bool(reasons) + result = "; ".join(reasons) if reasons else "Fresh" + if initial_grace: + result = f"Initial grace until {never_ran_grace_until.isoformat()}" + evidence = [ + describe_run(runs[label]) if label in runs else "Inspection failed" + for label in ("Last attempt", "Last completed success") + ] + rows.append((workflow, max_age_hours, result, *evidence)) + write_report(report, rows, default_branch) + return 1 if failed else 0 class SelfTests(unittest.TestCase): @@ -173,15 +223,6 @@ class SelfTests(unittest.TestCase): self.assertIsNotNone(stale_reason(past_limit, self.NOW, 36)) self.assertIsNotNone(stale_reason(None, self.NOW, 36)) - def test_never_ran_grace_only_covers_missing_runs(self) -> None: - future_grace = self.NOW + timedelta(hours=1) - past_grace = self.NOW - timedelta(seconds=1) - self.assertIsNone(stale_reason(None, self.NOW, 36, future_grace)) - self.assertIsNone(stale_reason(None, self.NOW, 36, self.NOW)) - self.assertIsNotNone(stale_reason(None, self.NOW, 36, past_grace)) - stale_run = {"created_at": "2026-08-20T23:59:59Z"} - self.assertIsNotNone(stale_reason(stale_run, self.NOW, 36, future_grace)) - def test_config_rejects_duplicate_and_invalid_entries(self) -> None: with tempfile.TemporaryDirectory() as tmp: path = Path(tmp) / "validations.json" @@ -238,63 +279,205 @@ class SelfTests(unittest.TestCase): ], ) - def test_check_reports_missing_runs(self) -> None: + @staticmethod + def run_fixture(**overrides: object) -> dict[str, object]: + return { + "status": "completed", + "conclusion": "success", + "event": "schedule", + "head_branch": "release/current", + "created_at": "2026-08-22T00:00:00Z", + "html_url": "https://github.test/rustfs/rustfs/actions/runs/1", + **overrides, + } + + def check_payloads( + self, payloads: list[object], *, grace: str | None = None, workflows: int = 1 + ) -> tuple[int, str, list]: with tempfile.TemporaryDirectory() as tmp: root = Path(tmp) config = root / "validations.json" report = root / "report.md" - config.write_text( - json.dumps( - [ - {"workflow": ".github/workflows/ci.yml", "max_age_hours": 36}, - {"workflow": ".github/workflows/fuzz.yml", "max_age_hours": 36}, - {"workflow": ".github/workflows/mint.yml", "max_age_hours": 36}, - ] - ) - ) - with mock.patch( - __name__ + ".fetch_latest_scheduled_run", - side_effect=[ - {"created_at": "2999-01-01T00:00:00Z"}, - None, - RuntimeError("API unavailable"), - ], + entries = [ + {"workflow": f".github/workflows/check-{index}.yml", "max_age_hours": 36} + for index in range(workflows) + ] + if grace is not None: + entries[0]["never_ran_grace_until"] = grace + config.write_text(json.dumps(entries)) + responses = [] + for payload in payloads: + if isinstance(payload, dict) and isinstance(payload.get("workflow_runs"), list): + payload = {"total_count": len(payload["workflow_runs"]), **payload} + responses.append(payload if isinstance(payload, Exception) else io.StringIO(json.dumps(payload))) + with ( + mock.patch(__name__ + ".urlopen", side_effect=responses) as request, + mock.patch(__name__ + ".datetime", wraps=datetime) as clock, ): - self.assertEqual( - check_freshness( - config, - report, - "rustfs/rustfs", - "token", - "https://api.github.test", - ), - 1, + clock.now.return_value = self.NOW + status = check_freshness( + config, report, "rustfs/rustfs", "test-token", + "https://api.github.test", "release/current", ) - contents = report.read_text() - self.assertIn(".github/workflows/fuzz.yml", contents) - self.assertIn("inspection failed: API unavailable", contents) - self.assertNotIn(".github/workflows/ci.yml`", contents) + return status, report.read_text(), request.call_args_list - config.write_text( - json.dumps( - [{"workflow": ".github/workflows/ci.yml", "max_age_hours": 36}] + def test_requests_filter_schedule_default_branch_and_success_on_server(self) -> None: + attempt = self.run_fixture(status="in_progress", conclusion=None) + success = self.run_fixture(html_url="https://github.test/rustfs/rustfs/actions/runs/2") + status, report, calls = self.check_payloads([ + {"workflow_runs": [attempt], "total_count": 1001}, + {"workflow_runs": [success], "total_count": 1}, + ]) + self.assertEqual(status, 0) + self.assertEqual(len(calls), 2) + for call, successful in zip(calls, (False, True)): + request = call.args[0] + url = urlsplit(request.full_url) + self.assertEqual(url.path, "/repos/rustfs/rustfs/actions/workflows/check-0.yml/runs") + expected = {"event": ["schedule"], "branch": ["release/current"], "per_page": ["1"]} + if successful: + expected["status"] = ["success"] + self.assertEqual(parse_qs(url.query), expected) + self.assertEqual(request.get_header("Authorization"), "Bearer test-token") + self.assertEqual(call.kwargs, {"timeout": 15}) + self.assertIn("[in_progress]", report) + self.assertIn(str(attempt["html_url"]), report) + self.assertIn(str(success["html_url"]), report) + + def test_cancelled_attempt_cannot_refresh_expired_success(self) -> None: + attempt = self.run_fixture(conclusion="cancelled") + success = self.run_fixture( + created_at="2026-08-20T23:59:59Z", updated_at="2026-08-22T11:59:59Z", + html_url="https://github.test/rustfs/rustfs/actions/runs/2", + ) + status, report, _ = self.check_payloads([ + {"workflow_runs": [attempt]}, {"workflow_runs": [success]}, + ]) + self.assertEqual(status, 1) + self.assertIn("Last completed success: last scheduled run is", report) + self.assertIn("[completed/cancelled]", report) + for run in (attempt, success): + self.assertIn(str(run["html_url"]), report) + self.assertIn(str(run["created_at"]), report) + + def test_attempt_outcome_does_not_replace_recent_success(self) -> None: + success = self.run_fixture(created_at="2026-08-21T00:00:00Z") + for state, conclusion in ( + ("completed", "failure"), ("completed", "cancelled"), + ("completed", "timed_out"), ("completed", "success"), + ("queued", None), ("in_progress", None), + ): + with self.subTest(state=state, conclusion=conclusion): + status, report, _ = self.check_payloads([ + {"workflow_runs": [self.run_fixture(status=state, conclusion=conclusion)]}, + {"workflow_runs": [success]}, + ]) + self.assertEqual(status, 0) + self.assertIn(f"[{state}" + (f"/{conclusion}" if conclusion else "") + "]", report) + self.assertIn("Fresh", report) + self.assertNotIn("All critical scheduled validations", report) + + def test_grace_requires_two_successful_queries_with_no_history(self) -> None: + for attempt, success, grace, expected in ( + (None, None, "2026-08-22T12:00:00Z", 0), + (None, None, "2026-08-22T11:59:59Z", 1), + (self.run_fixture(conclusion="failure"), None, "2026-08-23T00:00:00Z", 1), + (self.run_fixture(status="queued", conclusion=None), None, "2026-08-23T00:00:00Z", 1), + (None, self.run_fixture(), "2026-08-23T00:00:00Z", 1), + ): + with self.subTest(attempt=attempt, success=success, grace=grace): + status, report, _ = self.check_payloads([ + {"workflow_runs": [] if attempt is None else [attempt]}, + {"workflow_runs": [] if success is None else [success]}, + ], grace=grace) + self.assertEqual(status, expected) + self.assertEqual("Initial grace until" in report, expected == 0) + + def test_api_failures_preserve_other_evidence_and_never_enter_grace(self) -> None: + good = {"workflow_runs": [self.run_fixture()]} + for first, second in ( + (RuntimeError("API unavailable"), good), + (good, RuntimeError("API unavailable")), + (RuntimeError("API unavailable"), {"workflow_runs": []}), + ): + with self.subTest(first=first, second=second): + status, report, calls = self.check_payloads( + [first, second], grace="2026-08-23T00:00:00Z" ) - ) - with mock.patch( - __name__ + ".fetch_latest_scheduled_run", - return_value={"created_at": "2999-01-01T00:00:00Z"}, - ): - self.assertEqual( - check_freshness( - config, - report, - "rustfs/rustfs", - "token", - "https://api.github.test", - ), - 0, - ) - self.assertIn("All critical scheduled validations", report.read_text()) + self.assertEqual(status, 1) + self.assertEqual(len(calls), 2) + self.assertIn("inspection failed: API unavailable", report) + self.assertNotIn("Initial grace until", report) + if first is good or second is good: + self.assertIn(str(self.run_fixture()["html_url"]), report) + + def test_invalid_api_evidence_fails_closed(self) -> None: + malformed = [ + [], {}, {"workflow_runs": {}}, {"workflow_runs": [None]}, + {"workflow_runs": [], "total_count": 1}, + {"workflow_runs": [], "total_count": -1}, + {"workflow_runs": [], "total_count": None}, + {"workflow_runs": [], "total_count": True}, + *({"workflow_runs": [self.run_fixture(**override)]} for override in ( + {"event": "workflow_dispatch"}, {"head_branch": "other"}, + {"created_at": "invalid"}, {"created_at": "2026-08-22T00:00:00"}, + {"status": None}, {"conclusion": None}, {"conclusion": 1}, + {"html_url": ""}, + )), + ] + for payload in malformed: + for index, label in enumerate(("Last attempt", "Last completed success")): + with self.subTest(payload=payload, label=label): + payloads = [{"workflow_runs": [self.run_fixture()]} for _ in range(2)] + payloads[index] = payload + status, report, _ = self.check_payloads(payloads, grace="2026-08-23T00:00:00Z") + self.assertEqual(status, 1) + self.assertIn(f"{label}: inspection failed", report) + self.assertNotIn("Initial grace until", report) + self.assertIn(str(self.run_fixture()["html_url"]), report) + for state, conclusion in (("in_progress", "success"), ("completed", "failure"), ("completed", "skipped")): + with self.subTest(state=state, conclusion=conclusion): + status, report, _ = self.check_payloads([ + {"workflow_runs": [self.run_fixture()]}, + {"workflow_runs": [self.run_fixture(status=state, conclusion=conclusion)]}, + ]) + self.assertEqual(status, 1) + self.assertIn("without a completed success", report) + + def test_report_retains_every_workflow(self) -> None: + status, report, calls = self.check_payloads([ + {"workflow_runs": [self.run_fixture()]}, {"workflow_runs": [self.run_fixture()]}, + {"workflow_runs": []}, {"workflow_runs": []}, + RuntimeError("API unavailable"), {"workflow_runs": [self.run_fixture()]}, + ], workflows=3) + self.assertEqual(status, 1) + self.assertEqual(len(calls), 6) + for index in range(3): + self.assertEqual(report.count(f"`.github/workflows/check-{index}.yml`"), 1) + self.assertIn("No recorded run", report) + self.assertIn("Inspection failed", report) + + def test_cli_requires_the_repository_default_branch(self) -> None: + from check_test_wiring import yaml_block + + workflow = (ROOT / ".github/workflows/scheduled-validation-freshness.yml").read_text().splitlines() + job = yaml_block(workflow, "check-freshness", 2) + self.assertIsNotNone(job) + start = job.index(" - name: Check latest scheduled runs") + end = next((index for index in range(start + 1, len(job)) if job[index].startswith(" - ")), len(job)) + environment = yaml_block(job[start:end], "env", 8) + self.assertIsNotNone(environment) + self.assertIn(" RUSTFS_DEFAULT_BRANCH: ${{ github.event.repository.default_branch }}", environment) + + with ( + mock.patch.dict(os.environ, {"GITHUB_REPOSITORY": "rustfs/rustfs", "GH_TOKEN": "test-token"}, clear=True), + mock.patch.object(sys, "argv", ["checker", "--report", "unused.md"]), + mock.patch("sys.stderr", new=io.StringIO()) as stderr, + self.assertRaises(SystemExit) as error, + ): + main() + self.assertEqual(error.exception.code, 2) + self.assertIn("RUSTFS_DEFAULT_BRANCH", stderr.getvalue()) def main() -> int: @@ -318,11 +501,14 @@ def main() -> int: repository = os.environ.get("GITHUB_REPOSITORY", "") token = os.environ.get("GH_TOKEN", "") api_url = os.environ.get("GITHUB_API_URL", "https://api.github.com") + default_branch = os.environ.get("RUSTFS_DEFAULT_BRANCH", "") if not re.fullmatch(r"[^/\s]+/[^/\s]+", repository): parser.error("GITHUB_REPOSITORY must be owner/repository") if not token: parser.error("GH_TOKEN is required") - return check_freshness(args.config, args.report, repository, token, api_url) + if not default_branch or any(character.isspace() for character in default_branch): + parser.error("RUSTFS_DEFAULT_BRANCH is required and must name the repository default branch") + return check_freshness(args.config, args.report, repository, token, api_url, default_branch) if __name__ == "__main__":