diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index ee45c5509..5fba72685 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -1234,6 +1234,7 @@ impl LocalDisk { } #[async_recursion::async_recursion] + #[allow(clippy::too_many_arguments)] async fn scan_dir( &self, mut current: String, @@ -1242,6 +1243,7 @@ impl LocalDisk { out: &mut MetacacheWriter, objs_returned: &mut i32, skip_current_dir_object: bool, + multipart_dir_to_skip: Option>, ) -> 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>)> = 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>)>, + dir_name: String, + skip_object: bool, + dir_to_skip: Option>| { + 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 = 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 = 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";