From 673031eea143a56c49a4d72bbed442e0416e9f7a Mon Sep 17 00:00:00 2001 From: Hauser Date: Fri, 2 Oct 2026 15:34:07 +0800 Subject: [PATCH] fix(s3): honor sparse ListParts markers and verify termination Resume part listings at the first part above the numeric marker, even when that marker is absent. Use binary search over the sorted part numbers and retain the existing exact-tail empty-slice path. Add storage, XML, and real AWS SDK paginator regressions for empty and terminal pages, sparse markers, and multipart completion integrity. Co-Authored-By: heihutu Co-Authored-By: zhi22915 --- .../src/multipart_storage_class_test.rs | 124 ++++++++++++++++++ crates/ecstore/src/set_disk/mod.rs | 26 +++- crates/ecstore/src/set_disk/ops/multipart.rs | 48 +++++++ rustfs/src/storage/s3_api/multipart.rs | 49 +++++++ 4 files changed, 243 insertions(+), 4 deletions(-) diff --git a/crates/e2e_test/src/multipart_storage_class_test.rs b/crates/e2e_test/src/multipart_storage_class_test.rs index 21a5e52d0..413bfc394 100644 --- a/crates/e2e_test/src/multipart_storage_class_test.rs +++ b/crates/e2e_test/src/multipart_storage_class_test.rs @@ -65,6 +65,130 @@ mod tests { Ok(()) } + #[tokio::test] + async fn list_parts_sdk_paginator_terminates_and_preserves_completion() -> Result<(), Box> + { + init_logging(); + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(Vec::new()).await?; + let client = env.create_s3_client(); + let bucket = "multipart-pagination"; + let key = "sparse-parts.bin"; + env.create_test_bucket(bucket).await?; + let upload = client.create_multipart_upload().bucket(bucket).key(key).send().await?; + let upload_id = upload.upload_id().ok_or("CreateMultipartUpload returned no upload ID")?; + + let mut empty_pages = client + .list_parts() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .into_paginator() + .stop_on_duplicate_token(false) + .send(); + let empty = tokio::time::timeout(std::time::Duration::from_secs(30), empty_pages.next()) + .await? + .ok_or("empty upload must return one page")??; + assert!(empty.parts().is_empty()); + assert_eq!(empty.is_truncated(), Some(false)); + assert_eq!(empty.next_part_number_marker(), None); + assert!( + tokio::time::timeout(std::time::Duration::from_secs(30), empty_pages.next()) + .await? + .is_none() + ); + + let mut expected_body = Vec::new(); + let mut completed_parts = Vec::new(); + for (part_number, body) in [(1, vec![b'a'; PART_SIZE]), (3, vec![b'b'; PART_SIZE]), (10, b"tail".to_vec())] { + expected_body.extend_from_slice(&body); + let uploaded = client + .upload_part() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .part_number(part_number) + .body(ByteStream::from(body)) + .send() + .await?; + completed_parts.push( + CompletedPart::builder() + .part_number(part_number) + .set_e_tag(uploaded.e_tag().map(str::to_owned)) + .build(), + ); + } + + for (page_size, expected_pages) in [(1, 3), (2, 2), (3, 1), (1000, 1)] { + let mut pages = client + .list_parts() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .into_paginator() + .page_size(page_size) + // The server must terminate pagination even without the SDK's duplicate-token safeguard. + .stop_on_duplicate_token(false) + .send(); + let mut page_count = 0; + let mut listed_parts = Vec::new(); + while let Some(page) = tokio::time::timeout(std::time::Duration::from_secs(30), pages.next()).await? { + assert!(page_count < expected_pages, "paginator must terminate after the final page"); + let page = page?; + page_count += 1; + let truncated = page_count < expected_pages; + assert_eq!(page.is_truncated(), Some(truncated)); + let last_part = page.parts().last().and_then(|part| part.part_number()); + let expected_marker = truncated.then(|| last_part.expect("truncated page must contain parts").to_string()); + assert_eq!(page.next_part_number_marker(), expected_marker.as_deref()); + listed_parts.extend( + page.parts() + .iter() + .map(|part| part.part_number().expect("part number must be present")), + ); + } + assert_eq!(page_count, expected_pages); + assert_eq!(listed_parts, [1, 3, 10]); + } + + let mut resumed_pages = client + .list_parts() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .part_number_marker("2") + .into_paginator() + .page_size(1) + .stop_on_duplicate_token(false) + .send(); + for (part_number, next_marker) in [(3, Some("3")), (10, None)] { + let page = tokio::time::timeout(std::time::Duration::from_secs(30), resumed_pages.next()) + .await? + .ok_or("parts above a missing marker must still be listed")??; + assert_eq!(page.parts().len(), 1); + assert_eq!(page.parts()[0].part_number(), Some(part_number)); + assert_eq!(page.is_truncated(), Some(next_marker.is_some())); + assert_eq!(page.next_part_number_marker(), next_marker); + } + assert!( + tokio::time::timeout(std::time::Duration::from_secs(30), resumed_pages.next()) + .await? + .is_none() + ); + + client + .complete_multipart_upload() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .multipart_upload(CompletedMultipartUpload::builder().set_parts(Some(completed_parts)).build()) + .send() + .await?; + assert_completed_object(&client, bucket, key, "STANDARD", &expected_body).await?; + env.stop_server(); + Ok(()) + } + #[tokio::test] async fn multipart_upload_preserves_standard_and_rrs_across_retry_and_resume() -> Result<(), Box> { diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index bc6daedf6..59a71adc3 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -7171,10 +7171,10 @@ fn parts_after_marker(part_numbers: &[usize], part_number_marker: usize) -> Opti return Some(part_numbers); } - part_numbers - .iter() - .position(|&part_number| part_number != 0 && part_number == part_number_marker) - .map(|index| &part_numbers[index + 1..]) + // reduce_quorum_part_numbers returns sorted numbers; the marker need not exist. + part_numbers.last().filter(|&&last| part_number_marker <= last)?; + let index = part_numbers.partition_point(|&part_number| part_number <= part_number_marker); + Some(&part_numbers[index..]) } pub fn canonicalize_etag(etag: &str) -> String { @@ -13972,6 +13972,24 @@ mod tests { assert!(parts_after_marker(&part_numbers, 4).is_none()); } + #[test] + fn parts_after_marker_uses_exclusive_numeric_boundary_for_sparse_parts() { + let part_numbers = [1, 3, 10]; + for (marker, expected) in [ + (0, Some(&part_numbers[..])), + (1, Some(&part_numbers[1..])), + (2, Some(&part_numbers[1..])), + (3, Some(&part_numbers[2..])), + (9, Some(&part_numbers[2..])), + (10, Some(&part_numbers[3..])), + (11, None), + (usize::MAX, None), + ] { + assert_eq!(parts_after_marker(&part_numbers, marker), expected, "marker {marker}"); + } + assert_eq!(parts_after_marker(&[], 1), None); + } + #[test] fn delete_file_info_version_id_maps_explicit_null_version_to_stored_null() { assert_eq!(delete_file_info_version_id(Some(Uuid::nil())), None); diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 61e62ea59..b84458619 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -4505,6 +4505,54 @@ mod tests { } } + #[tokio::test] + async fn list_object_parts_pagination_preserves_sparse_markers_and_terminates() { + let (_temp_dirs, disks, set_disks) = hermetic_set_disks(4).await; + let bucket = "multipart-pagination"; + let object = "object"; + make_bucket_on_all(&disks, bucket).await; + let opts = ObjectOptions::default(); + let upload = set_disks + .new_multipart_upload(bucket, object, &opts) + .await + .expect("multipart upload should be created"); + + let empty = set_disks + .list_object_parts(bucket, object, &upload.upload_id, None, 2, &opts) + .await + .expect("empty upload should be listable"); + assert!(empty.parts.is_empty()); + assert!(!empty.is_truncated); + assert_eq!(empty.next_part_number_marker, None); + + for part_number in [10, 1, 3] { + put_test_part(&set_disks, bucket, object, &upload.upload_id, part_number, b"part", 4).await; + } + + for (marker, max_parts, expected_parts, next_marker) in [ + (None, 0, vec![], None), + (None, 1, vec![1], Some(1)), + (None, 2, vec![1, 3], Some(3)), + (None, 3, vec![1, 3, 10], None), + (None, MAX_PARTS_COUNT + 1, vec![1, 3, 10], None), + (Some(1), 2, vec![3, 10], None), + (Some(2), 1, vec![3], Some(3)), + (Some(3), 2, vec![10], None), + (Some(10), 2, vec![], None), + (Some(11), 2, vec![], None), + ] { + let page = set_disks + .list_object_parts(bucket, object, &upload.upload_id, marker, max_parts, &opts) + .await + .expect("uploaded parts should be listable"); + assert_eq!(page.part_number_marker, marker.unwrap_or_default()); + assert_eq!(page.max_parts, max_parts.min(MAX_PARTS_COUNT)); + assert_eq!(page.parts.iter().map(|part| part.part_num).collect::>(), expected_parts); + assert_eq!(page.is_truncated, next_marker.is_some()); + assert_eq!(page.next_part_number_marker, next_marker); + } + } + async fn object_transaction_epochs(disks: &[DiskStore], bucket: &str, object: &str) -> Vec> { let mut epochs = Vec::with_capacity(disks.len()); let deadline = tokio::time::Instant::now() + Duration::from_secs(30); diff --git a/rustfs/src/storage/s3_api/multipart.rs b/rustfs/src/storage/s3_api/multipart.rs index 8a336f77a..1124e31f6 100644 --- a/rustfs/src/storage/s3_api/multipart.rs +++ b/rustfs/src/storage/s3_api/multipart.rs @@ -259,6 +259,55 @@ mod tests { assert_eq!(output.initiator, Some(rustfs_initiator())); } + #[test] + fn test_list_parts_output_xml_omits_terminal_marker() { + for (marker, part_numbers) in [(0, vec![]), (3, vec![10]), (0, vec![1, 3])] { + let output = build_list_parts_output(ListPartsInfo { + part_number_marker: marker, + max_parts: 2, + parts: part_numbers + .into_iter() + .map(|part_num| PartInfo { + part_num, + ..Default::default() + }) + .collect(), + ..Default::default() + }); + + assert_eq!(output.is_truncated, Some(false)); + assert_eq!(output.next_part_number_marker, None); + let xml = String::from_utf8(crate::storage::storage_api::serialize(&output).expect("output should serialize")) + .expect("XML should be UTF-8"); + assert!(xml.contains("false")); + assert!(xml.contains(&format!("{marker}"))); + assert!(!xml.contains("NextPartNumberMarker"), "terminal page must omit the token: {xml}"); + } + } + + #[test] + fn test_list_parts_output_xml_preserves_truncated_marker() { + let output = build_list_parts_output(ListPartsInfo { + is_truncated: true, + next_part_number_marker: Some(3), + max_parts: 2, + parts: [1, 3] + .into_iter() + .map(|part_num| PartInfo { + part_num, + ..Default::default() + }) + .collect(), + ..Default::default() + }); + + let xml = String::from_utf8(crate::storage::storage_api::serialize(&output).expect("output should serialize")) + .expect("XML should be UTF-8"); + assert!(xml.contains("true")); + assert!(xml.contains("3")); + assert_eq!(output.parts.as_ref().expect("parts should be present").len(), 2); + } + #[test] fn test_list_parts_output_reports_logical_size_for_compressed_parts() { let input = ListPartsInfo {