mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 04:21:35 +00:00
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 <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -65,6 +65,130 @@ mod tests {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn list_parts_sdk_paginator_terminates_and_preserves_completion() -> Result<(), Box<dyn std::error::Error + Send + Sync>>
|
||||
{
|
||||
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<dyn std::error::Error + Send + Sync>> {
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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::<Vec<_>>(), 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<Option<Uuid>> {
|
||||
let mut epochs = Vec::with_capacity(disks.len());
|
||||
let deadline = tokio::time::Instant::now() + Duration::from_secs(30);
|
||||
|
||||
@@ -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("<IsTruncated>false</IsTruncated>"));
|
||||
assert!(xml.contains(&format!("<PartNumberMarker>{marker}</PartNumberMarker>")));
|
||||
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("<IsTruncated>true</IsTruncated>"));
|
||||
assert!(xml.contains("<NextPartNumberMarker>3</NextPartNumberMarker>"));
|
||||
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 {
|
||||
|
||||
Reference in New Issue
Block a user