From 2585855d23fe334f5f5fa04ae55da8967b7bec2c Mon Sep 17 00:00:00 2001 From: cxymds Date: Mon, 6 Jul 2026 14:04:53 +0800 Subject: [PATCH] fix(ecstore): hide deleted versioned folder prefixes (#4296) --- .../ecstore/src/cache_value/metacache_set.rs | 4 + crates/ecstore/src/disk/local.rs | 206 +++++++++++++++++- crates/ecstore/src/disk/mod.rs | 6 + crates/ecstore/src/store/list_objects.rs | 85 +++++++- 4 files changed, 299 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/cache_value/metacache_set.rs b/crates/ecstore/src/cache_value/metacache_set.rs index f7fe5c82f..56fb0e489 100644 --- a/crates/ecstore/src/cache_value/metacache_set.rs +++ b/crates/ecstore/src/cache_value/metacache_set.rs @@ -110,6 +110,7 @@ pub struct ListPathRawOptions { pub bucket: String, pub path: String, pub recursive: bool, + pub incl_deleted: bool, pub filter_prefix: Option, pub forward_to: Option, pub min_disks: usize, @@ -138,6 +139,7 @@ impl Clone for ListPathRawOptions { bucket: self.bucket.clone(), path: self.path.clone(), recursive: self.recursive, + incl_deleted: self.incl_deleted, filter_prefix: self.filter_prefix.clone(), forward_to: self.forward_to.clone(), min_disks: self.min_disks, @@ -251,6 +253,7 @@ async fn list_path_raw_inner( bucket: opts_clone.bucket.clone(), base_dir: opts_clone.path.clone(), recursive: opts_clone.recursive, + incl_deleted: opts_clone.incl_deleted, report_notfound: opts_clone.report_not_found, filter_prefix: opts_clone.filter_prefix.clone(), forward_to: opts_clone.forward_to.clone(), @@ -403,6 +406,7 @@ async fn list_path_raw_inner( bucket: opts_clone.bucket.clone(), base_dir: opts_clone.path.clone(), recursive: opts_clone.recursive, + incl_deleted: opts_clone.incl_deleted, report_notfound: opts_clone.report_not_found, filter_prefix: opts_clone.filter_prefix.clone(), forward_to: opts_clone.forward_to.clone(), diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 42b42e72f..e876f7a35 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -2110,7 +2110,12 @@ impl LocalDisk { // If dirObject, but no metadata (which is unexpected) we skip it. if !is_dir_obj && !is_empty_dir(self.get_object_path(&opts.bucket, &meta.name)?).await { meta.name.push_str(SLASH_SEPARATOR); - schedule_dir(&mut dir_stack, meta.name, false, None); + if opts.recursive + || opts.incl_deleted + || self.directory_has_visible_listing_entry(&opts.bucket, &meta.name).await? + { + schedule_dir(&mut dir_stack, meta.name, false, None); + } } continue; @@ -2161,6 +2166,76 @@ impl LocalDisk { Ok(()) } + + async fn directory_has_visible_listing_entry(&self, bucket: &str, dir_name: &str) -> Result { + let mut stack = vec![dir_name.trim_matches('/').to_owned()]; + + while let Some(current) = stack.pop() { + if current.is_empty() { + continue; + } + + let entries = match self.list_dir("", bucket, ¤t, -1).await { + Ok(entries) => entries, + Err(err) => { + if err == DiskError::VolumeNotFound || err == Error::FileNotFound { + continue; + } + + return Err(err); + } + }; + + let mut data_dirs_to_skip = HashSet::new(); + let mut child_dirs = Vec::new(); + + for entry in entries { + if entry == STORAGE_FORMAT_FILE { + let metadata_path = path_join_buf(&[current.as_str(), STORAGE_FORMAT_FILE]); + match self.read_metadata(bucket, metadata_path.as_str()).await { + Ok(metadata) => { + let file_meta = match FileMeta::load(&metadata) { + Ok(file_meta) => file_meta, + Err(_) => return Ok(true), + }; + + if file_meta_counts_toward_limit(&file_meta) { + return Ok(true); + } + + if let Ok(data_dirs) = file_meta.get_data_dirs() { + for data_dir in data_dirs.iter().flatten() { + data_dirs_to_skip.insert(data_dir.to_string()); + } + } + } + Err(err) => { + if err != Error::FileNotFound && err != Error::IsNotRegular { + return Err(err); + } + } + } + + continue; + } + + if entry.ends_with(SLASH_SEPARATOR) { + let child = entry.trim_end_matches(SLASH_SEPARATOR); + if !child.is_empty() { + child_dirs.push(child.to_owned()); + } + } + } + + for child in child_dirs { + if !data_dirs_to_skip.contains(&child) { + stack.push(path_join_buf(&[current.as_str(), child.as_str()])); + } + } + } + + Ok(false) + } } pub struct ScanGuard(pub Arc); @@ -5296,6 +5371,135 @@ mod test { assert_eq!(objs_returned, 1); } + #[tokio::test] + async fn test_scan_dir_nonrecursive_skips_dirs_with_only_hidden_delete_markers() { + use rustfs_filemeta::MetacacheReader; + use tempfile::tempdir; + + fn hidden_versioned_object_metadata(name: &str, delete_version_id: &str, object_version_id: &str) -> Vec { + let mut fm = FileMeta::default(); + fm.add_version({ + let mut fi = FileInfo::new(name, 1, 1); + fi.version_id = Some(Uuid::parse_str(object_version_id).expect("test version id should parse")); + fi.mod_time = Some(OffsetDateTime::now_utc() - time::Duration::seconds(1)); + fi + }) + .expect("object metadata should be valid"); + fm.add_version(FileInfo { + name: name.to_owned(), + deleted: true, + version_id: Some(Uuid::parse_str(delete_version_id).expect("test version id should parse")), + mod_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }) + .expect("delete marker metadata should be valid"); + fm.marshal_msg().expect("hidden metadata should encode") + } + + fn visible_object_metadata(name: &str, version_id: &str) -> Vec { + let mut fm = FileMeta::default(); + let mut fi = FileInfo::new(name, 1, 1); + fi.version_id = Some(Uuid::parse_str(version_id).expect("test version id should parse")); + fi.mod_time = Some(OffsetDateTime::now_utc()); + fm.add_version(fi).expect("object metadata should be valid"); + fm.marshal_msg().expect("visible metadata should encode") + } + + async fn scan_names(disk: &LocalDisk, bucket: &str, base_dir: &str, incl_deleted: bool) -> Vec { + let (reader, mut writer) = tokio::io::duplex(4096); + let mut out = MetacacheWriter::new(&mut writer); + let opts = WalkDirOptions { + bucket: bucket.to_string(), + base_dir: base_dir.to_string(), + recursive: false, + incl_deleted, + ..Default::default() + }; + let mut objs_returned = 0; + + disk.scan_dir(base_dir.to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None) + .await + .expect("scan_dir should succeed"); + out.close().await.expect("metacache writer should close"); + drop(out); + drop(writer); + + let mut reader = MetacacheReader::new(reader); + reader + .read_all() + .await + .expect("scan output should decode") + .into_iter() + .map(|entry| entry.name) + .collect() + } + + let dir = tempdir().expect("tempdir should be created"); + let bucket = "test-bucket"; + let bucket_dir = dir.path().join(bucket); + + let hidden_object = bucket_dir.join("hidden/deleted.txt"); + fs::create_dir_all(&hidden_object) + .await + .expect("hidden object dir should be created"); + fs::write( + hidden_object.join(STORAGE_FORMAT_FILE), + hidden_versioned_object_metadata( + "hidden/deleted.txt", + "11111111-1111-1111-1111-111111111111", + "22222222-2222-2222-2222-222222222222", + ), + ) + .await + .expect("hidden object metadata should be written"); + + let nested_hidden_object = bucket_dir.join("hidden/nested/deleted.txt"); + fs::create_dir_all(&nested_hidden_object) + .await + .expect("nested hidden object dir should be created"); + fs::write( + nested_hidden_object.join(STORAGE_FORMAT_FILE), + hidden_versioned_object_metadata( + "hidden/nested/deleted.txt", + "33333333-3333-3333-3333-333333333333", + "44444444-4444-4444-4444-444444444444", + ), + ) + .await + .expect("nested hidden object metadata should be written"); + + let visible_object = bucket_dir.join("visible/nested/object.txt"); + fs::create_dir_all(&visible_object) + .await + .expect("visible object dir should be created"); + fs::write( + visible_object.join(STORAGE_FORMAT_FILE), + visible_object_metadata("visible/nested/object.txt", "55555555-5555-5555-5555-555555555555"), + ) + .await + .expect("visible object metadata should be written"); + + let endpoint = + Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + + let root_names = scan_names(&disk, bucket, "", false).await; + assert!(!root_names.contains(&"hidden/".to_string())); + assert!(root_names.contains(&"visible/".to_string())); + + let hidden_names = scan_names(&disk, bucket, "hidden/", false).await; + assert!(!hidden_names.contains(&"hidden/nested/".to_string())); + + let visible_names = scan_names(&disk, bucket, "visible/", false).await; + assert!(visible_names.contains(&"visible/nested/".to_string())); + + let versioned_root_names = scan_names(&disk, bucket, "", true).await; + assert!(versioned_root_names.contains(&"hidden/".to_string())); + + let versioned_hidden_names = scan_names(&disk, bucket, "hidden/", true).await; + assert!(versioned_hidden_names.contains(&"hidden/nested/".to_string())); + } + #[cfg(unix)] #[tokio::test] async fn test_scan_dir_propagates_metadata_read_errors() { diff --git a/crates/ecstore/src/disk/mod.rs b/crates/ecstore/src/disk/mod.rs index 225fafa90..abba9a357 100644 --- a/crates/ecstore/src/disk/mod.rs +++ b/crates/ecstore/src/disk/mod.rs @@ -784,6 +784,10 @@ pub struct WalkDirOptions { // Do a full recursive scan. pub recursive: bool, + // Include entries hidden by delete markers. + #[serde(default)] + pub incl_deleted: bool, + // ReportNotFound will return errFileNotFound if all disks reports the BaseDir cannot be found. pub report_notfound: bool, @@ -1030,6 +1034,7 @@ mod tests { bucket: "test-bucket".to_string(), base_dir: "/path/to/dir".to_string(), recursive: true, + incl_deleted: false, report_notfound: false, filter_prefix: Some("prefix_".to_string()), forward_to: Some("object/path".to_string()), @@ -1041,6 +1046,7 @@ mod tests { assert_eq!(opts.bucket, "test-bucket"); assert_eq!(opts.base_dir, "/path/to/dir"); assert!(opts.recursive); + assert!(!opts.incl_deleted); assert!(!opts.report_notfound); assert_eq!(opts.filter_prefix, Some("prefix_".to_string())); assert_eq!(opts.forward_to, Some("object/path".to_string())); diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 43d7cf849..f4d947255 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -525,6 +525,7 @@ struct ListingSupplementOptions { bucket: String, path: String, recursive: bool, + incl_deleted: bool, filter_prefix: Option, forward_to: Option, per_disk_limit: i32, @@ -653,6 +654,7 @@ async fn read_fallback_listing_disk( bucket: options.bucket, base_dir: options.path, recursive: options.recursive, + incl_deleted: options.incl_deleted, report_notfound: false, filter_prefix: options.filter_prefix, forward_to: options.forward_to, @@ -1707,6 +1709,7 @@ impl ECStore { bucket: bucket.to_owned(), path: path.clone(), recursive: true, + incl_deleted: !opts.latest_only, filter_prefix: Some(filter_prefix.clone()), forward_to: opts.marker.clone(), per_disk_limit: bounded_usize_to_i32(opts.limit), @@ -1726,6 +1729,7 @@ impl ECStore { bucket: bucket.to_owned(), path, recursive: true, + incl_deleted: !opts.latest_only, filter_prefix: Some(filter_prefix), forward_to: opts.marker.clone(), min_disks: raw_min_disks, @@ -2032,7 +2036,7 @@ async fn gather_results( continue; } - if !opts.incl_deleted && entry.is_object() && entry.is_latest_delete_marker() && !entry.is_object_dir() { + if !opts.incl_deleted && entry.is_object() && entry.is_latest_delete_marker() { continue; } @@ -2824,6 +2828,7 @@ impl Sets { bucket: bucket.to_owned(), path: path.clone(), recursive: true, + incl_deleted: !opts.latest_only, filter_prefix: Some(filter_prefix.clone()), forward_to: opts.marker.clone(), per_disk_limit: bounded_usize_to_i32(opts.limit), @@ -2843,6 +2848,7 @@ impl Sets { bucket: bucket.to_owned(), path, recursive: true, + incl_deleted: !opts.latest_only, filter_prefix: Some(filter_prefix), forward_to: opts.marker.clone(), min_disks: raw_min_disks, @@ -3673,6 +3679,7 @@ impl SetDisks { bucket: bucket.clone(), path: opts.base_dir.clone(), recursive: opts.recursive, + incl_deleted: opts.incl_deleted, filter_prefix: opts.filter_prefix.clone(), forward_to: opts.marker.clone(), per_disk_limit: limit, @@ -3692,6 +3699,7 @@ impl SetDisks { bucket: opts.bucket, path: opts.base_dir, recursive: opts.recursive, + incl_deleted: opts.incl_deleted, filter_prefix: opts.filter_prefix, forward_to: opts.marker, min_disks: raw_min_disks, @@ -3917,6 +3925,27 @@ mod test { } } + fn test_delete_marker_meta_entry(name: &str) -> MetaCacheEntry { + let mut meta = FileMeta::new(); + meta.add_version(FileInfo { + volume: "bucket".to_owned(), + name: name.to_owned(), + deleted: true, + version_id: Some(Uuid::from_u128(1)), + mod_time: Some(time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")), + ..Default::default() + }) + .expect("test metadata should accept delete marker version"); + let metadata = meta.marshal_msg().expect("test metadata should marshal"); + + MetaCacheEntry { + name: name.to_owned(), + metadata, + cached: Some(meta), + reusable: false, + } + } + #[test] fn fallback_entries_for_object_filters_claimed_physical_disks() { let mut entries = FallbackListingEntries::new(); @@ -3948,6 +3977,7 @@ mod test { bucket: "bucket".to_string(), path: String::new(), recursive: true, + incl_deleted: false, filter_prefix: None, forward_to: None, per_disk_limit: 100, @@ -3961,6 +3991,7 @@ mod test { bucket: "bucket".to_string(), path: String::new(), recursive: true, + incl_deleted: false, filter_prefix: None, forward_to: None, per_disk_limit: 0, @@ -4007,6 +4038,7 @@ mod test { bucket: "bucket".to_string(), path: String::new(), recursive: true, + incl_deleted: false, filter_prefix: None, forward_to: None, per_disk_limit: 0, @@ -4234,6 +4266,57 @@ mod test { assert!(!cancel.is_cancelled()); } + #[tokio::test] + async fn list_path_gather_results_skips_directory_delete_marker_by_default() { + let (entry_tx, entry_rx) = mpsc::channel(4); + let (result_tx, mut result_rx) = mpsc::channel(1); + let cancel = CancellationToken::new(); + + entry_tx + .send(test_delete_marker_meta_entry("folder/")) + .await + .expect("directory delete marker should be queued"); + entry_tx + .send(test_object_meta_entry("visible")) + .await + .expect("visible object should be queued"); + drop(entry_tx); + + let handle = tokio::spawn(gather_results( + cancel.clone(), + ListPathOptions { + bucket: "bucket".to_owned(), + separator: Some("/".to_owned()), + include_directories: true, + limit: 2, + ..Default::default() + }, + entry_rx, + result_tx, + )); + + let result = timeout(Duration::from_secs(1), result_rx.recv()) + .await + .expect("eof result should be sent promptly") + .expect("eof result should be present"); + let entries = result.entries.expect("result entries should be present"); + let names = entries + .entries() + .into_iter() + .map(|entry| entry.name.clone()) + .collect::>(); + + assert_eq!(names, ["visible".to_string()]); + + let state = timeout(Duration::from_secs(1), handle) + .await + .expect("gather_results should finish after input closes") + .expect("gather_results task should not panic") + .expect("gather_results should succeed"); + assert_eq!(state, GatherResultsState::InputClosed); + assert!(!cancel.is_cancelled()); + } + #[test] fn list_path_forward_past_is_idempotent_for_same_marker() { let mut first_page = sorted_entries(&["obj-0001", "obj-0002", "obj-0003", "obj-0004"]);