diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index 480a8af3e..5a4d2d220 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -16,6 +16,7 @@ use crate::common::{ RustFSTestEnvironment, awscurl_available, awscurl_post_sts_form_urlencoded, init_logging, local_http_client, }; use aws_sdk_s3::config::{Credentials, Region}; +use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::primitives::ByteStream; use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration}; use aws_sdk_s3::{Client, Config}; @@ -701,7 +702,7 @@ async fn wait_for_object_on_target( return Ok(body); } Err(err) => { - if err.to_string().contains("NoSuchKey") || err.to_string().contains("NotFound") { + if matches!(err.code(), Some("NoSuchKey" | "NotFound" | "NoSuchVersion")) { sleep(Duration::from_millis(250)).await; continue; } diff --git a/crates/ecstore/src/bucket/bucket_target_sys.rs b/crates/ecstore/src/bucket/bucket_target_sys.rs index 707acb33b..5b1c53aeb 100644 --- a/crates/ecstore/src/bucket/bucket_target_sys.rs +++ b/crates/ecstore/src/bucket/bucket_target_sys.rs @@ -1343,7 +1343,30 @@ impl TargetClient { .await { Ok(_) => Ok(()), - Err(e) => Err(e.into()), + Err(e) => match e { + SdkError::ServiceError(service_err) => { + let err = service_err.into_err(); + let meta = err.meta(); + let error = match (meta.code(), meta.message()) { + (Some(code), Some(message)) => format!("put_object failed: {code}: {message}"), + (Some(code), None) => format!("put_object failed: {code}"), + (None, Some(message)) => format!("put_object failed: {message}"), + (None, None) => format!("put_object failed: {err:?}"), + }; + Err(S3ClientError::with_metadata( + error, + None, + meta.code().map(ToOwned::to_owned), + meta.message().map(ToOwned::to_owned), + )) + } + SdkError::DispatchFailure(dispatch_err) => Err(S3ClientError::new(format!( + "put_object dispatch failure for bucket:{bucket} object:{object}: {dispatch_err:?}" + ))), + other => Err(S3ClientError::new(format!( + "put_object request failed for bucket:{bucket} object:{object}: {other:?}" + ))), + }, } } diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index a8fcd6efb..753027644 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -2734,6 +2734,10 @@ fn site_replication_bucket_target_for_peer( let port = parsed.port_or_known_default().ok_or_else(|| { S3Error::with_message(S3ErrorCode::InvalidRequest, format!("peer endpoint missing port: {}", peer.endpoint)) })?; + let region = get_global_region() + .map(|region| region.to_string()) + .filter(|region| !region.is_empty()) + .unwrap_or_else(|| "us-east-1".to_string()); let arn = arn_override.unwrap_or_else(|| { ARN::new( BucketTargetType::ReplicationService, @@ -2756,6 +2760,7 @@ fn site_replication_bucket_target_for_peer( target_bucket: bucket.to_string(), secure: parsed.scheme().eq_ignore_ascii_case("https"), arn, + region, target_type: BucketTargetType::ReplicationService, deployment_id: peer.deployment_id.clone(), ..Default::default() @@ -3464,6 +3469,31 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> { ) .await?; } + + if item.r#type == "version-config" + && metadata_sys::get_versioning_config(&item.bucket) + .await + .ok() + .is_some_and(|(config, _)| config.enabled()) + && let Some(runtime) = runtime_site_replication_targets().await? + { + ensure_site_replication_bucket_targets( + &item.bucket, + &runtime.state, + &runtime.local_peer, + replication_config.as_ref(), + &runtime.service_account_secret_key, + ) + .await?; + ensure_site_replication_bucket_replication_config( + &item.bucket, + &runtime.state, + &runtime.local_peer, + &runtime.service_account_secret_key, + ) + .await?; + } + Ok(()) } @@ -5489,6 +5519,7 @@ mod tests { assert_eq!(target.target_bucket, "photos"); assert_eq!(target.deployment_id, "remote"); assert_eq!(target.arn, "arn:rustfs:replication::remote:photos"); + assert_eq!(target.region, "us-east-1"); let credentials = target .credentials .as_ref() diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 9a5da2b4c..42399d00f 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -1132,6 +1132,23 @@ fn response_storage_class(info: &ObjectInfo, metadata: &HashMap) .map(StorageClass::from) } +fn response_storage_class_for_object_attributes( + info: &ObjectInfo, + metadata: &HashMap, + requested: bool, +) -> Option { + if !requested { + return None; + } + + info.storage_class + .clone() + .or_else(|| metadata.get(AMZ_STORAGE_CLASS).cloned()) + .or_else(|| Some(storageclass::STANDARD.to_string())) + .filter(|storage_class| !storage_class.is_empty()) + .map(StorageClass::from) +} + async fn apply_put_request_object_lock_opts( bucket: &str, object_lock_legal_hold_status: Option, @@ -2878,7 +2895,6 @@ impl DefaultObjectUsecase { validate_ssec_for_read(&info.user_defined, sse_customer_key.as_ref(), sse_customer_key_md5.as_ref())?; let metadata_map = info.user_defined.clone(); - let storage_class = response_storage_class(&info, &metadata_map); debug!( "GetObjectAttributes raw object_attributes={:?}", @@ -2886,6 +2902,8 @@ impl DefaultObjectUsecase { ); let requested = |name: &'static str| -> bool { object_attributes_requested(&object_attributes, name) }; + let storage_class = + response_storage_class_for_object_attributes(&info, &metadata_map, requested(ObjectAttributes::STORAGE_CLASS)); let e_tag = if requested(ObjectAttributes::ETAG) { info.etag.as_ref().map(|etag| to_s3s_etag(etag)) @@ -5804,6 +5822,38 @@ mod tests { ); } + #[test] + fn response_storage_class_for_object_attributes_defaults_to_standard_when_requested() { + let metadata = HashMap::new(); + let info = ObjectInfo { + storage_class: None, + user_defined: Arc::new(metadata.clone()), + ..Default::default() + }; + + assert_eq!( + response_storage_class_for_object_attributes(&info, &metadata, true) + .as_ref() + .map(StorageClass::as_str), + Some(storageclass::STANDARD) + ); + } + + #[test] + fn response_storage_class_for_object_attributes_skips_value_when_not_requested() { + let metadata = HashMap::new(); + let info = ObjectInfo { + storage_class: Some(storageclass::STANDARD_IA.to_string()), + user_defined: Arc::new(metadata.clone()), + ..Default::default() + }; + + assert!( + response_storage_class_for_object_attributes(&info, &metadata, false).is_none(), + "StorageClass must only be returned when explicitly requested" + ); + } + #[tokio::test] async fn build_get_object_output_context_returns_content_disposition() { let mut metadata = HashMap::new();