fix(ecstore): hide deleted versioned folder prefixes (#4296)

This commit is contained in:
cxymds
2026-07-06 14:04:53 +08:00
committed by GitHub
parent 005197140e
commit 2585855d23
4 changed files with 299 additions and 2 deletions
@@ -110,6 +110,7 @@ pub struct ListPathRawOptions {
pub bucket: String,
pub path: String,
pub recursive: bool,
pub incl_deleted: bool,
pub filter_prefix: Option<String>,
pub forward_to: Option<String>,
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(),
+205 -1
View File
@@ -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<bool> {
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, &current, -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<AtomicU32>);
@@ -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<u8> {
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<u8> {
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<String> {
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() {
+6
View File
@@ -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()));
+84 -1
View File
@@ -525,6 +525,7 @@ struct ListingSupplementOptions {
bucket: String,
path: String,
recursive: bool,
incl_deleted: bool,
filter_prefix: Option<String>,
forward_to: Option<String>,
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::<Vec<_>>();
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"]);