From 8c49671161d90aa43af9ad6cda12ebd62bce870c Mon Sep 17 00:00:00 2001 From: "Sergei Z." <47861714+SamuraJey@users.noreply.github.com> Date: Wed, 13 May 2026 19:36:23 +0500 Subject: [PATCH] Fix #2775 recursive list handling in LocalDisk::scan_dir() (#2923) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: houseme Co-authored-by: 安正超 --- .../src/list_objects_v2_pagination_test.rs | 125 ++++++++++++++++ crates/ecstore/src/disk/local.rs | 136 +++++++++++++++--- 2 files changed, 238 insertions(+), 23 deletions(-) diff --git a/crates/e2e_test/src/list_objects_v2_pagination_test.rs b/crates/e2e_test/src/list_objects_v2_pagination_test.rs index 339e92cf1..624d21f67 100644 --- a/crates/e2e_test/src/list_objects_v2_pagination_test.rs +++ b/crates/e2e_test/src/list_objects_v2_pagination_test.rs @@ -31,6 +31,7 @@ mod tests { use aws_sdk_s3::Client; use aws_sdk_s3::primitives::ByteStream; use serial_test::serial; + use std::collections::HashSet; use tracing::info; /// 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 = None; + let mut last_key: Option = 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 = 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 /// /// This is the core bug from issue #1596: the server was returning diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 34884e6c9..61d0c18d8 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -1247,29 +1247,16 @@ impl LocalDisk { W: AsyncWrite + Unpin + Send, { let forward = { - opts.forward_to.as_ref().filter(|v| v.starts_with(&*current)).map(|v| { - let forward = v.trim_start_matches(&*current); - if let Some(idx) = forward.find('/') { - 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 { - // "" - // } + opts.forward_to + .as_ref() + .and_then(|v| v.strip_prefix(¤t)) + .map(|forward| { + if let Some(idx) = forward.find('/') { + forward[..idx].to_owned() + } else { + forward.to_owned() + } + }) }; 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); } + #[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, 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 = 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] async fn test_make_volume() { let p = "./testv0";