fix(ecstore): list_object_v2 error when scanning multipart folder (#2946)

* fix(ecstore): list_object_v2 error when scanning multipart directory with prefix

Signed-off-by: w0od <dingboning02@163.com>

* test(ecstore): cover prefix dir scan with multipart folder support

Signed-off-by: w0od <dingboning02@163.com>

* test(ecstore): harden test shape to improve regression detection

Signed-off-by: w0od <dingboning02@163.com>

* refactor(ecstore): move multipart dir filter into recursion to reduce I/O

Signed-off-by: w0od <dingboning02@163.com>

* test(ecstore): replace scan_dir test with walk_dir integration coverage

Signed-off-by: w0od <dingboning02@163.com>

* refactor(ecstore): remove unnecessary .clone() calls

Signed-off-by: w0od <dingboning02@163.com>

---------

Signed-off-by: w0od <dingboning02@163.com>
This commit is contained in:
wood
2026-05-18 19:01:20 +08:00
committed by GitHub
parent 4c9fd789ea
commit 0e888cfef2
+219 -16
View File
@@ -1234,6 +1234,7 @@ impl LocalDisk {
}
#[async_recursion::async_recursion]
#[allow(clippy::too_many_arguments)]
async fn scan_dir<W>(
&self,
mut current: String,
@@ -1242,6 +1243,7 @@ impl LocalDisk {
out: &mut MetacacheWriter<W>,
objs_returned: &mut i32,
skip_current_dir_object: bool,
multipart_dir_to_skip: Option<HashSet<String>>,
) -> Result<()>
where
W: AsyncWrite + Unpin + Send,
@@ -1298,6 +1300,14 @@ impl LocalDisk {
if opts.limit > 0 && *objs_returned >= opts.limit {
return Ok(());
}
// check multipart dir
if skip_current_dir_object
&& let Some(ref dir_to_skip) = multipart_dir_to_skip
&& dir_to_skip.contains(entry.trim_end_matches(SLASH_SEPARATOR))
{
*item = "".to_owned();
continue;
}
// check prefix
if !prefix.is_empty() && !entry.starts_with(prefix.as_str()) {
*item = "".to_owned();
@@ -1367,15 +1377,25 @@ impl LocalDisk {
}
}
let mut dir_stack: Vec<(String, bool)> = Vec::with_capacity(5);
let mut dir_stack: Vec<(String, bool, Option<HashSet<String>>)> = Vec::with_capacity(5);
// Explicit directory markers and real directories can resolve to the same logical path.
let schedule_dir = |dir_stack: &mut Vec<(String, bool)>, dir_name: String, skip_object: bool| {
if let Some((last_dir_name, existing_skip_object)) = dir_stack.last_mut()
let schedule_dir = |dir_stack: &mut Vec<(String, bool, Option<HashSet<String>>)>,
dir_name: String,
skip_object: bool,
dir_to_skip: Option<HashSet<String>>| {
if let Some((last_dir_name, existing_skip_object, existing_dir_to_skip)) = dir_stack.last_mut()
&& *last_dir_name == dir_name
{
*existing_skip_object |= skip_object;
if let Some(existing_dir_to_skip) = existing_dir_to_skip {
if let Some(new_dir_to_skip) = &dir_to_skip {
existing_dir_to_skip.extend(new_dir_to_skip.iter().cloned());
}
} else {
*existing_dir_to_skip = dir_to_skip;
}
} else {
dir_stack.push((dir_name, skip_object));
dir_stack.push((dir_name, skip_object, dir_to_skip));
}
};
prefix = "".to_owned();
@@ -1391,9 +1411,10 @@ impl LocalDisk {
let name = path_join_buf(&[current.as_str(), entry.as_str()]);
while let Some((pop, skip_object)) = dir_stack.last().cloned()
&& pop < name
while let Some((last_name, _, _)) = dir_stack.last()
&& *last_name < name
{
let (pop, skip_object, dir_to_skip) = dir_stack.pop().unwrap();
out.write_obj(&MetaCacheEntry {
name: pop.clone(),
..Default::default()
@@ -1401,11 +1422,11 @@ impl LocalDisk {
.await?;
if opts.recursive
&& let Err(er) = Box::pin(self.scan_dir(pop, prefix.clone(), opts, out, objs_returned, skip_object)).await
&& let Err(er) =
Box::pin(self.scan_dir(pop, prefix.clone(), opts, out, objs_returned, skip_object, dir_to_skip)).await
{
error!("scan_dir err {:?}", er);
}
dir_stack.pop();
}
let mut meta = MetaCacheEntry {
@@ -1442,11 +1463,24 @@ impl LocalDisk {
// }
if opts.recursive {
let mut dir_to_skip = HashSet::new();
if let Ok(file_meta) = FileMeta::load(&res)
&& let Ok(data_dirs) = file_meta.get_data_dirs()
{
for data_dir in data_dirs.iter().flatten() {
dir_to_skip.insert(data_dir.to_string());
}
}
let mut dir_name = meta.name.clone();
if !dir_name.ends_with(SLASH_SEPARATOR) {
dir_name.push_str(SLASH_SEPARATOR);
}
schedule_dir(&mut dir_stack, dir_name, true);
schedule_dir(
&mut dir_stack,
dir_name,
true,
if dir_to_skip.is_empty() { None } else { Some(dir_to_skip) },
);
}
}
Err(err) => {
@@ -1455,7 +1489,7 @@ 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);
schedule_dir(&mut dir_stack, meta.name, false, None);
}
}
@@ -1464,7 +1498,7 @@ impl LocalDisk {
};
}
while let Some((dir, skip_object)) = dir_stack.pop() {
while let Some((dir, skip_object, dir_to_skip)) = dir_stack.pop() {
if opts.limit > 0 && *objs_returned >= opts.limit {
return Ok(());
}
@@ -1476,7 +1510,8 @@ impl LocalDisk {
.await?;
if opts.recursive
&& let Err(er) = Box::pin(self.scan_dir(dir, prefix.clone(), opts, out, objs_returned, skip_object)).await
&& let Err(er) =
Box::pin(self.scan_dir(dir, prefix.clone(), opts, out, objs_returned, skip_object, dir_to_skip)).await
{
warn!("scan_dir err {:?}", &er);
}
@@ -2328,6 +2363,7 @@ impl DiskAPI for LocalDisk {
let mut objs_returned = 0;
let mut skip_current_dir_object = false;
let mut multipart_dir_to_skip: HashSet<String> = HashSet::new();
if opts.base_dir.ends_with(SLASH_SEPARATOR) {
if let Ok(data) = self
.read_metadata(
@@ -2351,10 +2387,23 @@ impl DiskAPI for LocalDisk {
let fpath =
self.get_object_path(&opts.bucket, path_join_buf(&[opts.base_dir.as_str(), STORAGE_FORMAT_FILE]).as_str())?;
if let Ok(meta) = tokio::fs::metadata(fpath).await
if let Ok(meta) = tokio::fs::metadata(&fpath).await
&& meta.is_file()
{
skip_current_dir_object = true;
if let Ok(meta_bytes) = self
.read_metadata(
opts.bucket.as_str(),
path_join_buf(&[opts.base_dir.as_str(), STORAGE_FORMAT_FILE]).as_str(),
)
.await
&& let Ok(file_meta) = FileMeta::load(&meta_bytes)
&& let Ok(data_dirs) = file_meta.get_data_dirs()
{
for data_dir in data_dirs.iter().flatten() {
multipart_dir_to_skip.insert(data_dir.to_string());
}
}
}
}
}
@@ -2366,6 +2415,11 @@ impl DiskAPI for LocalDisk {
&mut out,
&mut objs_returned,
skip_current_dir_object,
if multipart_dir_to_skip.is_empty() {
None
} else {
Some(multipart_dir_to_skip)
},
)
.await?;
@@ -3151,7 +3205,7 @@ mod test {
};
let mut objs_returned = 0;
disk.scan_dir("".to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false)
disk.scan_dir("".to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None)
.await
.unwrap();
out.close().await.unwrap();
@@ -3206,7 +3260,7 @@ mod test {
};
let mut objs_returned = 0;
disk.scan_dir("marker/".to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false)
disk.scan_dir("marker/".to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None)
.await
.unwrap();
out.close().await.unwrap();
@@ -3266,7 +3320,7 @@ mod test {
};
let mut objs_returned = 0;
disk.scan_dir(base_dir.to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false)
disk.scan_dir(base_dir.to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None)
.await
.unwrap();
out.close().await.unwrap();
@@ -3327,6 +3381,155 @@ mod test {
assert_eq!(double_count as usize, double_names.len());
}
#[tokio::test]
async fn test_walk_dir_ignore_multipart_dirs() {
use rustfs_filemeta::MetacacheReader;
use tempfile::tempdir;
const UUID_MULTIPART_1: &str = "8b262d24-fcf9-473d-a4cd-f9b27f24f60e";
const UUID_MULTIPART_2: &str = "fbf3183c-63be-45cc-b3bf-424ddb7f95f8";
const UUID_OBJ: &str = "db8b9b74-9016-4f9e-83e9-82a772947d28";
const VER_ID_1: &str = "c683f9f8-c0a1-4bc5-8a67-0faafa839a1a";
const VER_ID_2: &str = "a4b84f6e-c8ba-461b-8f9d-43feb0893efb";
const VER_ID_3: &str = "892c9ae7-2bb3-44ee-9a71-bc7ddf08d765";
const BASE_DIR: &str = "dir1/obj/";
const MULTIPART_DIR: &str = "multipart-file";
const DIR_IN_MULTIPART_DIR: &str = "dir-in-multipart";
const EMPTY_STR: &str = "";
let parse_uuid = |s: &str| Uuid::parse_str(s).unwrap();
let create_file_info = |version_id: &str, data_dir: &str| FileInfo {
version_id: Some(parse_uuid(version_id)),
data_dir: Some(parse_uuid(data_dir)),
mod_time: Some(OffsetDateTime::now_utc()),
..Default::default()
};
let dir = tempdir().unwrap();
let obj_base = dir.path().join("test-bucket").join(BASE_DIR);
let multipart_base = obj_base.join(MULTIPART_DIR);
let dir_in_multipart_base = multipart_base.join(DIR_IN_MULTIPART_DIR);
fs::create_dir_all(&multipart_base).await.unwrap();
for uuid in &[UUID_MULTIPART_1, UUID_MULTIPART_2] {
fs::create_dir_all(multipart_base.join(uuid)).await.unwrap();
fs::write(multipart_base.join(uuid).join("part.1"), b"part").await.unwrap();
}
fs::create_dir_all(obj_base.join(UUID_OBJ)).await.unwrap();
fs::write(obj_base.join(UUID_OBJ).join("part.1"), b"part").await.unwrap();
fs::create_dir_all(&dir_in_multipart_base).await.unwrap();
fs::write(dir_in_multipart_base.join(STORAGE_FORMAT_FILE), b"meta")
.await
.unwrap();
let mut fm = FileMeta::default();
fm.add_version(create_file_info(VER_ID_1, UUID_MULTIPART_1)).unwrap();
fm.add_version(create_file_info(VER_ID_2, UUID_MULTIPART_2)).unwrap();
fs::write(multipart_base.join(STORAGE_FORMAT_FILE), fm.marshal_msg().unwrap())
.await
.unwrap();
let mut fm = FileMeta::default();
fm.add_version(create_file_info(VER_ID_3, UUID_OBJ)).unwrap();
fs::write(obj_base.join(STORAGE_FORMAT_FILE), fm.marshal_msg().unwrap())
.await
.unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let (reader, mut writer) = tokio::io::duplex(4096);
disk.walk_dir(
WalkDirOptions {
bucket: "test-bucket".to_string(),
base_dir: BASE_DIR.to_string(),
recursive: true,
filter_prefix: Some(EMPTY_STR.to_string()),
..Default::default()
},
&mut writer,
)
.await
.unwrap();
MetacacheWriter::new(&mut writer).close().await.unwrap();
let mut reader = MetacacheReader::new(reader);
let entries = reader.read_all().await.unwrap();
let names: Vec<String> = entries.into_iter().map(|entry| entry.name).collect();
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}", BASE_DIR, MULTIPART_DIR))
.count(),
1
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}/", BASE_DIR, MULTIPART_DIR))
.count(),
1
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}/{}", BASE_DIR, MULTIPART_DIR, DIR_IN_MULTIPART_DIR))
.count(),
1
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}/{}/", BASE_DIR, MULTIPART_DIR, DIR_IN_MULTIPART_DIR))
.count(),
1
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}/{}", BASE_DIR, MULTIPART_DIR, UUID_MULTIPART_1))
.count(),
0
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}/{}", BASE_DIR, MULTIPART_DIR, UUID_MULTIPART_2))
.count(),
0
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}", BASE_DIR, UUID_OBJ))
.count(),
0
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}/{}/", BASE_DIR, MULTIPART_DIR, UUID_MULTIPART_1))
.count(),
0
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}/{}/", BASE_DIR, MULTIPART_DIR, UUID_MULTIPART_2))
.count(),
0
);
assert_eq!(
names
.iter()
.filter(|name| *name == &format!("{}{}/", BASE_DIR, UUID_OBJ))
.count(),
0
);
}
#[tokio::test]
async fn test_make_volume() {
let p = "./testv0";