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 {