From 2d6f9417ffea2ec72f137f5655341b0b65a2ff6b Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 7 Sep 2026 14:27:59 +0800 Subject: [PATCH] feat(scanner): persist raw enumeration cancellation cursor (#7358) Return a partial data-usage cache when raw filesystem enumeration is cancelled before object progress can be written. The partial cache now carries a validated raw enumeration cursor for V2 checkpoint scans while keeping the snapshot incomplete and clearing older frontier/checkpoint metadata. Co-authored-by: zhi22915 --- crates/scanner/src/scanner_folder.rs | 115 +++++++++++++++++- crates/scanner/src/scanner_folder/tests.rs | 88 ++++++++++++++ .../tests/enumeration_restart.rs | 39 ++++-- 3 files changed, 227 insertions(+), 15 deletions(-) diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 05e04b0ac..3518fa002 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -21,8 +21,8 @@ use std::time::{Duration, Instant, SystemTime}; use crate::ReplTargetSizeSummary; use crate::data_usage_define::{ DATA_USAGE_SCAN_CHECKPOINT_VERSION, DataUsageCache, DataUsageCacheInfo, DataUsageEntry, DataUsageHash, DataUsageHashMap, - DataUsageScanCheckpoint, DataUsageScanCheckpointReason, PendingScannerHeal, PendingScannerHealKind, ScannerSizeSummaryExt, - SizeReconciliationEntry, SizeSummary, hash_path, + DataUsageRawEnumerationCursor, DataUsageScanCheckpoint, DataUsageScanCheckpointReason, PendingScannerHeal, + PendingScannerHealKind, ScannerSizeSummaryExt, SizeReconciliationEntry, SizeSummary, hash_path, }; use crate::error::ScannerError; use crate::runtime_config::{ @@ -55,6 +55,7 @@ use rustfs_scanner_metrics::metrics::{ UpdateCurrentPathFn, current_path_updater, global_metrics, }; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; +use sha2::{Digest, Sha256}; use time::OffsetDateTime; use tokio::select; use tokio::sync::mpsc; @@ -733,6 +734,7 @@ pub struct FolderScanner { coverage_frontier: Option, resume_frontier: Option, coverage_gap: bool, + raw_enumeration_progress: Vec, pending_heal_sync_deferred: bool, pending_heal_batch_dirty: bool, #[cfg(test)] @@ -744,6 +746,50 @@ pub struct FolderScanner { list_path_raw_options_observer: Option>, } +struct RawEnumerationProgress { + parent: String, + last_entry: Option, + entries_seen: u64, + digest: Sha256, +} + +impl RawEnumerationProgress { + fn new(parent: &str) -> Self { + let mut digest = Sha256::new(); + update_raw_enumeration_digest(&mut digest, b"parent", parent.as_bytes()); + Self { + parent: parent.to_string(), + last_entry: None, + entries_seen: 0, + digest, + } + } + + fn record_entry(&mut self, entry: &str) { + update_raw_enumeration_digest(&mut self.digest, b"entry", entry.as_bytes()); + self.last_entry = Some(entry.to_string()); + self.entries_seen = self.entries_seen.saturating_add(1); + } + + fn into_cursor(self) -> Option { + if self.entries_seen == 0 { + return None; + } + Some(DataUsageRawEnumerationCursor::new( + self.parent, + self.last_entry, + self.entries_seen, + self.digest.finalize().into(), + )) + } +} + +fn update_raw_enumeration_digest(digest: &mut Sha256, label: &[u8], value: &[u8]) { + digest.update(label); + digest.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_le_bytes()); + digest.update(value); +} + fn size_reconciliation_entry_bytes(entry: &SizeReconciliationEntry) -> usize { entry.key.len() + entry.bucket.len() @@ -999,6 +1045,41 @@ impl FolderScanner { } } + fn record_raw_enumeration_entry(&mut self, parent: &str, entry: &str) { + if self.old_cache.info.scan_progress.is_none() { + return; + } + if let Some(position) = self + .raw_enumeration_progress + .iter() + .position(|progress| progress.parent == parent) + { + self.raw_enumeration_progress.truncate(position + 1); + } else { + self.raw_enumeration_progress.push(RawEnumerationProgress::new(parent)); + } + if let Some(progress) = self.raw_enumeration_progress.last_mut() { + progress.record_entry(entry); + } + } + + fn finish_raw_enumeration_parent(&mut self, parent: &str) { + self.raw_enumeration_progress.retain(|progress| { + progress.parent != parent + && !progress + .parent + .strip_prefix(parent) + .is_some_and(|suffix| suffix.starts_with(SLASH_SEPARATOR)) + }); + } + + fn take_raw_enumeration_cursor(&mut self) -> Option { + self.raw_enumeration_progress + .drain(..) + .next() + .and_then(RawEnumerationProgress::into_cursor) + } + fn carry_forward_old_children(&mut self, parent_hash: &DataUsageHash, entry: &mut DataUsageEntry) { if entry.compacted { // Compacted entries store child totals directly; child links would be flattened twice. @@ -1329,11 +1410,15 @@ impl FolderScanner { }; let mut pending_entry_progress = 0_u64; let mut last_entry_progress = Instant::now(); + let mut raw_enumeration_complete = false; loop { let entry = match dir_reader.next_entry().await { Ok(Some(entry)) => entry, - Ok(None) => break, + Ok(None) => { + raw_enumeration_complete = true; + break; + } Err(e) if e.kind() == ErrorKind::NotFound => { debug!( target: "rustfs::scanner::folder", @@ -1345,6 +1430,7 @@ impl FolderScanner { error = %e, "Scanner folder state updated" ); + raw_enumeration_complete = true; break; } Err(e) if e.kind() == ErrorKind::NotADirectory => { @@ -1358,6 +1444,7 @@ impl FolderScanner { error = %e, "Scanner folder state updated" ); + raw_enumeration_complete = true; break; } Err(e) => return Err(ScannerError::Io(e)), @@ -1376,6 +1463,7 @@ impl FolderScanner { if file_name.is_empty() || file_name == "." || file_name == ".." { continue; } + self.record_raw_enumeration_entry(&folder.name, &file_name); let is_storage_format_entry = file_name == STORAGE_FORMAT_FILE; let file_path = entry.path().to_string_lossy().to_string(); @@ -1686,6 +1774,9 @@ impl FolderScanner { } } self.budget.record_entries_visited(pending_entry_progress); + if raw_enumeration_complete { + self.finish_raw_enumeration_parent(&folder.name); + } let mut found_erasure_data_directory = false; if self.is_erasure_mode && !found_object_metadata { @@ -2533,6 +2624,7 @@ pub(crate) async fn scan_data_folder_scoped( coverage_gap: false, pending_heal_sync_deferred: false, pending_heal_batch_dirty: false, + raw_enumeration_progress: Vec::new(), #[cfg(test)] pending_heal_sync_count: 0, pending_size_reconciliation_keys: HashSet::new(), @@ -2593,6 +2685,7 @@ pub(crate) async fn scan_data_folder_scoped( let had_scan_checkpoint = cache.info.scan_checkpoint.is_some() || new_cache.info.scan_checkpoint.is_some(); new_cache.info.scan_resume_after = None; new_cache.info.scan_checkpoint = None; + new_cache.info.scan_raw_enumeration_cursor = None; new_cache.info.scan_coverage_receipt = None; if had_scan_checkpoint { global_metrics().record_scanner_checkpoint_cleared(); @@ -2610,6 +2703,9 @@ pub(crate) async fn scan_data_folder_scoped( let root_hash = hash_path(&cache.info.name); let root_has_progress = data_usage_root_has_progress(&root); let pending_heals_changed = scanner.pending_heals_changed; + let raw_enumeration_cursor = scanner.take_raw_enumeration_cursor(); + let carry_forward_cache = + (raw_enumeration_cursor.is_some() && !root_has_progress).then(|| scanner.old_cache.cache.clone()); if root_has_progress { scanner.carry_forward_old_children(&root_hash, &mut root); } @@ -2617,8 +2713,19 @@ pub(crate) async fn scan_data_folder_scoped( let new_cache = scanner.as_mut_new_cache(); if root_has_progress { new_cache.replace_hashed(&root_hash, &None, &root); + } else if let Some(cache) = carry_forward_cache { + new_cache.cache = cache; } - if partial_cache_is_useful(&root, pending_heals_changed) || !new_cache.info.size_reconciliation.is_empty() { + if raw_enumeration_cursor.is_some() { + new_cache.info.scan_raw_enumeration_cursor = raw_enumeration_cursor; + new_cache.info.scan_checkpoint = None; + new_cache.info.scan_resume_after = None; + new_cache.info.scan_coverage_receipt = None; + } + if partial_cache_is_useful(&root, pending_heals_changed) + || new_cache.info.scan_raw_enumeration_cursor.is_some() + || !new_cache.info.size_reconciliation.is_empty() + { if new_cache.root().is_some() { new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN); } diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index 962090a00..cc1add61f 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -353,6 +353,7 @@ async fn build_test_scanner() -> (FolderScanner, std::path::PathBuf) { coverage_frontier: None, resume_frontier: None, coverage_gap: false, + raw_enumeration_progress: Vec::new(), pending_heal_sync_deferred: false, pending_heal_batch_dirty: false, pending_heal_sync_count: 0, @@ -2637,6 +2638,93 @@ async fn test_scan_data_folder_returns_partial_cache_on_budget_cancel() { assert_eq!(budget.reason(), Some(crate::scanner_budget::ScannerCycleBudgetReason::Directories)); } +#[tokio::test] +#[serial] +async fn test_scan_data_folder_returns_raw_cursor_on_enumeration_cancel_without_root_progress() { + let (scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir.clone()), + }; + + let bucket_dir = temp_dir.join("bucket"); + tokio::fs::create_dir_all(&bucket_dir) + .await + .expect("failed to create bucket directory"); + for entry in ["entry-a", "entry-b", "entry-c"] { + tokio::fs::write(bucket_dir.join(entry), b"data") + .await + .expect("failed to create raw directory entry"); + } + + let plan = crate::data_usage_define::DataUsageScanPlanDigest([11; 32]); + let source = crate::data_usage_define::DataUsageCacheSource::new(1, 0); + let identity = crate::data_usage_define::DataUsageScanIdentity { + version: 1, + bucket_incarnation: Uuid::from_u128(7), + set_layout: crate::data_usage_define::DataUsageScanPlanDigest([12; 32]), + publication_epoch: 3, + tier_registry_generation: 0, + scan_mode: HealScanMode::Normal, + }; + let mut cache = DataUsageCache { + info: crate::data_usage_define::DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + ..Default::default() + }, + ..Default::default() + }; + assert_eq!( + cache.prepare_bucket_checkpoint("bucket", 7, 3, source, plan, identity), + crate::data_usage_define::DataUsageCachePrepareOutcome::Reset + ); + + let parent = CancellationToken::new(); + let budget = ScannerCycleBudget::new_with_progress_tracking(&parent, Default::default()); + let _raw_entry_budget = enumeration_restart::install_raw_entry_budget(scanner.local_disk.path(), 1); + + let result = scan_data_folder( + budget.token(), + budget.clone(), + vec![scanner.local_disk.clone()], + scanner.local_disk.clone(), + cache, + None, + HealScanMode::Normal, + SCANNER_SLEEPER.clone(), + ) + .await; + + let partial_cache = match result { + Err(ScannerError::PartialCache(partial_cache)) => partial_cache, + other => panic!("expected raw enumeration partial cache after cancellation, got {other:?}"), + }; + + assert!( + partial_cache + .root() + .is_none_or(|root| root.objects == 0 && root.versions == 0 && root.size == 0), + "raw cursor writer must not invent object progress" + ); + assert!(partial_cache.info.last_update.is_some()); + assert_eq!(partial_cache.info.next_cycle, 7); + assert!(!partial_cache.info.snapshot_complete); + assert!(partial_cache.info.scan_checkpoint.is_none()); + assert!(partial_cache.info.scan_resume_after.is_none()); + + let raw_cursor = partial_cache + .info + .scan_raw_enumeration_cursor + .as_ref() + .expect("raw enumeration cancellation should persist a cursor"); + assert_eq!(raw_cursor.parent, "bucket"); + assert_eq!(raw_cursor.entries_seen, 1); + assert!(raw_cursor.last_entry.is_some()); + assert_ne!(raw_cursor.page_digest, [0; 32]); + assert_eq!(partial_cache.validated_raw_enumeration_cursor(), Some(raw_cursor)); + assert_eq!(budget.reason(), Some(crate::scanner_budget::ScannerCycleBudgetReason::Runtime)); +} + #[tokio::test] #[serial] async fn test_scan_data_folder_reports_invalid_checkpoint_ignored_once() { diff --git a/crates/scanner/src/scanner_folder/tests/enumeration_restart.rs b/crates/scanner/src/scanner_folder/tests/enumeration_restart.rs index ccb37882b..f04759491 100644 --- a/crates/scanner/src/scanner_folder/tests/enumeration_restart.rs +++ b/crates/scanner/src/scanner_folder/tests/enumeration_restart.rs @@ -40,7 +40,7 @@ pub(in crate::scanner_folder) fn observe_raw_entry(dir: &str, name: &std::ffi::O } } -struct ObservationGuard; +pub(in crate::scanner_folder) struct ObservationGuard; impl Drop for ObservationGuard { fn drop(&mut self) { @@ -48,6 +48,18 @@ impl Drop for ObservationGuard { } } +pub(in crate::scanner_folder) fn install_raw_entry_budget(root: PathBuf, limit: u64) -> ObservationGuard { + *OBSERVATION.lock().expect("install raw-entry observation") = Some(Observation { + root, + limit, + entries: 0, + name_bytes: 0, + first_entry: None, + last_entry: None, + }); + ObservationGuard +} + #[derive(serde::Deserialize)] #[serde(deny_unknown_fields)] struct Request { @@ -84,8 +96,21 @@ async fn round(request: &Request) -> serde_json::Value { } let mut initial = DataUsageCache::default(); initial.info.name = "bucket".to_string(); + let source = crate::data_usage_define::DataUsageCacheSource::new(0, 0); + let plan = crate::data_usage_define::DataUsageScanPlanDigest([31; 32]); + let identity = crate::data_usage_define::DataUsageScanIdentity { + version: 1, + bucket_incarnation: Uuid::from_u128(31), + set_layout: crate::data_usage_define::DataUsageScanPlanDigest([32; 32]), + publication_epoch: 1, + tier_registry_generation: 0, + scan_mode: HealScanMode::Normal, + }; + assert_eq!( + initial.prepare_bucket_checkpoint("bucket", 1, 0, source, plan, identity), + crate::data_usage_define::DataUsageCachePrepareOutcome::Reset + ); initial.info.skip_healing = true; - initial.info.snapshot_complete = false; initial.replace("bucket", "", DataUsageEntry::default()); tokio::fs::write(&cache_path, initial.marshal_msg().expect("initial cache codec")) .await @@ -106,15 +131,7 @@ async fn round(request: &Request) -> serde_json::Value { .expect("open synthetic disk in this process"); let parent = CancellationToken::new(); let budget = ScannerCycleBudget::new_with_progress_tracking(&parent, Default::default()); - *OBSERVATION.lock().expect("install observation") = Some(Observation { - root: disk.path(), - limit: request.raw_entry_budget, - entries: 0, - name_bytes: 0, - first_entry: None, - last_entry: None, - }); - let _observation_guard = ObservationGuard; + let _observation_guard = install_raw_entry_budget(disk.path(), request.raw_entry_budget); let result = scan_data_folder( budget.token(), budget.clone(),