diff --git a/.config/e2e-smoke-selection.txt b/.config/e2e-smoke-selection.txt index 78a1d1484..2ac994976 100644 --- a/.config/e2e-smoke-selection.txt +++ b/.config/e2e-smoke-selection.txt @@ -1 +1 @@ -sha256=548a58d74f3cb3b5a1ab9dcd4d3f8625b9ddebe847a2e14de1d08c3b2e6292eb +sha256=5de6852b2e606eca4ee0d86dbb3b725a3eb8923143fd00df96c549acc36040bd diff --git a/crates/e2e_test/src/list_objects_duplicates_test.rs b/crates/e2e_test/src/list_objects_duplicates_test.rs index 9b8349dcc..bc8e33ab6 100644 --- a/crates/e2e_test/src/list_objects_duplicates_test.rs +++ b/crates/e2e_test/src/list_objects_duplicates_test.rs @@ -311,4 +311,233 @@ mod tests { env.stop_server(); } + + async fn collect_prefix_pages( + client: &Client, + bucket: &str, + prefix: &str, + delimiter: Option<&str>, + max_keys: i32, + start_after: Option<&str>, + ) -> (Vec, Vec) { + let mut objects = Vec::new(); + let mut prefixes = Vec::new(); + let mut token = None; + let mut seen_tokens = std::collections::HashSet::new(); + for _ in 0..64 { + let page = client + .list_objects_v2() + .bucket(bucket) + .prefix(prefix) + .set_delimiter(delimiter.map(ToOwned::to_owned)) + .max_keys(max_keys) + .set_start_after(start_after.map(ToOwned::to_owned)) + .set_continuation_token(token) + .send() + .await + .expect("prefix pagination request should succeed"); + let count = page.contents().len() + page.common_prefixes().len(); + assert!(count <= usize::try_from(max_keys).expect("positive page size")); + assert_eq!(page.key_count(), Some(i32::try_from(count).expect("small page"))); + objects.extend_from_slice(page.contents()); + prefixes.extend( + page.common_prefixes() + .iter() + .map(|item| item.prefix().expect("CommonPrefix must contain a prefix").to_owned()), + ); + if page.is_truncated() == Some(false) { + assert!(page.next_continuation_token().is_none(), "final page must not have a next token"); + return (objects, prefixes); + } + assert_eq!(page.is_truncated(), Some(true)); + let next = page.next_continuation_token().expect("truncated page requires a token"); + assert!(!next.is_empty()); + assert!(seen_tokens.insert(next.to_owned()), "continuation token must advance"); + token = Some(next.to_owned()); + } + panic!("prefix pagination exceeded its finite page budget"); + } + + fn object_keys(objects: &[aws_sdk_s3::types::Object]) -> Vec<&str> { + objects + .iter() + .map(|object| object.key().expect("listed object must have a key")) + .collect() + } + + /// Issue #8175: the prefix object must occur once, including across page boundaries. + #[tokio::test] + async fn test_list_objects_v2_prefix_marker_pagination() { + init_logging(); + let mut env = RustFSTestEnvironment::new().await.expect("create test environment"); + env.start_rustfs_server(vec![]).await.expect("start RustFS"); + let client = create_s3_client(&env); + let bucket = "test-prefix-marker-pagination"; + create_bucket(&client, bucket).await.expect("create bucket"); + let mut expected = vec!["content/".to_owned()]; + expected.extend((0..25).map(|index| format!("content/0123456789abcdef0123456789abcdef/{index:03}"))); + for key in &expected { + client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(if key == "content/" { b"" } else { b"x" })) + .send() + .await + .expect("create prefix fixture"); + } + for max_keys in [1000, 1, 2] { + let (objects, prefixes) = collect_prefix_pages(&client, bucket, "content/", None, max_keys, None).await; + assert_eq!(object_keys(&objects), expected, "all keys must appear once at MaxKeys={max_keys}"); + assert!(prefixes.is_empty()); + } + let (objects, prefixes) = collect_prefix_pages(&client, bucket, "content/", None, 1, Some("content/")).await; + assert_eq!(object_keys(&objects), expected[1..]); + assert!(prefixes.is_empty()); + + // V1 shares the storage listing path but resumes with a key marker. + let mut marker = Some("content/".to_owned()); + let mut v1_keys = Vec::new(); + let mut finished = false; + for _ in 0..32 { + let page = client + .list_objects() + .bucket(bucket) + .prefix("content/") + .max_keys(1) + .set_marker(marker.clone()) + .send() + .await + .expect("list V1 continuation page"); + assert!(page.contents().len() <= 1); + let keys = object_keys(page.contents()); + v1_keys.extend(keys.iter().map(|key| (*key).to_owned())); + if page.is_truncated() == Some(false) { + finished = true; + break; + } + assert_eq!(page.is_truncated(), Some(true)); + let next = page + .next_marker() + .or_else(|| keys.last().copied()) + .expect("V1 page must advance"); + assert!(marker.as_deref().is_none_or(|previous| next > previous)); + marker = Some(next.to_owned()); + } + assert!(finished, "V1 pagination must terminate"); + assert_eq!(v1_keys, expected[1..]); + env.stop_server(); + } + + /// An exact prefix match does not establish EOF, including for ordinary keys. + #[tokio::test] + async fn test_list_objects_v2_exact_prefix_pagination_boundaries() { + init_logging(); + let mut env = RustFSTestEnvironment::new().await.expect("create test environment"); + env.start_rustfs_server(vec![]).await.expect("start RustFS"); + let client = create_s3_client(&env); + let bucket = "test-exact-prefix-boundaries"; + create_bucket(&client, bucket).await.expect("create bucket"); + for key in [ + "a", + "ab", + "solo/", + "marker/", + "marker/file", + "marker/subdir/", + "marker/subdir/file", + ] { + client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(b"")) + .send() + .await + .expect("create boundary fixture"); + } + let (objects, prefixes) = collect_prefix_pages(&client, bucket, "a", None, 1, None).await; + assert_eq!(object_keys(&objects), vec!["a", "ab"]); + assert!(prefixes.is_empty()); + let first = client + .list_objects() + .bucket(bucket) + .prefix("a") + .max_keys(1) + .send() + .await + .expect("list V1 exact-prefix first page"); + assert_eq!(object_keys(first.contents()), vec!["a"]); + assert_eq!(first.is_truncated(), Some(true)); + let last = client + .list_objects() + .bucket(bucket) + .prefix("a") + .max_keys(1) + .marker(first.next_marker().unwrap_or("a")) + .send() + .await + .expect("list V1 exact-prefix final page"); + assert_eq!(object_keys(last.contents()), vec!["ab"]); + assert_eq!(last.is_truncated(), Some(false)); + + let solo = client + .list_objects_v2() + .bucket(bucket) + .prefix("solo/") + .max_keys(1) + .send() + .await + .expect("list isolated marker"); + assert_eq!(object_keys(solo.contents()), vec!["solo/"]); + assert_eq!(solo.is_truncated(), Some(false)); + assert!(solo.next_continuation_token().is_none()); + let (objects, prefixes) = collect_prefix_pages(&client, bucket, "marker/", Some("/"), 1, None).await; + assert_eq!(object_keys(&objects), vec!["marker/", "marker/file"]); + assert_eq!(prefixes, vec!["marker/subdir/"]); + env.stop_server(); + } + + /// A directory marker must retain its own metadata when a same-named object exists. + #[tokio::test] + async fn test_list_objects_v2_prefix_marker_preserves_metadata() { + init_logging(); + let mut env = RustFSTestEnvironment::new().await.expect("create test environment"); + env.start_rustfs_server(vec![]).await.expect("start RustFS"); + let client = create_s3_client(&env); + let bucket = "test-prefix-marker-metadata"; + create_bucket(&client, bucket).await.expect("create bucket"); + let plain = client + .put_object() + .bucket(bucket) + .key("content") + .body(ByteStream::from_static(b"plain object body")) + .send() + .await + .expect("put plain object"); + let directory = client + .put_object() + .bucket(bucket) + .key("content/") + .body(ByteStream::from_static(b"")) + .send() + .await + .expect("put directory marker"); + assert_ne!(plain.e_tag(), directory.e_tag(), "fixture ETags must distinguish the objects"); + for max_keys in [1, 1000] { + let (objects, prefixes) = collect_prefix_pages(&client, bucket, "content/", None, max_keys, None).await; + assert_eq!(object_keys(&objects), vec!["content/"]); + assert!(prefixes.is_empty()); + assert_eq!(objects[0].size(), Some(0)); + assert_eq!(objects[0].e_tag(), directory.e_tag()); + let (objects, prefixes) = collect_prefix_pages(&client, bucket, "content", None, max_keys, None).await; + assert_eq!(object_keys(&objects), vec!["content", "content/"]); + assert!(prefixes.is_empty()); + assert_eq!(objects[0].size(), Some(17)); + assert_eq!(objects[0].e_tag(), plain.e_tag()); + assert_eq!(objects[1].size(), Some(0)); + assert_eq!(objects[1].e_tag(), directory.e_tag()); + } + env.stop_server(); + } } diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 077c8a90c..aa8a2c39a 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -10447,28 +10447,30 @@ impl DiskAPI for LocalDisk { }; write_metacache_obj(&mut out, &meta).await?; objs_returned += 1; - } else { - let fpath = self - .io_get_object_path(&opts.bucket, path_join_buf(&[opts.base_dir.as_str(), STORAGE_FORMAT_FILE]).as_str())?; + } - if let Ok(meta) = with_walk_stall_deadline(stall, tokio::fs::metadata(&fpath)).await? - && meta.is_file() + // A plain object shares this directory with its children even when + // an explicit directory marker exists in the sibling encoded path. + let fpath = + self.io_get_object_path(&opts.bucket, path_join_buf(&[opts.base_dir.as_str(), STORAGE_FORMAT_FILE]).as_str())?; + + if let Ok(meta) = with_walk_stall_deadline(stall, tokio::fs::metadata(&fpath)).await? + && meta.is_file() + { + skip_current_dir_object = true; + if let Ok(meta_bytes) = with_walk_stall_deadline( + stall, + 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() { - skip_current_dir_object = true; - if let Ok(meta_bytes) = with_walk_stall_deadline( - stall, - 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()); - } + for data_dir in data_dirs.iter().flatten() { + multipart_dir_to_skip.insert(data_dir.to_string()); } } } @@ -18707,6 +18709,89 @@ mod test { ); } + async fn assert_walk_dir_prefix_marker_entries(with_plain_object: bool) { + use rustfs_filemeta::MetacacheReader; + use tempfile::tempdir; + + async fn write_metadata(root: &Path, name: &str, size: i64, data_dir: Option) -> Vec { + let object_dir = root.join(encode_dir_object(name)); + fs::create_dir_all(&object_dir) + .await + .expect("object directory should be created"); + let mut file_info = FileInfo::new(name, 1, 1); + file_info.size = size; + file_info.data_dir = data_dir; + file_info.mod_time = Some(OffsetDateTime::now_utc()); + let mut metadata = FileMeta::default(); + metadata.add_version(file_info).expect("object metadata should be valid"); + let bytes = metadata.marshal_msg().expect("object metadata should encode"); + fs::write(object_dir.join(STORAGE_FORMAT_FILE), &bytes) + .await + .expect("object metadata should be written"); + bytes + } + + let dir = tempdir().expect("temporary disk should be created"); + let bucket = "test-bucket"; + let bucket_dir = dir.path().join(bucket); + let marker_metadata = write_metadata(&bucket_dir, "content/", 0, None).await; + let child_metadata = write_metadata(&bucket_dir, "content/child", 7, None).await; + let data_dir = Uuid::parse_str("bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb").expect("data directory UUID should parse"); + if with_plain_object { + write_metadata(&bucket_dir, "content", 19, Some(data_dir)).await; + // A storage data directory can contain metadata-bearing subdirectories. + // The entire directory must be skipped rather than exposed as objects. + write_metadata(&bucket_dir, &format!("content/{data_dir}/segment"), 31, None).await; + fs::write(bucket_dir.join("content").join(data_dir.to_string()).join("part.1"), b"part") + .await + .expect("object part should be written"); + } + + let endpoint = + Endpoint::try_from(dir.path().to_str().expect("disk path should be UTF-8")).expect("disk endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should initialize"); + let (reader, mut writer) = tokio::io::duplex(65536); + disk.walk_dir( + WalkDirOptions { + bucket: bucket.to_owned(), + base_dir: "content/".to_owned(), + recursive: true, + ..Default::default() + }, + &mut writer, + ) + .await + .expect("prefix walk should succeed"); + let entries = MetacacheReader::new(reader) + .read_all() + .await + .expect("walk output should decode"); + assert!( + entries + .iter() + .all(|entry| !entry.name.starts_with(&format!("content/{data_dir}"))), + "plain object storage directories must not appear in the walk" + ); + let objects: Vec<_> = entries.into_iter().filter(MetaCacheEntry::is_object).collect(); + assert_eq!( + objects.iter().map(|entry| entry.name.as_str()).collect::>(), + vec!["content/", "content/child"], + "the marker and child must each appear exactly once" + ); + assert_eq!(objects[0].metadata, marker_metadata, "the marker must retain its own metadata"); + assert_eq!(objects[1].metadata, child_metadata, "the child must retain its own metadata"); + } + + #[tokio::test] + async fn test_walk_dir_prefix_marker_with_children_is_unique() { + assert_walk_dir_prefix_marker_entries(false).await; + } + + #[tokio::test] + async fn test_walk_dir_prefix_marker_skips_plain_object_and_parts() { + assert_walk_dir_prefix_marker_entries(true).await; + } + #[tokio::test] async fn test_scan_dir_reports_base_dir_object_metadata() { use rustfs_filemeta::MetacacheReader; diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 9c5342fac..fc9121c2f 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -33,8 +33,7 @@ use crate::core::sets::Sets; use crate::disk::error::DiskError; use crate::disk::{DiskAPI, DiskInfo, DiskStore, RUSTFS_META_BUCKET, WalkDirOptions}; use crate::error::{ - Error, Result, StorageError, is_all_disk_not_found, is_all_not_found, is_all_volume_not_found, is_err_bucket_not_found, - to_object_err, + Error, Result, StorageError, is_all_disk_not_found, is_all_not_found, is_all_volume_not_found, to_object_err, }; use crate::object_api::{ObjectInfo, ObjectOptions}; use crate::set_disk::SetDisks; @@ -4023,32 +4022,6 @@ impl ECStore { return Ok(result); } - // Optimization: use get for single object lookup with exact prefix - if !opts.prefix.is_empty() && max_keys == 1 && opts.marker.is_none() && !incl_deleted { - match self - .get_object_info( - &opts.bucket, - &opts.prefix, - &ObjectOptions { - no_lock: true, - ..Default::default() - }, - ) - .await - { - Ok(res) if !res.delete_marker && res.version_purge_status.is_empty() => { - return Ok(ListObjectsInfo { - objects: vec![res], - ..Default::default() - }); - } - Err(err) if is_err_bucket_not_found(&err) => { - return Err(err); - } - _ => {} - }; - }; - let mut list_result = self .clone() .list_path(&opts) @@ -5534,31 +5507,6 @@ impl Sets { // (notably `forward_past`) — see backlog#1047. opts.parse_marker(); - if !opts.prefix.is_empty() && max_keys == 1 && opts.marker.is_none() && !incl_deleted { - match self - .get_object_info( - &opts.bucket, - &opts.prefix, - &ObjectOptions { - no_lock: true, - ..Default::default() - }, - ) - .await - { - Ok(res) if !res.delete_marker && res.version_purge_status.is_empty() => { - return Ok(ListObjectsInfo { - objects: vec![res], - ..Default::default() - }); - } - Err(err) if is_err_bucket_not_found(&err) => { - return Err(err); - } - _ => {} - }; - } - let mut list_result = self .list_path(&opts) .await @@ -6448,33 +6396,6 @@ impl SetDisks { // (notably `forward_past`) — see backlog#1047. opts.parse_marker(); - if !opts.prefix.is_empty() && max_keys == 1 && opts.marker.is_none() { - match self - .get_object_info( - &opts.bucket, - &opts.prefix, - &ObjectOptions { - no_lock: true, - ..Default::default() - }, - ) - .await - { - Ok(res) if !res.delete_marker && res.version_purge_status.is_empty() => { - return Ok(ListObjectsInfo { - objects: vec![res], - ..Default::default() - }); - } - Ok(_) => {} - Err(err) => { - if is_err_bucket_not_found(&err) { - return Err(err); - } - } - }; - } - let mut list_result = self .list_path_result(&opts) .await @@ -9637,6 +9558,88 @@ mod test { } } + #[tokio::test] + async fn list_objects_exact_prefix_paginates_across_storage_layers() { + use crate::bucket::metadata_sys::{init_bucket_metadata_sys, test_support::isolated_store_over_temp_disks}; + use crate::object_api::{ObjectOptions, PutObjReader}; + use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; + use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _}; + + let (_dirs, store) = isolated_store_over_temp_disks().await; + let bucket = "exact-prefix-pagination-bucket"; + init_bucket_metadata_sys(store.clone(), Vec::new()).await; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("exact-prefix pagination bucket should be created"); + for name in ["a", "ab"] { + store.pools[0] + .put_object( + bucket, + name, + &mut PutObjReader::from_vec(b"pagination fixture".to_vec()), + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + .expect("pagination fixture object should be written"); + let info = store + .get_object_info( + bucket, + name, + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + .expect("exact-prefix fixture must be readable by object lookup"); + assert!(!info.delete_marker && info.version_purge_status.is_empty()); + } + + for layer in 0..3 { + for (prefix, expected) in [("a", vec!["a", "ab"]), ("ab", vec!["ab"]), ("missing", vec![])] { + let mut marker = None; + for page in 0..expected.len().max(1) { + let result = match layer { + 0 => { + store + .clone() + .list_objects_generic(bucket, prefix, marker.clone(), None, 1, false) + .await + } + 1 => { + store.pools[0] + .clone() + .list_objects_generic(bucket, prefix, marker.clone(), None, 1, false) + .await + } + _ => { + store.pools[0].disk_set[0] + .clone() + .list_objects_generic(bucket, prefix, marker.clone(), None, 1, false) + .await + } + } + .expect("exact-prefix page should list successfully"); + let names: Vec<_> = result.objects.iter().map(|object| object.name.as_str()).collect(); + let expected_page: Vec<_> = expected.get(page).copied().into_iter().collect(); + assert_eq!(names, expected_page, "layer {layer}, prefix {prefix}, page {page}"); + assert!(result.prefixes.is_empty(), "recursive listing should not return common prefixes"); + let has_more = page + 1 < expected.len(); + assert_eq!(result.is_truncated, has_more, "layer {layer}, prefix {prefix}, page {page}"); + assert_eq!(result.next_marker.is_some(), has_more, "only non-final pages should carry a marker"); + if has_more { + assert_ne!(result.next_marker, marker, "pagination marker must advance"); + } + marker = result.next_marker; + } + } + } + } + #[tokio::test] async fn list_objects_hides_pending_version_purge_across_walk_and_exact_prefix() { use crate::bucket::metadata_sys::{init_bucket_metadata_sys, test_support::isolated_store_over_temp_disks}; @@ -9685,7 +9688,7 @@ mod test { .expect("exact-prefix listing should succeed"); assert!( exact.objects.is_empty(), - "exact-prefix max_keys=1 shortcut should hide pending version-purge entries" + "exact-prefix max_keys=1 listing should hide pending version-purge entries" ); assert!(exact.prefixes.is_empty()); }