Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 18 additions & 5 deletions crates/cli/src/commands/cp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -618,6 +618,7 @@ fn requested_metadata_directive(args: &CpArgs) -> Option<MetadataDirective> {
fn transfer_copy_options(
args: &CpArgs,
source_version_id: Option<String>,
source_etag: Option<String>,
encryption: Option<&ObjectEncryptionRequest>,
) -> rc_core::Result<TransferCopyOptions> {
let metadata_directive = requested_metadata_directive(args);
Expand Down Expand Up @@ -645,6 +646,7 @@ fn transfer_copy_options(
customer_key: args.source_customer_key.clone(),
..TransferReadOptions::default()
},
source_etag,
metadata_directive,
tagging_directive,
destination,
Expand Down Expand Up @@ -697,7 +699,7 @@ fn validate_fidelity_directions(
)
});
if any_remote && target_remote {
let copy_options = transfer_copy_options(args, None, None)?;
let copy_options = transfer_copy_options(args, None, None, None)?;
if same_alias_remote_copy
&& matches!(
copy_options.metadata_directive,
Expand Down Expand Up @@ -1303,7 +1305,7 @@ async fn perform_planned_remote_copy(
)));
}
let options = multipart_options_from_source(&current)?;
let transfer = transfer_copy_options(args, current.version_id.clone(), encryption)?;
let transfer = transfer_copy_options(args, current.version_id.clone(), None, encryption)?;
let copied = source_client
.multipart_copy_with_transfer_options(
source,
Expand All @@ -1323,7 +1325,12 @@ async fn perform_planned_remote_copy(
object: copied.object,
});
}
let options = transfer_copy_options(args, source_info.version_id.clone(), encryption)?;
let options = transfer_copy_options(
args,
source_info.version_id.clone(),
source_info.etag.clone(),
encryption,
)?;
let copied = source_client
.copy_object_with_transfer_options(source, target, &options)
.await?;
Expand Down Expand Up @@ -4233,11 +4240,17 @@ mod tests {
args.fidelity.retain_until = Some("2099-01-02T03:04:05Z".to_string());
args.fidelity.legal_hold = Some("ON".to_string());

let options =
transfer_copy_options(&args, Some("source-v1".to_string()), None).expect("copy policy");
let options = transfer_copy_options(
&args,
Some("source-v1".to_string()),
Some("source-etag".to_string()),
None,
)
.expect("copy policy");
assert_eq!(options.metadata_directive, Some(MetadataDirective::Copy));
assert!(options.destination.attributes.is_none());
assert_eq!(options.source.version_id.as_deref(), Some("source-v1"));
assert_eq!(options.source_etag.as_deref(), Some("source-etag"));
assert!(options.destination.retention.is_some());
assert_eq!(
options.destination.legal_hold,
Expand Down
6 changes: 4 additions & 2 deletions crates/cli/src/commands/mv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -406,8 +406,10 @@ async fn copy_for_move(
}
None => {
let source_info = source_client.head_object(source).await?;
let copy_options =
CopyObjectOptions::for_source_version(source_info.version_id.clone())?;
let copy_options = CopyObjectOptions::for_source_identity(
source_info.version_id.clone(),
source_info.etag.clone(),
)?;
let object = source_client
.copy_object_with_options(source, target, &copy_options, encryption)
.await?;
Expand Down
1 change: 1 addition & 0 deletions crates/cli/src/commands/undo.rs
Original file line number Diff line number Diff line change
Expand Up @@ -361,6 +361,7 @@ async fn execute_plan_item<S: UndoStore>(
UndoAction::RestoreVersion { source_version_id } => {
let options = CopyObjectOptions {
source_version_id: Some(source_version_id.clone()),
source_etag: None,
};
store.copy_version(&path, &options).await.and_then(|info| {
let created_version_id = info.version_id.ok_or_else(|| {
Expand Down
4 changes: 4 additions & 0 deletions crates/cli/tests/recursive_remote_copy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1686,6 +1686,10 @@ fn recursive_copy_executes_the_planned_server_side_copy() {
copies[0].headers.get("x-amz-copy-source"),
Some(&"source/src/a.txt".to_string())
);
assert_eq!(
copies[0].headers.get("x-amz-copy-source-if-match"),
Some(&"\"source-etag\"".to_string())
);
}

#[test]
Expand Down
25 changes: 21 additions & 4 deletions crates/core/src/traits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -135,17 +135,35 @@ pub struct ObjectReadOptions {
pub struct CopyObjectOptions {
/// Exact historical source version to copy. `None` selects the current source.
pub source_version_id: Option<String>,
/// Copy only when the selected source still has this ETag.
pub source_etag: Option<String>,
}

impl CopyObjectOptions {
/// Build copy options while rejecting an ambiguous empty source version identifier.
pub fn for_source_version(source_version_id: Option<String>) -> Result<Self> {
Self::for_source_identity(source_version_id, None)
}

/// Build copy options pinned to an exact source version or ETag.
pub fn for_source_identity(
source_version_id: Option<String>,
source_etag: Option<String>,
) -> Result<Self> {
if source_version_id.as_deref().is_some_and(str::is_empty) {
return Err(Error::InvalidPath(
"Source version ID cannot be empty".to_string(),
));
}
Ok(Self { source_version_id })
if source_etag.as_deref().is_some_and(str::is_empty) {
return Err(Error::InvalidPath(
"Source ETag cannot be empty".to_string(),
));
}
Ok(Self {
source_version_id,
source_etag,
})
}
}

Expand Down Expand Up @@ -793,10 +811,9 @@ pub trait ObjectStore: Send + Sync {
options: &CopyObjectOptions,
encryption: Option<&ObjectEncryptionRequest>,
) -> Result<ObjectInfo> {
if options.source_version_id.is_some() {
if options.source_version_id.is_some() || options.source_etag.is_some() {
return Err(Error::UnsupportedFeature(
"Exact-version server-side copy is not implemented by this object store"
.to_string(),
"Exact-source server-side copy is not implemented by this object store".to_string(),
));
}
self.copy_object(src, dst, encryption).await
Expand Down
10 changes: 9 additions & 1 deletion crates/core/src/transfer_options.rs
Original file line number Diff line number Diff line change
Expand Up @@ -474,6 +474,8 @@ impl From<ObjectReadOptions> for TransferReadOptions {
pub struct TransferCopyOptions {
/// Source version, checksum, and SSE-C selection.
pub source: TransferReadOptions,
/// Copy only when the selected source still has this ETag.
pub source_etag: Option<String>,
/// Source metadata handling requested for the destination.
pub metadata_directive: Option<MetadataDirective>,
/// Source tag handling requested for the destination.
Expand All @@ -486,6 +488,11 @@ impl TransferCopyOptions {
/// Validate directives and their replacement payloads before mutation.
pub fn validate(&self) -> Result<()> {
self.source.validate()?;
if self.source_etag.as_deref().is_some_and(str::is_empty) {
return Err(Error::InvalidPath(
"Source ETag cannot be empty".to_string(),
));
}
self.destination.validate()?;
match self.metadata_directive {
None | Some(MetadataDirective::Copy) if self.destination.attributes.is_some() => {
Expand Down Expand Up @@ -529,7 +536,8 @@ impl TransferCopyOptions {
}
let source = self.source.legacy_read_options()?;
let encryption = self.destination.legacy_copy_encryption()?;
let copy = CopyObjectOptions::for_source_version(source.version_id)?;
let copy =
CopyObjectOptions::for_source_identity(source.version_id, self.source_etag.clone())?;
Ok((copy, encryption))
}

Expand Down
3 changes: 2 additions & 1 deletion crates/core/tests/undo_planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -173,8 +173,9 @@ fn overwrite_after_delete_marker_is_refused_instead_of_guessing() {
}

#[test]
fn copy_options_reject_empty_source_version_ids() {
fn copy_options_reject_empty_source_identities() {
assert!(CopyObjectOptions::for_source_version(Some(String::new())).is_err());
assert!(CopyObjectOptions::for_source_identity(None, Some(String::new())).is_err());
assert_eq!(
CopyObjectOptions::for_source_version(Some("data-v1".to_string()))
.expect("valid source version")
Expand Down
48 changes: 43 additions & 5 deletions crates/s3/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5501,17 +5501,18 @@ impl ObjectStore for S3Client {
encryption: Option<&ObjectEncryptionRequest>,
) -> Result<ObjectInfo> {
let copy_source = encoded_copy_source(src, options.source_version_id.as_deref());
let response = apply_object_encryption_to_copy_request(
let mut request = apply_object_encryption_to_copy_request(
self.inner
.copy_object()
.copy_source(&copy_source)
.bucket(&dst.bucket)
.key(&dst.key),
encryption,
)
.send()
.await
.map_err(|error| {
);
if let Some(source_etag) = options.source_etag.as_deref() {
request = request.copy_source_if_match(quoted_etag(source_etag));
}
let response = request.send().await.map_err(|error| {
Self::map_object_request_error(&error, src, options.source_version_id.as_deref())
})?;

Expand Down Expand Up @@ -5555,6 +5556,9 @@ impl ObjectStore for S3Client {
if matches!(options.metadata_directive, Some(MetadataDirective::Copy)) {
request = request.metadata_directive(aws_sdk_s3::types::MetadataDirective::Copy);
}
if let Some(source_etag) = options.source_etag.as_deref() {
request = request.copy_source_if_match(quoted_etag(source_etag));
}
request = apply_object_lock_to_copy_request(request, &options.destination)?;
let response = request.send().await.map_err(|error| {
if options.destination.retention.is_some() || options.destination.legal_hold.is_some() {
Expand Down Expand Up @@ -11389,6 +11393,35 @@ mod tests {
);
}

#[tokio::test]
async fn copy_object_with_options_conditions_on_source_etag() {
let response = http::Response::builder()
.status(412)
.header("x-amz-error-code", "PreconditionFailed")
.body(SdkBody::from(
"<Error><Code>PreconditionFailed</Code><Message>source changed</Message></Error>",
))
.expect("build precondition response");
let (client, request_receiver) = test_s3_client(Some(response));
let src = RemotePath::new("test", "source-bucket", "src.txt");
let dst = RemotePath::new("test", "destination-bucket", "dst.txt");
let options =
CopyObjectOptions::for_source_identity(None, Some("planned-source-etag".to_string()))
.expect("valid source identity");

let error = client
.copy_object_with_options(&src, &dst, &options, None)
.await
.expect_err("a changed source must block the copy");

let request = request_receiver.expect_request();
assert_eq!(
request.headers().get("x-amz-copy-source-if-match"),
Some("\"planned-source-etag\"")
);
assert!(matches!(error, Error::Conflict(_)));
}

#[tokio::test]
async fn transfer_copy_sends_explicit_metadata_copy_directive() {
let response = http::Response::builder()
Expand All @@ -11407,6 +11440,7 @@ mod tests {
&src,
&dst,
&TransferCopyOptions {
source_etag: Some("planned-source-etag".to_string()),
metadata_directive: Some(MetadataDirective::Copy),
destination: ObjectWriteOptions {
storage_class: Some("STANDARD".to_string()),
Expand All @@ -11426,6 +11460,10 @@ mod tests {
request.headers().get("x-amz-storage-class"),
Some("STANDARD")
);
assert_eq!(
request.headers().get("x-amz-copy-source-if-match"),
Some("\"planned-source-etag\"")
);
}

#[tokio::test]
Expand Down