From 542720a1f73ac06a54b6bb2d5bc40400b669ad06 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Thu, 4 Jun 2026 17:21:30 +0800 Subject: [PATCH] feat(scanner): add partial scan resume hints (#3207) * feat(scanner): add partial scan resume hints * test(scanner): cover clearing scan resume hints * fix(scanner): apply resume hint to combined child order --------- Co-authored-by: Henry Guo --- crates/scanner/src/data_usage_define.rs | 57 +++ crates/scanner/src/scanner_folder.rs | 452 ++++++++++++++++++++---- 2 files changed, 435 insertions(+), 74 deletions(-) diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 11afdb844..263a15ab2 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -227,6 +227,8 @@ pub struct DataUsageCacheInfo { pub replication: Option>, #[serde(default)] pub failed_objects: HashMap, + #[serde(default)] + pub scan_resume_after: Option, } /// Data usage cache @@ -1022,6 +1024,61 @@ mod tests { assert_eq!(decoded.failed_objects, 0); } + #[test] + fn test_data_usage_cache_info_deserialize_defaults_scan_resume_after() { + let value = serde_json::json!({ + "name": "bucket", + "next_cycle": 7, + "last_update": null, + "skip_healing": false, + "lifecycle": null, + "replication": null, + "failed_objects": {} + }); + + let decoded: DataUsageCacheInfo = serde_json::from_value(value).expect("Failed to deserialize cache info"); + + assert_eq!(decoded.name, "bucket"); + assert_eq!(decoded.next_cycle, 7); + assert!(decoded.scan_resume_after.is_none()); + } + + #[test] + fn test_data_usage_cache_info_unmarshal_old_msgpack_defaults_scan_resume_after() { + #[derive(Serialize)] + struct OldDataUsageCacheInfo { + name: String, + next_cycle: u64, + last_update: Option, + skip_healing: bool, + lifecycle: Option>, + replication: Option>, + failed_objects: HashMap, + } + + let old_info = OldDataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + last_update: None, + skip_healing: true, + lifecycle: None, + replication: None, + failed_objects: HashMap::from([("bad-object".to_string(), 11)]), + }; + let mut buf = Vec::new(); + old_info + .serialize(&mut rmp_serde::Serializer::new(&mut buf)) + .expect("Failed to serialize old cache info"); + + let decoded: DataUsageCacheInfo = rmp_serde::from_slice(&buf).expect("Failed to deserialize old cache info"); + + assert_eq!(decoded.name, "bucket"); + assert_eq!(decoded.next_cycle, 7); + assert!(decoded.skip_healing); + assert_eq!(decoded.failed_objects.get("bad-object"), Some(&11)); + assert!(decoded.scan_resume_after.is_none()); + } + #[test] fn test_data_usage_cache_mutations_update_in_place() { let mut cache = DataUsageCache { diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index d57ac5b3d..81cb0234a 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -142,6 +142,70 @@ fn should_yield_after_object(object_count: u64, yield_every: u64) -> bool { yield_every > 0 && object_count.is_multiple_of(yield_every) } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum FolderResumeMatch { + Exact, + Descendant, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum FolderScanSource { + New, + Existing, +} + +#[derive(Clone, Debug)] +struct QueuedFolder { + folder: CachedFolder, + source: FolderScanSource, +} + +fn folder_resume_match(folder_name: &str, resume_after: &str) -> Option { + if resume_after == folder_name { + return Some(FolderResumeMatch::Exact); + } + resume_after + .strip_prefix(folder_name) + .filter(|suffix| suffix.starts_with(SLASH_SEPARATOR)) + .map(|_| FolderResumeMatch::Descendant) +} + +fn order_items_for_resume(items: &mut [T], resume_after: Option<&str>, name: F) +where + F: Fn(&T) -> &str, +{ + items.sort_by(|left, right| name(left).cmp(name(right))); + + let Some(resume_after) = resume_after.filter(|resume_after| !resume_after.is_empty()) else { + return; + }; + + let Some((resume_index, resume_match)) = items + .iter() + .enumerate() + .find_map(|(index, item)| folder_resume_match(name(item), resume_after).map(|resume_match| (index, resume_match))) + else { + return; + }; + + let rotate_by = match resume_match { + FolderResumeMatch::Exact => resume_index + 1, + FolderResumeMatch::Descendant => resume_index, + }; + if rotate_by < items.len() { + items.rotate_left(rotate_by); + } +} + +#[cfg(test)] +fn order_folders_for_resume(folders: &mut [CachedFolder], resume_after: Option<&str>) { + order_items_for_resume(folders, resume_after, |folder| folder.name.as_str()); +} + +fn order_queued_folders_for_resume(folders: &mut [QueuedFolder], resume_after: Option<&str>) { + order_items_for_resume(folders, resume_after, |folder| folder.folder.name.as_str()); +} + fn should_alert_excessive_versions(remaining_versions: usize, cumulative_size: i64) -> (bool, bool) { let too_many_versions = remaining_versions as u64 >= scanner_excess_versions_threshold(); let too_large_versions = cumulative_size > 0 && cumulative_size as u64 >= scanner_excess_version_size_threshold(); @@ -781,6 +845,11 @@ impl FolderScanner { } } + fn record_scan_resume_hint(&mut self, folder: &str) { + self.new_cache.info.scan_resume_after = Some(folder.to_string()); + self.update_cache.info.scan_resume_after = Some(folder.to_string()); + } + fn alert_excessive_folders(&self, folder: &str, total_folders: usize) { let threshold = scanner_excess_folders_threshold(); if total_folders as u64 <= threshold { @@ -1155,37 +1224,69 @@ impl FolderScanner { } } - // Scan new folders - for folder_item in new_folders { + let scan_resume_after = self.old_cache.info.scan_resume_after.as_deref(); + let mut queued_folders = Vec::with_capacity(new_folders.len() + existing_folders.len()); + queued_folders.extend(new_folders.into_iter().map(|folder| QueuedFolder { + folder, + source: FolderScanSource::New, + })); + queued_folders.extend(existing_folders.into_iter().map(|folder| QueuedFolder { + folder, + source: FolderScanSource::Existing, + })); + order_queued_folders_for_resume(&mut queued_folders, scan_resume_after); + + // Scan child folders in the combined resume order. + for queued_folder in queued_folders { if ctx.is_cancelled() { return Err(ScannerError::Other("Operation cancelled".to_string())); } + let mut folder_item = queued_folder.folder; let h = hash_path(&folder_item.name); - // Add new folders to the update tree so totals update for these. - if !into.compacted { - let mut found_any = false; - let mut parent = this_hash.clone(); - let update_cache_name_hash = hash_path(&self.update_cache.info.name); - while parent != update_cache_name_hash { - let parent_key = parent.key(); - let e = self.update_cache.find(&parent_key); - if e.is_none_or(|v| v.compacted) { - found_any = true; - break; - } - if let Some(next) = self.update_cache.search_parent(&parent) { - parent = next; - } else { - found_any = true; - break; + match queued_folder.source { + FolderScanSource::New => { + // Add new folders to the update tree so totals update for these. + if !into.compacted { + let mut found_any = false; + let mut parent = this_hash.clone(); + let update_cache_name_hash = hash_path(&self.update_cache.info.name); + + while parent != update_cache_name_hash { + let parent_key = parent.key(); + let e = self.update_cache.find(&parent_key); + if e.is_none_or(|v| v.compacted) { + found_any = true; + break; + } + if let Some(next) = self.update_cache.search_parent(&parent) { + parent = next; + } else { + found_any = true; + break; + } + } + if !found_any { + // Add non-compacted empty entry. + self.update_cache + .replace_hashed(&h, &Some(this_hash.clone()), &DataUsageEntry::default()); + } } } - if !found_any { - // Add non-compacted empty entry. - self.update_cache - .replace_hashed(&h, &Some(this_hash.clone()), &DataUsageEntry::default()); + FolderScanSource::Existing => { + if !into.compacted && self.old_cache.is_compacted(&h) { + let next_cycle = self.old_cache.info.next_cycle as u32; + if !h.mod_(next_cycle, data_usage_update_dir_cycles()) { + // Transfer and add as child... + self.new_cache.copy_with_children(&self.old_cache, &h, &folder_item.parent); + into.add_child(&h); + self.record_scan_resume_hint(&folder_item.name); + continue; + } + + folder_item.object_heal_prob_div = data_usage_update_dir_cycles(); + } } } @@ -1195,6 +1296,7 @@ impl FolderScanner { // In compacted mode child totals are accumulated directly into the parent entry. let fut = Box::pin(self.scan_folder(ctx.clone(), folder_item.clone(), into)); fut.await.map_err(|e| ScannerError::Other(e.to_string()))?; + self.record_scan_resume_hint(&folder_item.name); tokio::task::yield_now().await; } else { let mut dst = DataUsageEntry::default(); @@ -1212,69 +1314,23 @@ impl FolderScanner { let h = DataUsageHash(folder_item.name.clone()); into.add_child(&h); + self.record_scan_resume_hint(&folder_item.name); // We scanned a folder, optionally send update. self.update_cache.delete_recursive(&h); self.update_cache.copy_with_children(&self.new_cache, &h, &folder_item.parent); self.send_update().await; } - if !into.compacted && self.update_cache.find(&this_hash.key()).is_some_and(|v| !v.compacted) { + if queued_folder.source == FolderScanSource::New + && !into.compacted + && self.update_cache.find(&this_hash.key()).is_some_and(|v| !v.compacted) + { self.update_cache.delete_recursive(&h); self.update_cache .copy_with_children(&self.new_cache, &h, &Some(this_hash.clone())); } } - // Scan existing folders - for mut folder_item in existing_folders { - if ctx.is_cancelled() { - return Err(ScannerError::Other("Operation cancelled".to_string())); - } - - let h = hash_path(&folder_item.name); - - if !into.compacted && self.old_cache.is_compacted(&h) { - let next_cycle = self.old_cache.info.next_cycle as u32; - if !h.mod_(next_cycle, data_usage_update_dir_cycles()) { - // Transfer and add as child... - self.new_cache.copy_with_children(&self.old_cache, &h, &folder_item.parent); - into.add_child(&h); - continue; - } - - folder_item.object_heal_prob_div = data_usage_update_dir_cycles(); - } - - (self.update_current_path)(&folder_item.name).await; - - if into.compacted { - // In compacted mode child totals are accumulated directly into the parent entry. - let fut = Box::pin(self.scan_folder(ctx.clone(), folder_item.clone(), into)); - fut.await.map_err(|e| ScannerError::Other(e.to_string()))?; - tokio::task::yield_now().await; - } else { - let mut dst = DataUsageEntry::default(); - - // Use Box::pin for recursive async call - let fut = Box::pin(self.scan_folder(ctx.clone(), folder_item.clone(), &mut dst)); - if let Err(e) = fut.await { - if ctx.is_cancelled() { - return Err(e); - } - warn!("scan_folder: failed to scan child folder {}: {}", folder_item.name, e); - continue; - } - tokio::task::yield_now().await; - - let h = DataUsageHash(folder_item.name.clone()); - into.add_child(&h); - // We scanned a folder, optionally send update. - self.update_cache.delete_recursive(&h); - self.update_cache.copy_with_children(&self.new_cache, &h, &folder_item.parent); - self.send_update().await; - } - } - // Scan for healing if abandoned_children.is_empty() || !self.should_heal().await { debug!("scan_folder: done for now abandoned children are empty or we are not healing"); @@ -1674,6 +1730,7 @@ pub async fn scan_data_folder( new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN); new_cache.info.last_update = Some(SystemTime::now()); new_cache.info.next_cycle = cache.info.next_cycle; + new_cache.info.scan_resume_after = None; close_disk().await; Ok(new_cache.clone()) @@ -1901,6 +1958,72 @@ mod tests { assert!(!should_yield_after_object(128, 0)); } + #[test] + fn test_order_folders_for_resume_rotates_after_exact_resume_hint() { + let mut folders = vec![ + CachedFolder { + name: "bucket/child-c".to_string(), + parent: None, + object_heal_prob_div: 1, + }, + CachedFolder { + name: "bucket/child-a".to_string(), + parent: None, + object_heal_prob_div: 1, + }, + CachedFolder { + name: "bucket/child-b".to_string(), + parent: None, + object_heal_prob_div: 1, + }, + ]; + + order_folders_for_resume(&mut folders, Some("bucket/child-b")); + + let names = folders.into_iter().map(|folder| folder.name).collect::>(); + assert_eq!( + names, + vec![ + "bucket/child-c".to_string(), + "bucket/child-a".to_string(), + "bucket/child-b".to_string() + ] + ); + } + + #[test] + fn test_order_folders_for_resume_prioritizes_descendant_resume_hint() { + let mut folders = vec![ + CachedFolder { + name: "bucket/child-c".to_string(), + parent: None, + object_heal_prob_div: 1, + }, + CachedFolder { + name: "bucket/child-a".to_string(), + parent: None, + object_heal_prob_div: 1, + }, + CachedFolder { + name: "bucket/child-b".to_string(), + parent: None, + object_heal_prob_div: 1, + }, + ]; + + order_folders_for_resume(&mut folders, Some("bucket/child-b/grandchild")); + + let names = folders.into_iter().map(|folder| folder.name).collect::>(); + assert_eq!( + names, + vec![ + "bucket/child-b".to_string(), + "bucket/child-c".to_string(), + "bucket/child-a".to_string() + ] + ); + } + #[tokio::test] #[serial] async fn test_record_failed_prunes_to_max_entries() { @@ -2284,6 +2407,187 @@ mod tests { assert_eq!(budget.reason(), Some(crate::scanner_budget::ScannerCycleBudgetReason::Directories)); } + #[tokio::test] + #[serial] + async fn test_scan_data_folder_resume_hint_prioritizes_next_existing_folder() { + let (scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir.clone()), + }; + + let bucket_dir = temp_dir.join("bucket"); + for child in ["child-a", "child-b", "child-c"] { + tokio::fs::create_dir_all(bucket_dir.join(child)) + .await + .expect("failed to create child directory"); + } + + let root_hash = hash_path("bucket"); + let mut cache = DataUsageCache { + info: crate::data_usage_define::DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 9, + scan_resume_after: Some("bucket/child-a".to_string()), + ..Default::default() + }, + ..Default::default() + }; + cache.replace_hashed(&root_hash, &None, &DataUsageEntry::default()); + for child in ["child-a", "child-b", "child-c"] { + cache.replace_hashed( + &hash_path(&format!("bucket/{child}")), + &Some(root_hash.clone()), + &DataUsageEntry::default(), + ); + } + + let parent = CancellationToken::new(); + let budget = ScannerCycleBudget::new( + &parent, + crate::scanner_budget::ScannerCycleBudgetConfig { + max_directories: Some(2), + ..Default::default() + }, + ); + + 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 partial cache after directory budget cancellation, got {other:?}"), + }; + + assert_eq!(partial_cache.info.scan_resume_after.as_deref(), Some("bucket/child-b")); + assert!( + partial_cache + .root() + .is_some_and(|root| root.children.contains(&hash_path("bucket/child-b").key())) + ); + assert_eq!(partial_cache.info.next_cycle, 9); + assert_eq!(budget.reason(), Some(crate::scanner_budget::ScannerCycleBudgetReason::Directories)); + } + + #[tokio::test] + #[serial] + async fn test_scan_data_folder_resume_hint_orders_across_new_and_existing_folders() { + let (scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir.clone()), + }; + + let bucket_dir = temp_dir.join("bucket"); + for child in ["child-a", "child-b", "child-c", "child-d"] { + tokio::fs::create_dir_all(bucket_dir.join(child)) + .await + .expect("failed to create child directory"); + } + + let root_hash = hash_path("bucket"); + let mut cache = DataUsageCache { + info: crate::data_usage_define::DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 9, + scan_resume_after: Some("bucket/child-b".to_string()), + ..Default::default() + }, + ..Default::default() + }; + cache.replace_hashed(&root_hash, &None, &DataUsageEntry::default()); + for child in ["child-b", "child-c"] { + cache.replace_hashed( + &hash_path(&format!("bucket/{child}")), + &Some(root_hash.clone()), + &DataUsageEntry::default(), + ); + } + + let parent = CancellationToken::new(); + let budget = ScannerCycleBudget::new( + &parent, + crate::scanner_budget::ScannerCycleBudgetConfig { + max_directories: Some(2), + ..Default::default() + }, + ); + + 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 partial cache after directory budget cancellation, got {other:?}"), + }; + + assert_eq!(partial_cache.info.scan_resume_after.as_deref(), Some("bucket/child-c")); + assert!( + partial_cache + .root() + .is_some_and(|root| root.children.contains(&hash_path("bucket/child-c").key())) + ); + assert_eq!(budget.reason(), Some(crate::scanner_budget::ScannerCycleBudgetReason::Directories)); + } + + #[tokio::test] + #[serial] + async fn test_scan_data_folder_success_clears_resume_hint() { + let (scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir.clone()), + }; + + tokio::fs::create_dir_all(temp_dir.join("bucket").join("child-a")) + .await + .expect("failed to create child directory"); + + let cache = DataUsageCache { + info: crate::data_usage_define::DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 11, + scan_resume_after: Some("bucket/child-a".to_string()), + ..Default::default() + }, + ..Default::default() + }; + + let parent = CancellationToken::new(); + let budget = ScannerCycleBudget::new(&parent, Default::default()); + + let result = scan_data_folder( + budget.token(), + budget, + vec![scanner.local_disk.clone()], + scanner.local_disk.clone(), + cache, + None, + HealScanMode::Normal, + SCANNER_SLEEPER.clone(), + ) + .await + .expect("scan should complete successfully"); + + assert!(result.info.scan_resume_after.is_none()); + assert_eq!(result.info.next_cycle, 11); + } + #[tokio::test] #[serial] #[cfg(unix)]