fix(replication): harden site-repl versioning sync (#3739)

This commit is contained in:
houseme
2026-06-22 19:18:42 +08:00
committed by GitHub
parent e57962d5e8
commit 64aef06c60
4 changed files with 108 additions and 3 deletions
@@ -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;
}
+24 -1
View File
@@ -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:?}"
))),
},
}
}
@@ -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()
+51 -1
View File
@@ -1132,6 +1132,23 @@ fn response_storage_class(info: &ObjectInfo, metadata: &HashMap<String, String>)
.map(StorageClass::from)
}
fn response_storage_class_for_object_attributes(
info: &ObjectInfo,
metadata: &HashMap<String, String>,
requested: bool,
) -> Option<StorageClass> {
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<ObjectLockLegalHoldStatus>,
@@ -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();