fix(s3): preserve prefix listing pagination and marker metadata

Merge the reviewed fix from pull request #8181.
This commit is contained in:
cxymds
2026-09-28 18:47:03 +08:00
committed by GitHub
parent 671238458c
commit e0a973bf4e
4 changed files with 419 additions and 102 deletions
+1 -1
View File
@@ -1 +1 @@
sha256=548a58d74f3cb3b5a1ab9dcd4d3f8625b9ddebe847a2e14de1d08c3b2e6292eb
sha256=5de6852b2e606eca4ee0d86dbb3b725a3eb8923143fd00df96c549acc36040bd
@@ -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<aws_sdk_s3::types::Object>, Vec<String>) {
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();
}
}
+105 -20
View File
@@ -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<Uuid>) -> Vec<u8> {
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<_>>(),
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;
+84 -81
View File
@@ -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());
}