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 <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-06-04 17:21:30 +08:00
committed by GitHub
parent a7be7c558d
commit 542720a1f7
2 changed files with 435 additions and 74 deletions
+57
View File
@@ -227,6 +227,8 @@ pub struct DataUsageCacheInfo {
pub replication: Option<Arc<ReplicationConfig>>,
#[serde(default)]
pub failed_objects: HashMap<String, u64>,
#[serde(default)]
pub scan_resume_after: Option<String>,
}
/// 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<SystemTime>,
skip_healing: bool,
lifecycle: Option<Arc<BucketLifecycleConfiguration>>,
replication: Option<Arc<ReplicationConfig>>,
failed_objects: HashMap<String, u64>,
}
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 {
+378 -74
View File
@@ -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<FolderResumeMatch> {
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<T, F>(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::<Vec<_>>();
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::<Vec<_>>();
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)]