mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-04 19:25:40 +00:00
Co-authored-by: houseme <housemecn@gmail.com> Co-authored-by: 安正超 <anzhengchao@gmail.com>
This commit is contained in:
@@ -31,6 +31,7 @@ mod tests {
|
|||||||
use aws_sdk_s3::Client;
|
use aws_sdk_s3::Client;
|
||||||
use aws_sdk_s3::primitives::ByteStream;
|
use aws_sdk_s3::primitives::ByteStream;
|
||||||
use serial_test::serial;
|
use serial_test::serial;
|
||||||
|
use std::collections::HashSet;
|
||||||
use tracing::info;
|
use tracing::info;
|
||||||
|
|
||||||
/// Helper function to create an S3 client for testing
|
/// Helper function to create an S3 client for testing
|
||||||
@@ -57,6 +58,130 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Test for Issue #2775: continuation forwarding must not
|
||||||
|
/// skip a child directory when the prefix component repeats in the key.
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn test_list_objects_v2_repeated_prefix_continuation() {
|
||||||
|
init_logging();
|
||||||
|
info!("Starting test: ListObjectsV2 repeated-prefix continuation");
|
||||||
|
|
||||||
|
const PAGE_SIZE: i32 = 2;
|
||||||
|
|
||||||
|
let mut env = RustFSTestEnvironment::new().await.expect("Failed to create test environment");
|
||||||
|
env.start_rustfs_server_with_env(vec![], &[("RUSTFS_CONSOLE_ENABLE", "false")])
|
||||||
|
.await
|
||||||
|
.expect("Failed to start RustFS");
|
||||||
|
|
||||||
|
let client = create_s3_client(&env);
|
||||||
|
let bucket = "test-repeated-prefix-small";
|
||||||
|
let prefix = "engineering/";
|
||||||
|
|
||||||
|
create_bucket(&client, bucket).await.expect("Failed to create bucket");
|
||||||
|
|
||||||
|
let expected_keys = vec![
|
||||||
|
format!("{prefix}alpha-000/artifact.txt"),
|
||||||
|
format!("{prefix}engineering/engineering/project-000/artifact.txt"),
|
||||||
|
format!("{prefix}engineering/engineering/project-001/artifact.txt"),
|
||||||
|
format!("{prefix}engineering/project-000/artifact.txt"),
|
||||||
|
format!("{prefix}engineering/project-001/artifact.txt"),
|
||||||
|
format!("{prefix}engineering/project-002/artifact.txt"),
|
||||||
|
format!("{prefix}zulu-000/artifact.txt"),
|
||||||
|
];
|
||||||
|
let noise_keys = [
|
||||||
|
"different/prefix/prefix/project-000/artifact.txt",
|
||||||
|
"engineering-other/project-000/artifact.txt",
|
||||||
|
"unrelated/engineering/project-000/artifact.txt",
|
||||||
|
];
|
||||||
|
|
||||||
|
for key in &expected_keys {
|
||||||
|
client
|
||||||
|
.put_object()
|
||||||
|
.bucket(bucket)
|
||||||
|
.key(key)
|
||||||
|
.body(ByteStream::from_static(b"x"))
|
||||||
|
.send()
|
||||||
|
.await
|
||||||
|
.expect("Failed to put object");
|
||||||
|
}
|
||||||
|
for key in noise_keys {
|
||||||
|
client
|
||||||
|
.put_object()
|
||||||
|
.bucket(bucket)
|
||||||
|
.key(key)
|
||||||
|
.body(ByteStream::from_static(b"x"))
|
||||||
|
.send()
|
||||||
|
.await
|
||||||
|
.expect("Failed to put noise object");
|
||||||
|
}
|
||||||
|
|
||||||
|
let mut listed_keys = Vec::new();
|
||||||
|
let mut continuation_token: Option<String> = None;
|
||||||
|
let mut last_key: Option<String> = None;
|
||||||
|
let mut page_count = 0;
|
||||||
|
|
||||||
|
loop {
|
||||||
|
let mut request = client.list_objects_v2().bucket(bucket).prefix(prefix).max_keys(PAGE_SIZE);
|
||||||
|
|
||||||
|
if let Some(token) = continuation_token.take() {
|
||||||
|
request = request.continuation_token(token);
|
||||||
|
}
|
||||||
|
|
||||||
|
let output = request.send().await.expect("Failed to list objects");
|
||||||
|
|
||||||
|
for obj in output.contents() {
|
||||||
|
if let Some(key) = obj.key() {
|
||||||
|
if let Some(previous) = &last_key {
|
||||||
|
assert!(
|
||||||
|
key > previous.as_str(),
|
||||||
|
"ListObjectsV2 did not preserve lexicographic order: {key} <= {previous}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
last_key = Some(key.to_string());
|
||||||
|
listed_keys.push(key.to_string());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
page_count += 1;
|
||||||
|
|
||||||
|
if output.is_truncated().unwrap_or(false) {
|
||||||
|
continuation_token = output.next_continuation_token().map(|s| s.to_string());
|
||||||
|
assert!(
|
||||||
|
continuation_token.is_some(),
|
||||||
|
"BUG: NextContinuationToken must be present when IsTruncated is true"
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
if page_count > 10 {
|
||||||
|
panic!("Too many pages, possible infinite loop due to pagination bug");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let seen: HashSet<String> = listed_keys.iter().cloned().collect();
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
listed_keys, expected_keys,
|
||||||
|
"Issue #2775 regression: repeated-prefix pagination must return exactly the expected keys in lexicographic order"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
listed_keys.len(),
|
||||||
|
expected_keys.len(),
|
||||||
|
"Issue #2775 regression: expected all {} repeated-prefix objects under {prefix}, got {}",
|
||||||
|
expected_keys.len(),
|
||||||
|
listed_keys.len()
|
||||||
|
);
|
||||||
|
assert_eq!(seen.len(), expected_keys.len(), "Listed keys must be unique");
|
||||||
|
|
||||||
|
for key in &expected_keys {
|
||||||
|
assert!(seen.contains(key), "Missing expected key after repeated-prefix pagination: {key}");
|
||||||
|
}
|
||||||
|
|
||||||
|
env.stop_server();
|
||||||
|
}
|
||||||
|
|
||||||
/// Test that IsTruncated is false when all objects fit within max_keys
|
/// Test that IsTruncated is false when all objects fit within max_keys
|
||||||
///
|
///
|
||||||
/// This is the core bug from issue #1596: the server was returning
|
/// This is the core bug from issue #1596: the server was returning
|
||||||
|
|||||||
@@ -1247,29 +1247,16 @@ impl LocalDisk {
|
|||||||
W: AsyncWrite + Unpin + Send,
|
W: AsyncWrite + Unpin + Send,
|
||||||
{
|
{
|
||||||
let forward = {
|
let forward = {
|
||||||
opts.forward_to.as_ref().filter(|v| v.starts_with(&*current)).map(|v| {
|
opts.forward_to
|
||||||
let forward = v.trim_start_matches(&*current);
|
.as_ref()
|
||||||
if let Some(idx) = forward.find('/') {
|
.and_then(|v| v.strip_prefix(¤t))
|
||||||
forward[..idx].to_owned()
|
.map(|forward| {
|
||||||
} else {
|
if let Some(idx) = forward.find('/') {
|
||||||
forward.to_owned()
|
forward[..idx].to_owned()
|
||||||
}
|
} else {
|
||||||
})
|
forward.to_owned()
|
||||||
// if let Some(forward_to) = &opts.forward_to {
|
}
|
||||||
|
})
|
||||||
// } else {
|
|
||||||
// None
|
|
||||||
// }
|
|
||||||
// if !opts.forward_to.is_empty() && opts.forward_to.starts_with(&*current) {
|
|
||||||
// let forward = opts.forward_to.trim_start_matches(&*current);
|
|
||||||
// if let Some(idx) = forward.find('/') {
|
|
||||||
// &forward[..idx]
|
|
||||||
// } else {
|
|
||||||
// forward
|
|
||||||
// }
|
|
||||||
// } else {
|
|
||||||
// ""
|
|
||||||
// }
|
|
||||||
};
|
};
|
||||||
|
|
||||||
if opts.limit > 0 && *objs_returned >= opts.limit {
|
if opts.limit > 0 && *objs_returned >= opts.limit {
|
||||||
@@ -3238,6 +3225,109 @@ mod test {
|
|||||||
assert_eq!(names.iter().filter(|name| *name == "marker/file.txt").count(), 1);
|
assert_eq!(names.iter().filter(|name| *name == "marker/file.txt").count(), 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_scan_dir_forward_to_repeated_prefix_component() {
|
||||||
|
use rustfs_filemeta::MetacacheReader;
|
||||||
|
use tempfile::tempdir;
|
||||||
|
|
||||||
|
let dir = tempdir().unwrap();
|
||||||
|
let bucket = "test-bucket";
|
||||||
|
let bucket_dir = dir.path().join(bucket);
|
||||||
|
|
||||||
|
for name in [
|
||||||
|
"different/prefix/prefix/repo-0000",
|
||||||
|
"different/prefix/prefix/repo-0001",
|
||||||
|
"different/prefix/prefix/repo-0002",
|
||||||
|
"engineering/alpha-0000",
|
||||||
|
"engineering/engineering/engineering/repo-0000",
|
||||||
|
"engineering/engineering/engineering/repo-0001",
|
||||||
|
"engineering/engineering/repo-0000",
|
||||||
|
"engineering/engineering/repo-0001",
|
||||||
|
"engineering/engineering/repo-0002",
|
||||||
|
"engineering/zulu-0000",
|
||||||
|
"unrelated/engineering/repo-0000",
|
||||||
|
] {
|
||||||
|
let object_dir = bucket_dir.join(name);
|
||||||
|
fs::create_dir_all(&object_dir).await.unwrap();
|
||||||
|
fs::write(object_dir.join(STORAGE_FORMAT_FILE), b"meta").await.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
|
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
|
||||||
|
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
|
||||||
|
|
||||||
|
async fn scan_names(disk: &LocalDisk, bucket: &str, base_dir: &str, forward_to: &str) -> (Vec<String>, i32) {
|
||||||
|
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: true,
|
||||||
|
forward_to: Some(forward_to.to_string()),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let mut objs_returned = 0;
|
||||||
|
|
||||||
|
disk.scan_dir(base_dir.to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
out.close().await.unwrap();
|
||||||
|
drop(out);
|
||||||
|
drop(writer);
|
||||||
|
|
||||||
|
let mut reader = MetacacheReader::new(reader);
|
||||||
|
let entries = reader.read_all().await.unwrap();
|
||||||
|
let names: Vec<String> = entries
|
||||||
|
.into_iter()
|
||||||
|
.filter(|entry| !entry.metadata.is_empty())
|
||||||
|
.map(|entry| entry.name)
|
||||||
|
.collect();
|
||||||
|
|
||||||
|
(names, objs_returned)
|
||||||
|
}
|
||||||
|
|
||||||
|
let (engineering_names, engineering_count) =
|
||||||
|
scan_names(&disk, bucket, "engineering/", "engineering/engineering/engineering/repo-0001").await;
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
engineering_names,
|
||||||
|
vec![
|
||||||
|
"engineering/engineering/engineering/repo-0001".to_string(),
|
||||||
|
"engineering/engineering/repo-0000".to_string(),
|
||||||
|
"engineering/engineering/repo-0001".to_string(),
|
||||||
|
"engineering/engineering/repo-0002".to_string(),
|
||||||
|
"engineering/zulu-0000".to_string(),
|
||||||
|
],
|
||||||
|
"forward_to must resume at the requested triply repeated prefix and preserve lexicographic order"
|
||||||
|
);
|
||||||
|
assert_eq!(engineering_count as usize, engineering_names.len());
|
||||||
|
|
||||||
|
let (different_names, different_count) =
|
||||||
|
scan_names(&disk, bucket, "different/", "different/prefix/prefix/repo-0001").await;
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
different_names,
|
||||||
|
vec![
|
||||||
|
"different/prefix/prefix/repo-0001".to_string(),
|
||||||
|
"different/prefix/prefix/repo-0002".to_string(),
|
||||||
|
],
|
||||||
|
"forward_to must also work for repeated components unrelated to the engineering prefix"
|
||||||
|
);
|
||||||
|
assert_eq!(different_count as usize, different_names.len());
|
||||||
|
|
||||||
|
let (double_names, double_count) = scan_names(&disk, bucket, "engineering/", "engineering/engineering/repo-0001").await;
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
double_names,
|
||||||
|
vec![
|
||||||
|
"engineering/engineering/repo-0001".to_string(),
|
||||||
|
"engineering/engineering/repo-0002".to_string(),
|
||||||
|
"engineering/zulu-0000".to_string(),
|
||||||
|
],
|
||||||
|
"forward_to must not skip a child directory whose name repeats the base prefix"
|
||||||
|
);
|
||||||
|
assert_eq!(double_count as usize, double_names.len());
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_make_volume() {
|
async fn test_make_volume() {
|
||||||
let p = "./testv0";
|
let p = "./testv0";
|
||||||
|
|||||||
Reference in New Issue
Block a user