mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-03 20:20:29 +00:00
test: cover overlapped multipart staging copy (#7865)
Add Harbor-style E2E coverage for a CopyObject racing with CompleteMultipartUpload. The test allows the copy to observe an unpublished source as NoSuchKey, but rejects 5xx leakage and proves retrying the completed source remains readable. Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -19,6 +19,8 @@ use aws_sdk_s3::error::ProvideErrorMetadata;
|
||||
use aws_sdk_s3::primitives::ByteStream;
|
||||
use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart, StorageClass};
|
||||
use std::error::Error;
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::Barrier;
|
||||
use tracing::info;
|
||||
|
||||
const BUCKET: &str = "multipart-copy-readiness";
|
||||
@@ -26,6 +28,9 @@ const SOURCE_KEY: &str = "docker/registry/v2/repositories/example/_uploads/uploa
|
||||
const TARGET_KEY: &str = "docker/registry/v2/blobs/sha256/c0/digest/data";
|
||||
const UNFINISHED_SOURCE_KEY: &str = "docker/registry/v2/repositories/example/_uploads/upload-id-unfinished/data";
|
||||
const UNFINISHED_TARGET_KEY: &str = "docker/registry/v2/blobs/sha256/c1/digest/data";
|
||||
const OVERLAP_SOURCE_KEY: &str = "docker/registry/v2/repositories/example/_uploads/upload-id-overlap/data";
|
||||
const OVERLAP_TARGET_KEY: &str = "docker/registry/v2/blobs/sha256/c2/digest/data";
|
||||
const OVERLAP_RETRY_TARGET_KEY: &str = "docker/registry/v2/blobs/sha256/c3/digest/data";
|
||||
|
||||
fn list_contains_key(output: &aws_sdk_s3::operation::list_objects_v2::ListObjectsV2Output, key: &str) -> bool {
|
||||
output
|
||||
@@ -34,6 +39,35 @@ fn list_contains_key(output: &aws_sdk_s3::operation::list_objects_v2::ListObject
|
||||
.any(|object| object.key().is_some_and(|candidate| candidate == key))
|
||||
}
|
||||
|
||||
async fn upload_one_part_mpu(
|
||||
client: &aws_sdk_s3::Client,
|
||||
bucket: &str,
|
||||
key: &str,
|
||||
payload: Vec<u8>,
|
||||
) -> Result<(String, CompletedMultipartUpload), Box<dyn Error + Send + Sync>> {
|
||||
let create = client.create_multipart_upload().bucket(bucket).key(key).send().await?;
|
||||
let upload_id = create.upload_id().ok_or("missing upload id")?.to_string();
|
||||
let part = client
|
||||
.upload_part()
|
||||
.bucket(bucket)
|
||||
.key(key)
|
||||
.upload_id(&upload_id)
|
||||
.part_number(1)
|
||||
.body(ByteStream::from(payload))
|
||||
.send()
|
||||
.await?;
|
||||
let completed = CompletedMultipartUpload::builder()
|
||||
.parts(
|
||||
CompletedPart::builder()
|
||||
.part_number(1)
|
||||
.set_e_tag(part.e_tag().map(str::to_string))
|
||||
.build(),
|
||||
)
|
||||
.build();
|
||||
|
||||
Ok((upload_id, completed))
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn harbor_style_multipart_staging_copy_object_boundaries() -> Result<(), Box<dyn Error + Send + Sync>> {
|
||||
init_logging();
|
||||
@@ -46,25 +80,7 @@ async fn harbor_style_multipart_staging_copy_object_boundaries() -> Result<(), B
|
||||
env.create_test_bucket(BUCKET).await?;
|
||||
|
||||
let payload = vec![0xAB; 273];
|
||||
let create = client.create_multipart_upload().bucket(BUCKET).key(SOURCE_KEY).send().await?;
|
||||
let upload_id = create.upload_id().ok_or("missing upload id")?.to_string();
|
||||
let part = client
|
||||
.upload_part()
|
||||
.bucket(BUCKET)
|
||||
.key(SOURCE_KEY)
|
||||
.upload_id(&upload_id)
|
||||
.part_number(1)
|
||||
.body(ByteStream::from(payload.clone()))
|
||||
.send()
|
||||
.await?;
|
||||
let completed = CompletedMultipartUpload::builder()
|
||||
.parts(
|
||||
CompletedPart::builder()
|
||||
.part_number(1)
|
||||
.set_e_tag(part.e_tag().map(str::to_string))
|
||||
.build(),
|
||||
)
|
||||
.build();
|
||||
let (upload_id, completed) = upload_one_part_mpu(&client, BUCKET, SOURCE_KEY, payload.clone()).await?;
|
||||
client
|
||||
.complete_multipart_upload()
|
||||
.bucket(BUCKET)
|
||||
@@ -99,6 +115,78 @@ async fn harbor_style_multipart_staging_copy_object_boundaries() -> Result<(), B
|
||||
let copied_body = copied.body.collect().await?.into_bytes();
|
||||
assert_eq!(copied_body.as_ref(), payload.as_slice());
|
||||
|
||||
let overlap_payload = vec![0xBC; 15_796];
|
||||
let (overlap_upload_id, overlap_completed) =
|
||||
upload_one_part_mpu(&client, BUCKET, OVERLAP_SOURCE_KEY, overlap_payload.clone()).await?;
|
||||
|
||||
let barrier = Arc::new(Barrier::new(2));
|
||||
let complete_client = client.clone();
|
||||
let complete_barrier = Arc::clone(&barrier);
|
||||
let complete_upload_id = overlap_upload_id.clone();
|
||||
let complete_task = tokio::spawn(async move {
|
||||
complete_barrier.wait().await;
|
||||
complete_client
|
||||
.complete_multipart_upload()
|
||||
.bucket(BUCKET)
|
||||
.key(OVERLAP_SOURCE_KEY)
|
||||
.upload_id(complete_upload_id)
|
||||
.multipart_upload(overlap_completed)
|
||||
.send()
|
||||
.await
|
||||
});
|
||||
|
||||
let copy_client = client.clone();
|
||||
let copy_barrier = Arc::clone(&barrier);
|
||||
let copy_task = tokio::spawn(async move {
|
||||
copy_barrier.wait().await;
|
||||
copy_client
|
||||
.copy_object()
|
||||
.bucket(BUCKET)
|
||||
.key(OVERLAP_TARGET_KEY)
|
||||
.copy_source(format!("/{BUCKET}/{OVERLAP_SOURCE_KEY}"))
|
||||
.storage_class(StorageClass::Standard)
|
||||
.send()
|
||||
.await
|
||||
});
|
||||
|
||||
complete_task.await??;
|
||||
match copy_task.await? {
|
||||
Ok(_) => {}
|
||||
Err(copy_err) => {
|
||||
assert_eq!(
|
||||
copy_err.raw_response().map(|response| response.status().as_u16()),
|
||||
Some(404),
|
||||
"overlapped CopyObject may race before publication, but must not leak a 5xx response: {copy_err:?}"
|
||||
);
|
||||
assert_eq!(
|
||||
copy_err.as_service_error().and_then(ProvideErrorMetadata::code),
|
||||
Some("NoSuchKey"),
|
||||
"overlapped CopyObject that wins before Complete must look like an unpublished ordinary object: {copy_err:?}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
client
|
||||
.copy_object()
|
||||
.bucket(BUCKET)
|
||||
.key(OVERLAP_RETRY_TARGET_KEY)
|
||||
.copy_source(format!("/{BUCKET}/{OVERLAP_SOURCE_KEY}"))
|
||||
.storage_class(StorageClass::Standard)
|
||||
.send()
|
||||
.await?;
|
||||
let overlap_copied = client
|
||||
.get_object()
|
||||
.bucket(BUCKET)
|
||||
.key(OVERLAP_RETRY_TARGET_KEY)
|
||||
.send()
|
||||
.await?;
|
||||
let overlap_copied_body = overlap_copied.body.collect().await?.into_bytes();
|
||||
assert_eq!(
|
||||
overlap_copied_body.as_ref(),
|
||||
overlap_payload.as_slice(),
|
||||
"a completed staging object must not become permanently unreadable after an overlapped copy attempt"
|
||||
);
|
||||
|
||||
let unfinished = client
|
||||
.create_multipart_upload()
|
||||
.bucket(BUCKET)
|
||||
@@ -155,6 +243,14 @@ async fn harbor_style_multipart_staging_copy_object_boundaries() -> Result<(), B
|
||||
.upload_id(unfinished_upload_id)
|
||||
.send()
|
||||
.await?;
|
||||
client
|
||||
.delete_object()
|
||||
.bucket(BUCKET)
|
||||
.key(OVERLAP_RETRY_TARGET_KEY)
|
||||
.send()
|
||||
.await?;
|
||||
let _ = client.delete_object().bucket(BUCKET).key(OVERLAP_TARGET_KEY).send().await;
|
||||
client.delete_object().bucket(BUCKET).key(OVERLAP_SOURCE_KEY).send().await?;
|
||||
client.delete_object().bucket(BUCKET).key(TARGET_KEY).send().await?;
|
||||
client.delete_object().bucket(BUCKET).key(SOURCE_KEY).send().await?;
|
||||
env.delete_test_bucket(BUCKET).await?;
|
||||
|
||||
Reference in New Issue
Block a user