From 7192eecbfbcc90bd7bcfd857a53c73cb9f3d9d41 Mon Sep 17 00:00:00 2001 From: Hauser Date: Mon, 14 Sep 2026 21:28:31 +0800 Subject: [PATCH] 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 --- .../src/multipart_copy_readiness_test.rs | 134 +++++++++++++++--- 1 file changed, 115 insertions(+), 19 deletions(-) diff --git a/crates/e2e_test/src/multipart_copy_readiness_test.rs b/crates/e2e_test/src/multipart_copy_readiness_test.rs index a668539b7..57be4679f 100644 --- a/crates/e2e_test/src/multipart_copy_readiness_test.rs +++ b/crates/e2e_test/src/multipart_copy_readiness_test.rs @@ -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, +) -> Result<(String, CompletedMultipartUpload), Box> { + 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> { 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?;