Compare commits

...

1 Commits

Author SHA1 Message Date
overtrue c568a54797 fix(replication): honor disabled version deletes 2026-07-30 13:37:01 +08:00
17 changed files with 1121 additions and 167 deletions
@@ -3237,6 +3237,116 @@ async fn test_bucket_replication_converges_delete_marker_and_version_purge() ->
.await?;
assert_eq!(retained.body.collect().await?.into_bytes().as_ref(), b"versioned replication payload v2");
source_client
.delete_object()
.bucket(source_bucket)
.key(object_key)
.version_id(delete_marker_version_id)
.send()
.await?;
assert_replication_converged(&source_client, source_bucket, &target_client, target_bucket).await?;
let target_state = list_replication_state(&target_client, target_bucket).await?;
assert!(
target_state.iter().all(|entry| entry.version_id != delete_marker_version_id),
"target retained the explicitly purged delete-marker version"
);
assert!(
target_state
.iter()
.any(|entry| !entry.delete_marker && entry.version_id == retained_version_id),
"target removed the retained object version while purging the delete marker: {target_state:?}"
);
Ok(())
}
#[tokio::test]
#[serial]
async fn test_bucket_replication_disabled_version_delete_preserves_target_versions_issue_5442() -> TestResult {
init_logging();
let mut source_env = RustFSTestEnvironment::new().await?;
let mut source_env_vars = replication_fast_env();
source_env_vars.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
source_env.start_rustfs_server_with_env(vec![], &source_env_vars).await?;
let mut target_env = RustFSTestEnvironment::new().await?;
target_env.start_rustfs_server_without_cleanup(vec![]).await?;
let source_bucket = "replication-no-version-delete-src";
let target_bucket = "replication-no-version-delete-dst";
let object_key = "retained-versions.txt";
let source_client = source_env.create_s3_client();
let target_client = target_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
target_client.create_bucket().bucket(target_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
enable_bucket_versioning(&target_env, target_bucket).await?;
let target_arn = set_replication_target(&source_env, source_bucket, &target_env, target_bucket).await?;
put_bucket_replication_with_delete_statuses(&source_env, source_bucket, &target_arn, "Enabled", None).await?;
let put = source_client
.put_object()
.bucket(source_bucket)
.key(object_key)
.body(ByteStream::from_static(b"permanent delete replication disabled"))
.send()
.await?;
let object_version_id = put.version_id().ok_or("source PUT omitted version ID")?.to_string();
assert_replication_converged(&source_client, source_bucket, &target_client, target_bucket).await?;
let delete = source_client
.delete_object()
.bucket(source_bucket)
.key(object_key)
.send()
.await?;
let delete_marker_version_id = delete
.version_id()
.ok_or("source DELETE omitted marker version ID")?
.to_string();
assert_eq!(delete.delete_marker(), Some(true));
assert_replication_converged(&source_client, source_bucket, &target_client, target_bucket).await?;
let expected_target_state = list_replication_state(&target_client, target_bucket).await?;
assert_eq!(expected_target_state.len(), 2);
assert!(
expected_target_state
.iter()
.any(|entry| !entry.delete_marker && entry.version_id == object_version_id)
);
assert!(
expected_target_state
.iter()
.any(|entry| entry.delete_marker && entry.version_id == delete_marker_version_id)
);
for version_id in [&object_version_id, &delete_marker_version_id] {
source_client
.delete_object()
.bucket(source_bucket)
.key(object_key)
.version_id(version_id)
.send()
.await?;
}
assert!(list_replication_state(&source_client, source_bucket).await?.is_empty());
let observation_deadline = tokio::time::Instant::now() + Duration::from_secs(10);
loop {
let target_state = list_replication_state(&target_client, target_bucket).await?;
assert_eq!(
target_state, expected_target_state,
"disabled permanent-delete replication changed target versions"
);
if tokio::time::Instant::now() >= observation_deadline {
break;
}
sleep(Duration::from_millis(100)).await;
}
Ok(())
}
@@ -1278,6 +1278,12 @@ fn build_remove_object_headers(version_id: Option<&str>, opts: &RemoveObjectOpti
fn resolve_delete_api_version_id(version_id: Option<String>, opts: &RemoveObjectOptions) -> Option<String> {
if opts.replication_request && opts.replication_delete_marker {
None
} else if opts.replication_request
&& version_id
.as_deref()
.is_some_and(|version_id| Uuid::parse_str(version_id).is_ok_and(|version_id| version_id.is_nil()))
{
Some("null".to_string())
} else {
version_id
}
@@ -2259,6 +2265,12 @@ mod tests {
assert_eq!(got.as_deref(), Some(vid.as_str()));
}
#[test]
fn null_version_purge_sends_s3_null_version_id() {
let got = resolve_delete_api_version_id(Some(Uuid::nil().to_string()), &remove_opts(true, false));
assert_eq!(got.as_deref(), Some("null"));
}
#[test]
fn delete_marker_propagation_omits_versionid_query_param() {
// Propagating a delete-marker CREATION (delete_marker=true): the target
@@ -13,6 +13,6 @@
// limitations under the License.
pub use rustfs_replication::{
ObjectOpts, ReplicationConfigurationExt, ReplicationTargetValidationError, replication_target_arns,
should_remove_replication_target, validate_replication_config_target_arns,
ObjectOpts, ReplicationConfigurationExt, ReplicationTargetValidationError, delete_replication_target_arns,
replication_target_arns, should_remove_replication_target, validate_replication_config_target_arns,
};
@@ -28,7 +28,7 @@ use super::replication_logging::{EVENT_RESYNC_CONFIG_LOOKUP_SKIPPED, LOG_COMPONE
use super::replication_metadata_boundary::ReplicationMetadataStore;
use super::replication_object_decision_boundary::{
MustReplicateOptions, ReplicationDeleteSource, ReplicationResyncTargetObject, delete_replication_missing_source_decision,
delete_replication_object_opts, resync_target_for_object,
delete_replication_object_opts, delete_replication_version_id, resync_target_for_object,
};
use super::replication_storage_boundary::{ObjectInfo, ObjectOptions, ObjectToDelete, object_to_delete_for_replication};
use super::replication_target_boundary::{BucketTargets, ReplicationTargetStore};
@@ -73,7 +73,7 @@ impl ReplicationConfig {
if oi.delete_marker {
let opts = ObjectOpts {
name: oi.name.clone(),
version_id: oi.version_id,
version_id: delete_replication_version_id(oi.delete_marker, oi.version_id, !oi.version_purge_status.is_empty()),
delete_marker: true,
op_type: ReplicationType::Delete,
existing_object: true,
@@ -299,8 +299,14 @@ pub(crate) async fn must_replicate(bucket: &str, object: &str, mopts: MustReplic
#[cfg(test)]
mod tests {
use s3s::dto::{Destination, ReplicationRule, ReplicationRuleStatus};
use s3s::dto::{
DeleteMarkerReplication, DeleteMarkerReplicationStatus, DeleteReplication, DeleteReplicationStatus, Destination,
ReplicationRule, ReplicationRuleStatus,
};
use uuid::Uuid;
use super::super::replication_filemeta_boundary::VersionPurgeStatusType;
use super::super::replication_target_boundary::BucketTarget;
use super::*;
fn replication_rule() -> ReplicationRule {
@@ -360,4 +366,56 @@ mod tests {
assert!(options.is_replication_request());
assert_eq!(options.user_tags(), "env=prod");
}
#[tokio::test]
async fn resync_distinguishes_delete_marker_creation_from_version_purge() {
let arn = "arn:rustfs:replication:us-east-1:target:bucket";
let mut rule = replication_rule();
rule.destination.bucket = arn.to_string();
rule.delete_marker_replication = Some(DeleteMarkerReplication {
status: Some(DeleteMarkerReplicationStatus::from_static(DeleteMarkerReplicationStatus::ENABLED)),
});
rule.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::DISABLED),
});
let config = ReplicationConfig::new(
Some(ReplicationConfiguration {
role: String::new(),
rules: vec![rule],
}),
Some(BucketTargets {
targets: vec![BucketTarget {
arn: arn.to_string(),
endpoint: "target.example".to_string(),
..Default::default()
}],
}),
);
let marker_version_id = Uuid::new_v4();
let marker = ObjectInfo {
bucket: "source".to_string(),
name: "object".to_string(),
delete_marker: true,
version_id: Some(marker_version_id),
replication_status: ReplicationStatusType::Pending,
..Default::default()
};
let creation = config
.resync(marker.clone(), ReplicateDecision::default(), &HashMap::new())
.await;
assert!(creation.targets.contains_key(arn));
let purge = config
.resync(
ObjectInfo {
version_purge_status: VersionPurgeStatusType::Pending,
..marker
},
ReplicateDecision::default(),
&HashMap::new(),
)
.await;
assert!(purge.targets.is_empty());
}
}
@@ -13,14 +13,14 @@
// limitations under the License.
pub use rustfs_replication::{
MustReplicateOptions, ReplicationDeleteScheduleInput, ReplicationDeleteStateSource, delete_replication_state_from_config,
delete_replication_version_id, should_schedule_delete_replication, should_use_existing_delete_replication_info,
should_use_existing_delete_replication_source,
MustReplicateOptions, ReplicationDeleteScheduleInput, ReplicationDeleteStateSource, delete_replication_parts,
delete_replication_state_from_config, delete_replication_version_id, should_schedule_delete_replication,
should_use_existing_delete_replication_info, should_use_existing_delete_replication_source,
};
pub(crate) use rustfs_replication::{
ReplicationDeleteSource, ReplicationMultipartPartInput, ReplicationResyncTargetObject,
delete_replication_missing_source_decision, delete_replication_object_opts, heal_uses_delete_replication_path,
is_retryable_delete_replication_head_error, is_version_delete_replication, replication_etags_match,
replication_multipart_complete_actual_size, replication_multipart_part_plan, resync_target_for_object,
should_retry_delete_marker_purge,
should_retry_delete_marker_purge, version_purge_target_missing,
};
@@ -67,6 +67,7 @@ const EVENT_REPLICATION_BACKPRESSURE: &str = "replication_backpressure";
const EVENT_REPLICATION_RESYNC_LOAD_SKIPPED: &str = "replication_resync_load_skipped";
const EVENT_REPLICATION_RESYNC_RECOVERED: &str = "replication_resync_recovered";
const EVENT_REPLICATION_MRF_QUEUE_UNAVAILABLE: &str = "replication_mrf_queue_unavailable";
const EVENT_REPLICATION_MRF_ENTRY_SKIPPED: &str = "replication_mrf_entry_skipped";
#[derive(Debug, Default)]
pub struct DurableMrfBacklog {
@@ -660,6 +661,21 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
for entry in entries.iter() {
match entry.op {
MrfOpKind::Delete => {
let Some(delete_parts) = entry.delete_parts_for_replay() else {
debug!(
event = EVENT_REPLICATION_MRF_ENTRY_SKIPPED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
bucket = %entry.bucket,
object = %entry.object,
reason = "invalid_delete_version_ids",
"Skipped invalid persisted replication delete"
);
continue;
};
let version_purge_id = delete_parts
.version_id
.or_else(|| delete_parts.delete_marker_version_id.filter(|_| !delete_parts.delete_marker));
// Reconstruct a heal delete and re-queue it. We do NOT call
// get_object_info here because the delete-marker or version may
// already be absent from the local store — that is expected.
@@ -674,15 +690,20 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
let oi = ObjectInfo {
bucket: entry.bucket.clone(),
name: entry.object.clone(),
version_id: entry.version_id,
delete_marker: entry.delete_marker,
version_id: version_purge_id,
delete_marker: delete_parts.delete_marker,
replication_status: if entry.replica {
ReplicationStatusType::Replica
} else {
ReplicationStatusType::Empty
},
..Default::default()
};
let dsc = check_replicate_delete(
&entry.bucket,
&ObjectToDelete {
object_name: entry.object.clone(),
version_id: entry.version_id,
version_id: version_purge_id,
..Default::default()
},
&oi,
@@ -708,9 +729,9 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
let dv = DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: entry.object.clone(),
version_id: entry.version_id,
delete_marker_version_id: entry.delete_marker_version_id,
delete_marker: entry.delete_marker,
version_id: delete_parts.version_id,
delete_marker_version_id: delete_parts.delete_marker_version_id,
delete_marker: delete_parts.delete_marker,
delete_marker_mtime,
replication_state: Some(rstate),
..Default::default()
@@ -2333,6 +2354,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
};
let second = MrfReplicateEntry {
@@ -2446,6 +2468,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
};
@@ -2479,6 +2502,7 @@ mod tests {
op: MrfOpKind::Delete,
delete_marker_version_id: Some(dm_vid),
delete_marker: true,
replica: false,
delete_marker_mtime: Some(mtime_nanos),
};
@@ -2512,6 +2536,7 @@ mod tests {
op: MrfOpKind::Delete,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
};
@@ -2540,6 +2565,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
},
MrfReplicateEntry {
@@ -2551,6 +2577,7 @@ mod tests {
op: MrfOpKind::Delete,
delete_marker_version_id: Some(del_dm_vid),
delete_marker: true,
replica: false,
delete_marker_mtime: None,
},
];
@@ -2580,6 +2607,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
};
assert_eq!(obj_entry.op, MrfOpKind::Object);
@@ -2594,6 +2622,7 @@ mod tests {
op: MrfOpKind::Delete,
delete_marker_version_id: Some(Uuid::new_v4()),
delete_marker: true,
replica: false,
delete_marker_mtime: None,
};
assert_eq!(del_entry.op, MrfOpKind::Delete);
@@ -2609,6 +2638,7 @@ mod tests {
op: MrfOpKind::default(),
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
};
assert_eq!(legacy_entry.op, MrfOpKind::Object, "legacy default must be Object");
@@ -2672,6 +2702,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
}];
let encoded = encode_mrf_file(&entries).expect("durable MRF backlog should encode");
@@ -2702,6 +2733,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
}])
.expect("invalid persisted entry should still encode for boundary testing");
@@ -13,7 +13,7 @@
// limitations under the License.
use super::replication_bandwidth_boundary;
use super::replication_config_boundary::{ObjectOpts, ReplicationConfigurationExt as _};
use super::replication_config_boundary::{ObjectOpts, ReplicationConfigurationExt as _, delete_replication_target_arns};
use super::replication_config_store::ReplicationConfigStore;
use super::replication_error_boundary::{Result, is_err_object_not_found, is_err_version_not_found};
use super::replication_event_sink::{EventArgs, send_event, send_local_event};
@@ -29,9 +29,10 @@ use super::replication_metadata_boundary::ReplicationMetadataStore;
use super::replication_msgp_boundary::ReplicationMsgpCodec;
use super::replication_object_config::{ReplicationConfig, check_replicate_delete, get_replication_config, must_replicate};
use super::replication_object_decision_boundary::{
MustReplicateOptions, ReplicationMultipartPartInput, heal_uses_delete_replication_path,
MustReplicateOptions, ReplicationMultipartPartInput, delete_replication_parts, heal_uses_delete_replication_path,
is_retryable_delete_replication_head_error, is_version_delete_replication, replication_etags_match,
replication_multipart_complete_actual_size, replication_multipart_part_plan, should_retry_delete_marker_purge,
version_purge_target_missing,
};
use super::replication_queue_boundary::DeletedObjectReplicationInfo;
use super::replication_resync_boundary::ResyncStatusType;
@@ -49,7 +50,7 @@ use super::replication_target_boundary::{
PutObjectOptions, PutObjectPartOptions, ReplicationTargetStore, TargetClient, replication_action_for_target_head,
replication_complete_multipart_options, replication_delete_marker_purge_remove_options, replication_delete_remove_options,
replication_force_delete_remove_options, replication_object_is_ssec_encrypted, replication_put_object_header_size,
replication_put_object_options, replication_target_head_is_newer_null_version,
replication_put_object_options, replication_target_head_is_newer_null_version, replication_target_version_id,
};
use super::replication_versioning_boundary::ReplicationVersioningStore;
use super::runtime_boundary as runtime_sources;
@@ -70,9 +71,8 @@ use rustfs_utils::http::{
AMZ_TAGGING_DIRECTIVE, SUFFIX_REPLICATION_RESET, SUFFIX_REPLICATION_STATUS, has_internal_suffix, insert_str,
};
use rustfs_utils::{DEFAULT_SIP_HASH_KEY, sip_hash};
#[cfg(test)]
use s3s::dto::ReplicationConfiguration;
use std::collections::HashMap;
use std::collections::{HashMap, HashSet};
use std::fmt::Display;
use std::sync::Arc;
use time::OffsetDateTime;
@@ -97,6 +97,8 @@ const EVENT_RESYNC_TASK_FAILED: &str = "replication_resync_task_failed";
const EVENT_RESYNC_TARGET_OPERATION_FAILED: &str = "replication_resync_target_operation_failed";
const EVENT_RESYNC_RUNTIME_CHANNEL_FAILED: &str = "replication_resync_runtime_channel_failed";
const ERR_REPLICATION_METADATA_COPY_UNSUPPORTED: &str = "metadata-only replication is not implemented";
const ERR_VERSION_PURGE_TARGET_STILL_EXISTS: &str = "target version still exists after replication purge";
const ERR_VERSION_PURGE_MISSING_VERSION_ID: &str = "version purge record is missing version ID";
const REPLICATION_TARGET_OFFLINE_ERROR_MARKERS: &[&str] = &[
"dispatch failure",
"timeouterror",
@@ -786,19 +788,46 @@ impl ReplicationResyncer {
}
if roi.delete_marker || !roi.version_purge_status.is_empty() {
let (version_id, dm_version_id) = if roi.version_purge_status.is_empty() {
(None, roi.version_id)
} else {
(roi.version_id, None)
let Some(parts) =
delete_replication_parts(roi.delete_marker, roi.version_id, !roi.version_purge_status.is_empty())
else {
debug!(
event = EVENT_RESYNC_OBJECT_PROCESSED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
bucket = %bucket_name,
object = %roi.name,
reason = "version_purge_missing_version_id",
"Failed to process resync version purge"
);
let status = TargetReplicationResyncStatus {
bucket: roi.bucket.clone(),
object: roi.name.clone(),
failed_count: 1,
error: Some(ERR_VERSION_PURGE_MISSING_VERSION_ID.to_string()),
..Default::default()
};
if let Err(err) = results_tx.send(status).await {
error!(
event = EVENT_RESYNC_RUNTIME_CHANNEL_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
bucket = %bucket_name,
reason = "status_channel_send_failed",
error = %err,
"Failed to send resync status"
);
}
continue;
};
let doi = DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: roi.name.clone(),
delete_marker_version_id: dm_version_id,
version_id,
delete_marker_version_id: parts.delete_marker_version_id,
version_id: parts.version_id,
replication_state: roi.replication_state.clone(),
delete_marker: roi.delete_marker,
delete_marker: parts.delete_marker,
delete_marker_mtime: roi.mod_time,
..Default::default()
},
@@ -821,21 +850,41 @@ impl ReplicationResyncer {
};
let reset_id = target_client.reset_id.clone();
let is_version_purge = !roi.version_purge_status.is_empty();
let head_result = head_object_with_proxy_stats(
&bucket_name,
target_client.as_ref(),
&target_client.bucket,
&roi.name,
roi.version_id.map(|v| v.to_string()),
replication_target_version_id(roi.version_id, is_version_purge),
)
.await;
let (size, err) = match head_result {
Ok(_) if is_version_purge => {
st.failed_count += 1;
(0, Some(ERR_VERSION_PURGE_TARGET_STILL_EXISTS.to_string()))
}
Ok(_) => {
st.replicated_count += 1;
st.replicated_size += roi.size;
(roi.size, None)
}
Err(err) if is_version_purge => {
let (is_not_found, code) = err
.as_service_error()
.map(|service_err| (service_err.is_not_found(), service_err.code()))
.unwrap_or((false, None));
let raw_status = err.raw_response().map(|response| response.status().as_u16());
let missing = version_purge_target_missing(is_not_found, code, raw_status);
if missing {
st.replicated_count += 1;
(0, None)
} else {
st.failed_count += 1;
(0, resync_target_error_detail(&err))
}
}
Err(err) if roi.delete_marker => {
// Verifying a replicated delete marker: only a
// definitive 404/NoSuchKey or 405/MethodNotAllowed
@@ -852,7 +901,7 @@ impl ReplicationResyncer {
};
if retryable {
st.failed_count += 1;
(0, Some(err))
(0, resync_target_error_detail(&err))
} else {
st.replicated_count += 1;
(0, None)
@@ -871,17 +920,17 @@ impl ReplicationResyncer {
}
Ok(None) => {
st.failed_count += 1;
(0, Some(err))
(0, resync_target_error_detail(&err))
}
Err(e2) => {
st.failed_count += 1;
(0, Some(e2))
(0, resync_target_error_detail(&e2))
}
}
}
Err(err) => {
st.failed_count += 1;
(0, Some(err))
(0, resync_target_error_detail(&err))
}
};
@@ -911,7 +960,7 @@ impl ReplicationResyncer {
"Processed resync object"
);
}
st.error = err.as_ref().and_then(resync_target_error_detail);
st.error = err;
if cancel_token.is_cancelled() {
return;
@@ -1065,22 +1114,27 @@ pub async fn get_heal_replicate_object_info(oi: &ObjectInfo, rcfg: &ReplicationC
}
let dsc = if heal_uses_delete_replication_path(oi.delete_marker, &oi.version_purge_status) {
check_replicate_delete(
oi.bucket.as_str(),
&ObjectToDelete {
object_name: oi.name.clone(),
version_id: oi.version_id,
..Default::default()
},
&oi,
&ObjectOptions {
versioned: ReplicationVersioningStore::prefix_enabled(&oi.bucket, &oi.name).await,
version_suspended: ReplicationVersioningStore::prefix_suspended(&oi.bucket, &oi.name).await,
..Default::default()
},
None,
)
.await
match delete_replication_parts(oi.delete_marker, oi.version_id, !oi.version_purge_status.is_empty()) {
Some(parts) => {
check_replicate_delete(
oi.bucket.as_str(),
&ObjectToDelete {
object_name: oi.name.clone(),
version_id: parts.version_id,
..Default::default()
},
&oi,
&ObjectOptions {
versioned: ReplicationVersioningStore::prefix_enabled(&oi.bucket, &oi.name).await,
version_suspended: ReplicationVersioningStore::prefix_suspended(&oi.bucket, &oi.name).await,
..Default::default()
},
None,
)
.await
}
None => ReplicateDecision::default(),
}
} else {
must_replicate(
oi.bucket.as_str(),
@@ -1150,7 +1204,6 @@ pub async fn replicate_delete<S: ReplicationStorage>(dobj: DeletedObjectReplicat
} else {
dobj.delete_object.version_id
};
let _rcfg = match get_replication_config(&bucket).await {
Ok(Some(config)) => config,
Ok(None) => {
@@ -1469,8 +1522,8 @@ pub async fn replicate_delete<S: ReplicationStorage>(dobj: DeletedObjectReplicat
delete_marker_version_id,
)
.await
&& replicate_delete_marker_purge_to_targets(&bucket_clone, &dobj_clone, &dsc_clone).await
{
replicate_delete_marker_purge_to_targets(&bucket_clone, &dobj_clone, &dsc_clone).await;
break;
}
tokio::time::sleep(TokioDuration::from_secs(1)).await;
@@ -1600,9 +1653,50 @@ async fn source_delete_marker_missing<S: EcstoreObjectOperations>(
}
}
async fn replicate_delete_marker_purge_to_targets(bucket: &str, dobj: &DeletedObjectReplicationInfo, dsc: &ReplicateDecision) {
fn delete_marker_purge_target_arns(config: &ReplicationConfiguration, dobj: &DeletedObjectReplicationInfo) -> HashSet<String> {
let replica = dobj
.delete_object
.replication_state
.as_ref()
.is_some_and(|state| state.replica_status == ReplicationStatusType::Replica);
delete_replication_target_arns(config, &dobj.delete_object.object_name, replica)
}
#[derive(Debug, PartialEq, Eq)]
enum DeleteMarkerPurgeConfig {
Apply(HashSet<String>),
Stop,
Retry,
}
fn delete_marker_purge_config<E>(
result: std::result::Result<Option<ReplicationConfiguration>, E>,
dobj: &DeletedObjectReplicationInfo,
) -> DeleteMarkerPurgeConfig {
match result {
Ok(Some(config)) => DeleteMarkerPurgeConfig::Apply(delete_marker_purge_target_arns(&config, dobj)),
Ok(None) => DeleteMarkerPurgeConfig::Stop,
Err(_) => DeleteMarkerPurgeConfig::Retry,
}
}
async fn replicate_delete_marker_purge_to_targets(
bucket: &str,
dobj: &DeletedObjectReplicationInfo,
dsc: &ReplicateDecision,
) -> bool {
let Some(delete_marker_version_id) = dobj.delete_object.delete_marker_version_id else {
return;
return true;
};
let marker_creation_purge_targets = if dobj.delete_object.delete_marker {
match delete_marker_purge_config(get_replication_config(bucket).await, dobj) {
DeleteMarkerPurgeConfig::Apply(targets) => Some(targets),
DeleteMarkerPurgeConfig::Stop => return true,
DeleteMarkerPurgeConfig::Retry => return false,
}
} else {
None
};
for tgt_entry in dsc.targets_map.values() {
@@ -1612,6 +1706,12 @@ async fn replicate_delete_marker_purge_to_targets(bucket: &str, dobj: &DeletedOb
if !dobj.target_arn.is_empty() && dobj.target_arn != tgt_entry.arn {
continue;
}
if marker_creation_purge_targets
.as_ref()
.is_some_and(|targets| !targets.contains(&tgt_entry.arn))
{
continue;
}
let Some(tgt_client) = ReplicationTargetStore::remote_target_client(bucket, &tgt_entry.arn).await else {
continue;
};
@@ -1625,6 +1725,7 @@ async fn replicate_delete_marker_purge_to_targets(bucket: &str, dobj: &DeletedOb
)
.await;
}
true
}
async fn replicate_force_delete_to_targets<S: ReplicationStorage>(dobj: &DeletedObjectReplicationInfo, storage: Arc<S>) {
@@ -1851,10 +1952,10 @@ async fn replicate_force_delete_to_targets<S: ReplicationStorage>(dobj: &Deleted
}
async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_client: Arc<TargetClient>) -> ReplicatedTargetInfo {
let version_id = if let Some(version_id) = &dobj.delete_object.delete_marker_version_id {
version_id.to_owned()
let version_id = if let Some(version_id) = dobj.delete_object.delete_marker_version_id {
Some(version_id)
} else {
dobj.delete_object.version_id.unwrap_or_default()
dobj.delete_object.version_id
};
let mut rinfo = dobj
@@ -1889,11 +1990,7 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli
return rinfo;
}
let version_id = if version_id.is_nil() {
None
} else {
Some(version_id.to_string())
};
let version_id = replication_target_version_id(version_id, is_version_purge);
if dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() {
match head_object_with_proxy_stats(
@@ -3143,6 +3240,10 @@ async fn replicate_object_with_multipart<S: ReplicationObjectIO>(ctx: MultipartR
#[cfg(test)]
mod tests {
use super::*;
use s3s::dto::{
DeleteReplication, DeleteReplicationStatus, Destination, ReplicaModifications, ReplicaModificationsStatus,
ReplicationRule, ReplicationRuleAndOperator, ReplicationRuleFilter, ReplicationRuleStatus, SourceSelectionCriteria, Tag,
};
use std::collections::HashMap;
use time::OffsetDateTime;
use uuid::Uuid;
@@ -3452,6 +3553,110 @@ mod tests {
);
}
#[test]
fn test_delete_marker_purge_targets_follow_delete_and_replica_modification_rules() {
fn rule(arn: &str, delete_status: &'static str) -> ReplicationRule {
ReplicationRule {
delete_marker_replication: None,
delete_replication: Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(delete_status),
}),
destination: Destination {
bucket: arn.to_string(),
..Default::default()
},
existing_object_replication: None,
filter: None,
id: Some(arn.to_string()),
prefix: Some("logs/".to_string()),
priority: Some(1),
source_selection_criteria: None,
status: ReplicationRuleStatus::from_static(ReplicationRuleStatus::ENABLED),
}
}
let enabled_arn = "arn:rustfs:replication:us-east-1:target:enabled";
let disabled_arn = "arn:rustfs:replication:us-east-1:target:disabled";
let delete_marker_version_id = Uuid::new_v4();
let mut config = ReplicationConfiguration {
role: String::new(),
rules: vec![
rule(enabled_arn, DeleteReplicationStatus::ENABLED),
rule(disabled_arn, DeleteReplicationStatus::DISABLED),
],
};
let mut dobj = DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: "logs/object.txt".to_string(),
delete_marker: true,
delete_marker_version_id: Some(delete_marker_version_id),
..Default::default()
},
..Default::default()
};
assert_eq!(delete_marker_purge_target_arns(&config, &dobj), HashSet::from([enabled_arn.to_string()]));
assert_eq!(
delete_marker_purge_config::<()>(Ok(Some(config.clone())), &dobj),
DeleteMarkerPurgeConfig::Apply(HashSet::from([enabled_arn.to_string()]))
);
dobj.delete_object.replication_state = Some(Default::default());
dobj.delete_object
.replication_state
.as_mut()
.expect("test replication state")
.replica_status = ReplicationStatusType::Replica;
assert!(delete_marker_purge_target_arns(&config, &dobj).is_empty());
config.rules[0].source_selection_criteria = Some(SourceSelectionCriteria {
replica_modifications: Some(ReplicaModifications {
status: ReplicaModificationsStatus::from_static(ReplicaModificationsStatus::ENABLED),
}),
sse_kms_encrypted_objects: None,
});
assert_eq!(delete_marker_purge_target_arns(&config, &dobj), HashSet::from([enabled_arn.to_string()]));
config.rules[0].prefix = None;
config.rules[0].filter = Some(ReplicationRuleFilter {
tag: Some(Tag {
key: Some("env".to_string()),
value: Some("prod".to_string()),
}),
..Default::default()
});
assert!(delete_marker_purge_target_arns(&config, &dobj).is_empty());
config.rules[0].filter = Some(ReplicationRuleFilter {
and: Some(ReplicationRuleAndOperator {
prefix: Some("logs/".to_string()),
tags: Some(vec![Tag {
key: Some("env".to_string()),
value: Some("prod".to_string()),
}]),
}),
..Default::default()
});
assert!(delete_marker_purge_target_arns(&config, &dobj).is_empty());
}
#[test]
fn test_delete_marker_purge_config_errors_are_retryable() {
let delete_marker_version_id = Uuid::new_v4();
let dobj = DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: "object.txt".to_string(),
delete_marker: true,
delete_marker_version_id: Some(delete_marker_version_id),
..Default::default()
},
..Default::default()
};
assert_eq!(delete_marker_purge_config::<()>(Ok(None), &dobj), DeleteMarkerPurgeConfig::Stop);
assert_eq!(delete_marker_purge_config::<()>(Err(()), &dobj), DeleteMarkerPurgeConfig::Retry);
}
#[test]
fn test_is_retryable_delete_replication_head_error_allows_delete_marker_head_responses() {
assert!(
@@ -31,10 +31,13 @@ use rustfs_utils::http::{
};
use time::OffsetDateTime;
use time::format_description::well_known::Rfc3339;
use uuid::Uuid;
pub(crate) use crate::bucket::bucket_target_sys::{
AdvancedPutOptions, PutObjectOptions, PutObjectPartOptions, RemoveObjectOptions, TargetClient,
};
#[cfg(test)]
pub(crate) use crate::bucket::target::BucketTarget;
pub(crate) use crate::bucket::target::BucketTargets;
use super::replication_config_store::ReplicationConfigStore;
@@ -305,6 +308,14 @@ pub(crate) fn replication_target_head_is_newer_null_version(object_info: &Object
target_is_newer_than_source_null_version(&replication_source_object(object_info), &replication_target_object(target))
}
pub(crate) fn replication_target_version_id(version_id: Option<Uuid>, version_purge: bool) -> Option<String> {
match version_id {
Some(version_id) if version_id.is_nil() && version_purge => Some("null".to_string()),
Some(version_id) if !version_id.is_nil() => Some(version_id.to_string()),
_ => None,
}
}
pub(crate) fn replication_delete_remove_options(
delete_marker: bool,
replication_mtime: Option<OffsetDateTime>,
@@ -519,6 +530,18 @@ mod tests {
assert!(force.replication_request);
}
#[test]
fn replication_target_version_id_preserves_null_purges() {
assert_eq!(replication_target_version_id(Some(Uuid::nil()), true).as_deref(), Some("null"));
assert_eq!(replication_target_version_id(Some(Uuid::nil()), false), None);
let version_id = Uuid::new_v4();
assert_eq!(
replication_target_version_id(Some(version_id), true).as_deref(),
Some(version_id.to_string().as_str())
);
}
#[test]
fn replication_complete_multipart_options_sets_actual_size() {
let options = replication_complete_multipart_options("1024".to_string());
+295 -35
View File
@@ -18,7 +18,9 @@ use crate::rule::ReplicationRuleExt as _;
use s3s::dto::DeleteMarkerReplicationStatus;
use s3s::dto::DeleteReplicationStatus;
use s3s::dto::Destination;
use s3s::dto::{ExistingObjectReplicationStatus, ReplicationConfiguration, ReplicationRuleStatus, ReplicationRules};
use s3s::dto::{
ExistingObjectReplicationStatus, ReplicationConfiguration, ReplicationRule, ReplicationRuleStatus, ReplicationRules,
};
use serde::{Deserialize, Serialize};
use std::collections::HashSet;
use uuid::Uuid;
@@ -45,6 +47,86 @@ pub trait ReplicationConfigurationExt {
fn filter_target_arns(&self, obj: &ObjectOpts) -> Vec<String>;
}
pub fn delete_replication_target_arns(config: &ReplicationConfiguration, object_name: &str, replica: bool) -> HashSet<String> {
let role = config.role.trim();
if !role.is_empty() && active_replication_rule_destination_arns(config).len() > 1 {
return HashSet::new();
}
let mut targets = HashSet::new();
let mut targets_with_unknown_tags = HashSet::new();
for rule in &config.rules {
if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED)
|| !object_name.starts_with(rule.prefix())
{
continue;
}
let arn = if role.is_empty() {
rule.destination.bucket.trim()
} else {
role
};
if arn.is_empty() {
continue;
}
targets.insert(arn.to_string());
if rule.filter.as_ref().is_some_and(|filter| {
filter.tag.is_some()
|| filter
.and
.as_ref()
.and_then(|and| and.tags.as_ref())
.is_some_and(|tags| !tags.is_empty())
}) {
targets_with_unknown_tags.insert(arn.to_string());
}
}
targets
.into_iter()
.filter(|arn| !targets_with_unknown_tags.contains(arn))
.filter(|arn| {
config.replicate(&ObjectOpts {
name: object_name.to_string(),
target_arn: arn.clone(),
version_id: Some(Uuid::nil()),
delete_marker: true,
op_type: ReplicationType::Delete,
replica,
..Default::default()
})
})
.collect()
}
fn rule_replicates(rule: &ReplicationRule, obj: &ObjectOpts) -> bool {
if let Some(status) = &rule.existing_object_replication
&& obj.existing_object
&& status.status == ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::DISABLED)
{
return false;
}
if obj.op_type == ReplicationType::Delete {
if !rule.metadata_replicate(obj) {
return false;
}
if obj.version_id.is_some() {
return rule
.delete_replication
.clone()
.is_some_and(|d| d.status == DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED));
}
return rule.delete_marker_replication.clone().is_some_and(|d| {
d.status == Some(DeleteMarkerReplicationStatus::from_static(DeleteMarkerReplicationStatus::ENABLED))
});
}
rule.metadata_replicate(obj)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReplicationTargetValidationError {
RoleWithMultipleDestinations,
@@ -204,37 +286,7 @@ impl ReplicationConfigurationExt for ReplicationConfiguration {
continue;
}
if let Some(status) = &rule.existing_object_replication
&& obj.existing_object
&& status.status == ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::DISABLED)
{
return false;
}
if obj.op_type == ReplicationType::Delete {
if !rule.metadata_replicate(obj) {
return false;
}
if obj.version_id.is_some() {
if obj.delete_marker {
return rule.delete_marker_replication.clone().is_some_and(|d| {
d.status == Some(DeleteMarkerReplicationStatus::from_static(DeleteMarkerReplicationStatus::ENABLED))
});
}
return rule
.delete_replication
.clone()
.is_some_and(|d| d.status == DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED));
} else {
return rule.delete_marker_replication.clone().is_some_and(|d| {
d.status == Some(DeleteMarkerReplicationStatus::from_static(DeleteMarkerReplicationStatus::ENABLED))
});
}
}
// Regular object/metadata replication
return rule.metadata_replicate(obj);
return rule_replicates(rule, obj);
}
false
}
@@ -305,7 +357,10 @@ impl ReplicationConfigurationExt for ReplicationConfiguration {
#[cfg(test)]
mod tests {
use super::*;
use s3s::dto::{DeleteMarkerReplication, Destination, ExistingObjectReplication, ReplicationRule};
use s3s::dto::{
DeleteMarkerReplication, DeleteReplication, Destination, ExistingObjectReplication, ReplicationRule,
ReplicationRuleFilter, Tag,
};
fn replication_rule(id: &str, arn: &str) -> ReplicationRule {
ReplicationRule {
@@ -327,6 +382,37 @@ mod tests {
}
}
#[test]
fn delete_replication_target_arns_uses_highest_priority_matching_rule() {
let arn = "arn:target:a";
let mut lower_priority = replication_rule("lower", arn);
lower_priority.priority = Some(1);
lower_priority.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
let mut higher_priority = replication_rule("higher", arn);
higher_priority.priority = Some(2);
higher_priority.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::DISABLED),
});
let mut config = ReplicationConfiguration {
role: String::new(),
rules: vec![lower_priority, higher_priority],
};
let targets = delete_replication_target_arns(&config, "object", false);
assert!(targets.is_empty(), "the higher-priority disabled rule must suppress the target");
config.rules[0].delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::DISABLED),
});
config.rules[1].delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
assert_eq!(delete_replication_target_arns(&config, "object", false), HashSet::from([arn.to_string()]));
}
#[test]
fn filter_target_arns_uses_role_when_role_is_present() {
let config = ReplicationConfiguration {
@@ -337,15 +423,151 @@ mod tests {
],
};
let arns = config.filter_target_arns(&ObjectOpts {
let opts = ObjectOpts {
name: "object".to_string(),
op_type: ReplicationType::Object,
..Default::default()
});
};
let arns = config.filter_target_arns(&opts);
assert_eq!(arns, vec!["arn:legacy:target".to_string()]);
}
#[test]
fn delete_replication_target_arns_uses_role_when_role_is_present() {
let mut rule = replication_rule("rule", "arn:target:a");
rule.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
let config = ReplicationConfiguration {
role: " arn:legacy:target ".to_string(),
rules: vec![rule],
};
assert_eq!(
delete_replication_target_arns(&config, "object", false),
HashSet::from(["arn:legacy:target".to_string()])
);
}
#[test]
fn delete_replication_target_arns_ignores_disjoint_prefix_rules() {
let arn = "arn:target:a";
let mut matching = replication_rule("matching", arn);
matching.prefix = None;
matching.filter = Some(ReplicationRuleFilter {
prefix: Some("logs/".to_string()),
..Default::default()
});
matching.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::DISABLED),
});
let mut unrelated = replication_rule("unrelated", arn);
unrelated.prefix = None;
unrelated.filter = Some(ReplicationRuleFilter {
prefix: Some("archive/".to_string()),
..Default::default()
});
unrelated.priority = Some(2);
unrelated.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
let mut config = ReplicationConfiguration {
role: String::new(),
rules: vec![matching, unrelated],
};
assert!(delete_replication_target_arns(&config, "logs/object", false).is_empty());
config.rules[0].delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
config.rules[1].delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::DISABLED),
});
assert_eq!(
delete_replication_target_arns(&config, "logs/object", false),
HashSet::from([arn.to_string()])
);
}
#[test]
fn delete_replication_target_arns_fails_closed_for_unknown_tag_rules() {
let arn = "arn:target:a";
let mut known = replication_rule("known", arn);
known.priority = Some(2);
known.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
let mut unknown = replication_rule("unknown", arn);
unknown.priority = Some(1);
unknown.prefix = None;
unknown.filter = Some(ReplicationRuleFilter {
tag: Some(Tag {
key: Some("env".to_string()),
value: Some("prod".to_string()),
}),
..Default::default()
});
unknown.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
let config = ReplicationConfiguration {
role: String::new(),
rules: vec![known, unknown],
};
assert!(delete_replication_target_arns(&config, "object", false).is_empty());
}
#[test]
fn delete_replication_target_arns_reuses_full_destination_rule_order() {
let arn = "arn:target:a";
let mut first = replication_rule("first", arn);
first.priority = Some(1);
first.destination.account = Some("account-a".to_string());
first.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::DISABLED),
});
let mut second = replication_rule("second", arn);
second.priority = Some(2);
second.destination.account = Some("account-b".to_string());
second.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
let config = ReplicationConfiguration {
role: String::new(),
rules: vec![first, second],
};
let opts = ObjectOpts {
name: "object".to_string(),
target_arn: arn.to_string(),
version_id: Some(Uuid::new_v4()),
delete_marker: true,
op_type: ReplicationType::Delete,
..Default::default()
};
assert!(!config.replicate(&opts));
assert!(delete_replication_target_arns(&config, "object", false).is_empty());
}
#[test]
fn delete_replication_target_arns_rejects_role_with_multiple_destinations() {
let mut first = replication_rule("first", "arn:target:a");
first.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
let mut second = replication_rule("second", "arn:target:b");
second.delete_replication = first.delete_replication.clone();
let config = ReplicationConfiguration {
role: "arn:legacy:target".to_string(),
rules: vec![first, second],
};
assert!(delete_replication_target_arns(&config, "object", false).is_empty());
}
#[test]
fn filter_target_arns_falls_back_to_role_when_destination_is_empty() {
let config = ReplicationConfiguration {
@@ -605,4 +827,42 @@ mod tests {
"highest-priority rule disables delete-marker replication, so the delete marker must not replicate"
);
}
#[test]
fn delete_marker_version_purge_requires_delete_replication() {
let arn = "arn:rustfs:replication:us-east-1:target:bucket";
let mut config = ReplicationConfiguration {
role: String::new(),
rules: vec![delete_marker_rule("delete-markers-only", arn, "", 1, true)],
};
let opts = ObjectOpts {
name: "object.txt".to_string(),
op_type: ReplicationType::Delete,
delete_marker: true,
version_id: Some(Uuid::new_v4()),
..Default::default()
};
assert!(
!config.replicate(&opts),
"permanently deleting a delete-marker version must not use the delete-marker replication setting"
);
config.rules[0].delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::DISABLED),
});
assert!(
!config.replicate(&opts),
"an explicitly disabled permanent-delete setting must not replicate"
);
config.rules[0].delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
assert!(
config.replicate(&opts),
"permanently deleting a delete-marker version should replicate when delete replication is enabled"
);
}
}
+31 -3
View File
@@ -15,7 +15,7 @@
use std::any::Any;
use crate::storage_api::DeletedObject;
use crate::{MrfOpKind, MrfReplicateEntry, ReplicationType, ReplicationWorkerOperation};
use crate::{MrfOpKind, MrfReplicateEntry, ReplicationStatusType, ReplicationType, ReplicationWorkerOperation};
#[derive(Debug, Clone, Default)]
pub struct DeletedObjectReplicationInfo {
@@ -42,6 +42,11 @@ impl ReplicationWorkerOperation for DeletedObjectReplicationInfo {
op: MrfOpKind::Delete,
delete_marker_version_id: self.delete_object.delete_marker_version_id,
delete_marker: self.delete_object.delete_marker,
replica: self
.delete_object
.replication_state
.as_ref()
.is_some_and(|state| state.replica_status == ReplicationStatusType::Replica),
// Persist the original delete-marker mtime as Unix nanoseconds so replay after a
// restart stamps the replica with the source timestamp rather than the replay time
// (backlog#867). None when unknown; replay then falls back to the current time.
@@ -85,14 +90,22 @@ pub fn is_retryable_delete_replication_head_error(is_not_found: bool, code: Opti
!(is_not_found || matches!(code, Some("MethodNotAllowed" | "405")))
}
pub fn version_purge_target_missing(is_not_found: bool, code: Option<&str>, raw_status: Option<u16>) -> bool {
if matches!(code, Some("MethodNotAllowed" | "405")) || raw_status == Some(405) {
return false;
}
is_not_found || matches!(code, Some("NoSuchVersion" | "NoSuchKey" | "NotFound" | "404")) || raw_status == Some(404)
}
#[cfg(test)]
mod tests {
use super::{
DeletedObjectReplicationInfo, is_retryable_delete_replication_head_error, is_version_delete_replication,
should_retry_delete_marker_purge,
should_retry_delete_marker_purge, version_purge_target_missing,
};
use crate::storage_api::DeletedObject;
use crate::{MrfOpKind, ReplicationType, ReplicationWorkerOperation};
use crate::{MrfOpKind, ReplicationState, ReplicationStatusType, ReplicationType, ReplicationWorkerOperation};
use uuid::Uuid;
#[test]
@@ -109,6 +122,10 @@ mod tests {
delete_marker_version_id: Some(delete_marker_version_id),
delete_marker: true,
delete_marker_mtime: Some(mtime),
replication_state: Some(ReplicationState {
replica_status: ReplicationStatusType::Replica,
..Default::default()
}),
..Default::default()
},
..Default::default()
@@ -122,6 +139,7 @@ mod tests {
assert_eq!(entry.delete_marker_version_id, Some(delete_marker_version_id));
assert_eq!(entry.op, MrfOpKind::Delete);
assert!(entry.delete_marker);
assert!(entry.replica);
// The original mtime must be persisted (as Unix nanos) so replay keeps the source
// timestamp instead of stamping the replica with the replay time (backlog#867).
assert_eq!(
@@ -212,4 +230,14 @@ mod tests {
assert!(!is_retryable_delete_replication_head_error(true, Some("NoSuchKey")));
assert!(is_retryable_delete_replication_head_error(false, Some("AccessDenied")));
}
#[test]
fn version_purge_target_missing_requires_not_found() {
assert!(version_purge_target_missing(true, Some("NoSuchVersion"), Some(404)));
assert!(version_purge_target_missing(false, Some("NoSuchVersion"), None));
assert!(version_purge_target_missing(false, None, Some(404)));
assert!(!version_purge_target_missing(false, None, None));
assert!(!version_purge_target_missing(true, Some("MethodNotAllowed"), Some(405)));
assert!(!version_purge_target_missing(true, Some("405"), None));
}
}
+7
View File
@@ -587,6 +587,11 @@ pub struct MrfReplicateEntry {
#[serde(rename = "deleteMarker", default)]
pub delete_marker: bool,
// For delete entries: whether the operation originated from a replica.
// Old files lack this field and therefore default to a local-source delete.
#[serde(rename = "replica", default)]
pub replica: bool,
// For delete entries: the original delete-marker mtime, persisted as Unix nanoseconds so
// replay stamps replicas with the source timestamp instead of the replay time. Old files
// lack this key; default=None means "unknown", and replay falls back to the current time
@@ -798,6 +803,7 @@ impl ReplicationWorkerOperation for ReplicateObjectInfo {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
}
}
@@ -852,6 +858,7 @@ impl ReplicateObjectInfo {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
}
}
+8 -6
View File
@@ -30,11 +30,12 @@ pub mod tagging;
pub use config::{
ObjectOpts, ReplicationConfigurationExt, ReplicationTargetValidationError, active_replication_rule_destination_arns,
replication_target_arns, should_remove_replication_target, validate_replication_config_target_arns,
delete_replication_target_arns, replication_target_arns, should_remove_replication_target,
validate_replication_config_target_arns,
};
pub use delete::{
DeletedObjectReplicationInfo, is_retryable_delete_replication_head_error, is_version_delete_replication,
should_retry_delete_marker_purge,
should_retry_delete_marker_purge, version_purge_target_missing,
};
pub use filemeta::{
REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, REPLICATE_HEAL, REPLICATE_HEAL_DELETE, REPLICATE_INCOMING,
@@ -54,10 +55,11 @@ pub use object::{
replication_etags_match, target_is_newer_than_source_null_version,
};
pub use operation::{
MustReplicateOptions, ReplicationDeleteScheduleInput, ReplicationDeleteSource, ReplicationDeleteStateSource,
ReplicationResyncTargetObject, delete_replication_missing_source_decision, delete_replication_object_opts,
delete_replication_state_from_config, delete_replication_version_id, heal_uses_delete_replication_path, is_ssec_encrypted,
resync_target_for_object, should_schedule_delete_replication, should_use_existing_delete_replication_info,
MustReplicateOptions, ReplicationDeleteParts, ReplicationDeleteScheduleInput, ReplicationDeleteSource,
ReplicationDeleteStateSource, ReplicationResyncTargetObject, delete_replication_missing_source_decision,
delete_replication_object_opts, delete_replication_parts, delete_replication_state_from_config,
delete_replication_version_id, heal_uses_delete_replication_path, is_ssec_encrypted, resync_target_for_object,
should_schedule_delete_replication, should_use_existing_delete_replication_info,
should_use_existing_delete_replication_source,
};
pub use queue::{
+70 -1
View File
@@ -13,6 +13,7 @@
// limitations under the License.
use byteorder::{ByteOrder, LittleEndian};
use uuid::Uuid;
use crate::{Error, Result};
@@ -21,6 +22,31 @@ pub use crate::filemeta::{MrfOpKind, MrfReplicateEntry};
pub const MRF_META_FORMAT: u16 = 1;
pub const MRF_META_VERSION: u16 = 1;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MrfDeleteParts {
pub version_id: Option<Uuid>,
pub delete_marker_version_id: Option<Uuid>,
pub delete_marker: bool,
}
impl MrfReplicateEntry {
pub fn delete_parts_for_replay(&self) -> Option<MrfDeleteParts> {
match (self.version_id, self.delete_marker_version_id) {
(Some(version_id), None) => Some(MrfDeleteParts {
version_id: Some(version_id),
delete_marker_version_id: None,
delete_marker: false,
}),
(None, Some(delete_marker_version_id)) => Some(MrfDeleteParts {
version_id: None,
delete_marker_version_id: Some(delete_marker_version_id),
delete_marker: self.delete_marker,
}),
_ => None,
}
}
}
pub fn encode_mrf_file(entries: &[MrfReplicateEntry]) -> Result<Vec<u8>> {
let payload = rmp_serde::to_vec_named(entries).map_err(|e| Error::Other(e.to_string()))?;
let mut data = Vec::with_capacity(4 + payload.len());
@@ -54,7 +80,6 @@ pub fn decode_mrf_file(data: &[u8]) -> Result<Vec<MrfReplicateEntry>> {
#[cfg(test)]
mod tests {
use super::*;
use uuid::Uuid;
#[test]
fn mrf_file_round_trips_object_and_delete_entries() {
@@ -70,6 +95,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
},
MrfReplicateEntry {
@@ -81,6 +107,7 @@ mod tests {
op: MrfOpKind::Delete,
delete_marker_version_id: Some(del_vid),
delete_marker: true,
replica: true,
delete_marker_mtime: Some(1_705_312_200_123_456_789),
},
];
@@ -95,6 +122,7 @@ mod tests {
assert_eq!(decoded[1].delete_marker_version_id, Some(del_vid));
assert_eq!(decoded[1].op, MrfOpKind::Delete);
assert!(decoded[1].delete_marker);
assert!(decoded[1].replica);
assert_eq!(
decoded[1].delete_marker_mtime,
Some(1_705_312_200_123_456_789),
@@ -102,6 +130,46 @@ mod tests {
);
}
#[test]
fn legacy_version_purge_replay_clears_delete_marker_creation() {
let entry = MrfReplicateEntry {
bucket: "bucket".to_string(),
object: "object".to_string(),
version_id: Some(Uuid::new_v4()),
retry_count: 0,
size: 0,
delete_marker: true,
op: MrfOpKind::Delete,
delete_marker_version_id: None,
replica: false,
delete_marker_mtime: None,
};
let encoded = encode_mrf_file(&[entry]).expect("legacy MRF entry should encode");
let decoded = decode_mrf_file(&encoded).expect("legacy MRF entry should decode");
assert_eq!(decoded[0].delete_parts_for_replay().map(|parts| parts.delete_marker), Some(false));
}
#[test]
fn invalid_delete_id_shapes_fail_closed_on_replay() {
for (version_id, delete_marker_version_id) in [(None, None), (Some(Uuid::new_v4()), Some(Uuid::new_v4()))] {
let entry = MrfReplicateEntry {
bucket: "bucket".to_string(),
object: "object".to_string(),
version_id,
retry_count: 0,
size: 0,
op: MrfOpKind::Delete,
delete_marker_version_id,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
};
assert_eq!(entry.delete_parts_for_replay(), None);
}
}
#[test]
fn mrf_legacy_file_without_op_decodes_as_object() {
let mut payload = Vec::new();
@@ -129,6 +197,7 @@ mod tests {
assert_eq!(decoded[0].retry_count, 2);
assert_eq!(decoded[0].size, 100);
assert_eq!(decoded[0].op, MrfOpKind::Object);
assert!(!decoded[0].replica);
// Old files lack the deleteMarkerMtime key; it must default to None so replay keeps the
// pre-#867 fallback to the current time.
assert_eq!(decoded[0].delete_marker_mtime, None);
+117 -27
View File
@@ -151,6 +151,11 @@ pub fn delete_replication_state_from_config(
let pending_status = decision.pending_status();
let mut state = ReplicationState {
replica_status: if source.replica {
ReplicationStatusType::Replica
} else {
ReplicationStatusType::Empty
},
replicate_decision_str: decision.to_string(),
..Default::default()
};
@@ -174,11 +179,31 @@ pub struct ReplicationDeleteScheduleInput<'a> {
pub deleted_delete_marker_version: bool,
}
fn delete_version_purge_source_status(status: &ReplicationStatusType) -> bool {
status == &ReplicationStatusType::Replica
|| status == &ReplicationStatusType::Pending
|| status == &ReplicationStatusType::Completed
|| status == &ReplicationStatusType::Failed
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ReplicationDeleteParts {
pub delete_marker: bool,
pub version_id: Option<Uuid>,
pub delete_marker_version_id: Option<Uuid>,
}
pub fn delete_replication_parts(
source_delete_marker: bool,
source_version_id: Option<Uuid>,
version_purge: bool,
) -> Option<ReplicationDeleteParts> {
if version_purge {
return source_version_id.map(|version_id| ReplicationDeleteParts {
delete_marker: false,
version_id: Some(version_id),
delete_marker_version_id: None,
});
}
Some(ReplicationDeleteParts {
delete_marker: source_delete_marker,
version_id: None,
delete_marker_version_id: source_version_id,
})
}
pub fn should_schedule_delete_replication(input: ReplicationDeleteScheduleInput<'_>) -> bool {
@@ -186,14 +211,11 @@ pub fn should_schedule_delete_replication(input: ReplicationDeleteScheduleInput<
return false;
}
if input.version_id_requested && !input.deleted_delete_marker_version && !input.source_delete_marker {
return delete_version_purge_source_status(input.source_replication_status);
if input.version_id_requested {
return input.source_version_purge_status == &VersionPurgeStatusType::Pending;
}
input.source_replication_status == &ReplicationStatusType::Replica
|| input.source_replication_status == &ReplicationStatusType::Pending
|| input.source_version_purge_status == &VersionPurgeStatusType::Pending
|| (input.deleted_delete_marker_version && input.source_replication_status == &ReplicationStatusType::Completed)
input.source_replication_status == &ReplicationStatusType::Pending
}
pub fn delete_replication_version_id(
@@ -312,19 +334,20 @@ pub fn resync_target_for_object(
#[cfg(test)]
mod tests {
use super::{
MustReplicateOptions, ReplicationDeleteScheduleInput, ReplicationDeleteSource, ReplicationDeleteStateSource,
ReplicationResyncTargetObject, delete_replication_missing_source_decision, delete_replication_object_opts,
delete_replication_state_from_config, delete_replication_version_id, heal_uses_delete_replication_path,
is_ssec_encrypted, resync_target_for_object, should_schedule_delete_replication,
should_use_existing_delete_replication_info, should_use_existing_delete_replication_source,
MustReplicateOptions, ReplicationDeleteParts, ReplicationDeleteScheduleInput, ReplicationDeleteSource,
ReplicationDeleteStateSource, ReplicationResyncTargetObject, delete_replication_missing_source_decision,
delete_replication_object_opts, delete_replication_parts, delete_replication_state_from_config,
delete_replication_version_id, heal_uses_delete_replication_path, is_ssec_encrypted, resync_target_for_object,
should_schedule_delete_replication, should_use_existing_delete_replication_info,
should_use_existing_delete_replication_source,
};
use crate::http::{AMZ_BUCKET_REPLICATION_STATUS, SSEC_ALGORITHM_HEADER};
use crate::storage_api::ObjectToDelete;
use crate::{ReplicationStatusType, ReplicationType, VersionPurgeStatusType, target_reset_header};
use s3s::dto::{
DeleteMarkerReplication, DeleteMarkerReplicationStatus, Destination, ExistingObjectReplication,
ExistingObjectReplicationStatus, ReplicaModifications, ReplicaModificationsStatus, ReplicationConfiguration,
ReplicationRule, ReplicationRuleStatus, SourceSelectionCriteria,
DeleteMarkerReplication, DeleteMarkerReplicationStatus, DeleteReplication, DeleteReplicationStatus, Destination,
ExistingObjectReplication, ExistingObjectReplicationStatus, ReplicaModifications, ReplicaModificationsStatus,
ReplicationConfiguration, ReplicationRule, ReplicationRuleStatus, SourceSelectionCriteria,
};
use std::collections::HashMap;
use time::{Duration, OffsetDateTime};
@@ -474,6 +497,7 @@ mod tests {
.expect("replica delete marker should be forwarded to downstream targets");
let pending = format!("{arn}=PENDING;");
assert_eq!(state.replica_status, ReplicationStatusType::Replica);
assert_eq!(state.replication_status_internal.as_deref(), Some(pending.as_str()));
assert_eq!(state.replicate_decision_str, format!("{arn}=true;false;{arn};"));
assert!(state.targets.contains_key(arn));
@@ -500,9 +524,13 @@ mod tests {
#[test]
fn delete_replication_state_tracks_delete_marker_version_purges() {
let arn = "arn:aws:s3:::target-bucket";
let mut rule = delete_replication_rule(arn, false);
rule.delete_replication = Some(DeleteReplication {
status: DeleteReplicationStatus::from_static(DeleteReplicationStatus::ENABLED),
});
let config = ReplicationConfiguration {
role: arn.to_string(),
rules: vec![delete_replication_rule(arn, false)],
rules: vec![rule],
};
let source = ReplicationDeleteStateSource {
name: "test/object.txt".to_string(),
@@ -513,9 +541,10 @@ mod tests {
};
let state = delete_replication_state_from_config(&config, &source)
.expect("delete-marker version purge should honor delete-marker replication rules");
.expect("delete-marker version purge should honor delete replication rules");
let pending = format!("{arn}=PENDING;");
assert_eq!(state.replica_status, ReplicationStatusType::Empty);
assert_eq!(state.version_purge_status_internal.as_deref(), Some(pending.as_str()));
assert_eq!(state.replicate_decision_str, format!("{arn}=true;false;{arn};"));
assert!(state.purge_targets.contains_key(arn));
@@ -534,13 +563,13 @@ mod tests {
}
#[test]
fn delete_replication_schedule_keeps_marker_and_version_purges() {
fn delete_replication_schedule_uses_pending_state_from_current_delete_rules() {
assert!(should_schedule_delete_replication(ReplicationDeleteScheduleInput {
replication_request: false,
version_id_requested: true,
source_delete_marker: true,
source_replication_status: &ReplicationStatusType::Completed,
source_version_purge_status: &VersionPurgeStatusType::Empty,
source_replication_status: &ReplicationStatusType::Empty,
source_version_purge_status: &VersionPurgeStatusType::Pending,
deleted_delete_marker_version: true,
}));
assert!(should_schedule_delete_replication(ReplicationDeleteScheduleInput {
@@ -548,19 +577,54 @@ mod tests {
version_id_requested: true,
source_delete_marker: false,
source_replication_status: &ReplicationStatusType::Completed,
source_version_purge_status: &VersionPurgeStatusType::Empty,
source_version_purge_status: &VersionPurgeStatusType::Pending,
deleted_delete_marker_version: false,
}));
assert!(should_schedule_delete_replication(ReplicationDeleteScheduleInput {
replication_request: false,
version_id_requested: false,
source_delete_marker: false,
source_replication_status: &ReplicationStatusType::Empty,
source_version_purge_status: &VersionPurgeStatusType::Pending,
source_replication_status: &ReplicationStatusType::Pending,
source_version_purge_status: &VersionPurgeStatusType::Empty,
deleted_delete_marker_version: false,
}));
}
#[test]
fn delete_replication_schedule_skips_non_pending_states() {
for replication_status in [
ReplicationStatusType::Empty,
ReplicationStatusType::Replica,
ReplicationStatusType::Completed,
ReplicationStatusType::CompletedLegacy,
ReplicationStatusType::Failed,
] {
assert!(!should_schedule_delete_replication(ReplicationDeleteScheduleInput {
replication_request: false,
version_id_requested: false,
source_delete_marker: false,
source_replication_status: &replication_status,
source_version_purge_status: &VersionPurgeStatusType::Empty,
deleted_delete_marker_version: false,
}));
}
for version_purge_status in [
VersionPurgeStatusType::Empty,
VersionPurgeStatusType::Complete,
VersionPurgeStatusType::Failed,
] {
assert!(!should_schedule_delete_replication(ReplicationDeleteScheduleInput {
replication_request: false,
version_id_requested: true,
source_delete_marker: true,
source_replication_status: &ReplicationStatusType::Completed,
source_version_purge_status: &version_purge_status,
deleted_delete_marker_version: true,
}));
}
}
#[test]
fn delete_replication_version_id_splits_marker_creation_and_purge() {
let version_id = Uuid::new_v4();
@@ -569,6 +633,32 @@ mod tests {
assert_eq!(delete_replication_version_id(true, Some(version_id), true), Some(version_id));
}
#[test]
fn delete_replication_parts_fail_closed_without_purge_version() {
assert_eq!(delete_replication_parts(true, None, true), None);
for version_id in [Uuid::nil(), Uuid::new_v4()] {
assert_eq!(
delete_replication_parts(true, Some(version_id), true),
Some(ReplicationDeleteParts {
delete_marker: false,
version_id: Some(version_id),
delete_marker_version_id: None,
})
);
}
let marker_version_id = Uuid::new_v4();
assert_eq!(
delete_replication_parts(true, Some(marker_version_id), false),
Some(ReplicationDeleteParts {
delete_marker: true,
version_id: None,
delete_marker_version_id: Some(marker_version_id),
})
);
}
#[test]
fn delete_replication_source_selection_prefers_existing_marker_source_only_for_replica_requests() {
assert!(should_use_existing_delete_replication_source(true, true, true));
+25 -13
View File
@@ -17,8 +17,8 @@ use std::any::Any;
use crate::storage_api::DeletedObject;
use crate::{
DeletedObjectReplicationInfo, MrfReplicateEntry, REPLICATE_EXISTING, REPLICATE_HEAL, REPLICATE_HEAL_DELETE,
ReplicateObjectInfo, ReplicationStatusType, ReplicationType, ReplicationWorkerOperation, ResyncDecision,
VersionPurgeStatusType,
ReplicateObjectInfo, ReplicationDeleteParts, ReplicationStatusType, ReplicationType, ReplicationWorkerOperation,
ResyncDecision, VersionPurgeStatusType, delete_replication_parts,
};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
@@ -109,7 +109,11 @@ pub fn replication_heal_queue_action(roi: &mut ReplicateObjectInfo) -> Replicati
}
if roi.delete_marker || !roi.version_purge_status.is_empty() {
let delete_info = heal_deleted_object_replication_info(roi);
let Some(parts) = delete_replication_parts(roi.delete_marker, roi.version_id, !roi.version_purge_status.is_empty())
else {
return ReplicationHealQueueAction::Skip;
};
let delete_info = heal_deleted_object_replication_info(roi, parts);
if is_pending_or_failed_object_heal(roi) || is_pending_or_failed_version_purge(roi) {
return ReplicationHealQueueAction::QueueDelete(delete_info);
@@ -144,21 +148,18 @@ pub fn replication_heal_queue_action(roi: &mut ReplicateObjectInfo) -> Replicati
ReplicationHealQueueAction::Skip
}
fn heal_deleted_object_replication_info(roi: &ReplicateObjectInfo) -> DeletedObjectReplicationInfo {
let (version_id, delete_marker_version_id) = if roi.version_purge_status.is_empty() {
(None, roi.version_id)
} else {
(roi.version_id, None)
};
fn heal_deleted_object_replication_info(
roi: &ReplicateObjectInfo,
parts: ReplicationDeleteParts,
) -> DeletedObjectReplicationInfo {
DeletedObjectReplicationInfo {
delete_object: DeletedObject {
object_name: roi.name.clone(),
delete_marker_version_id,
version_id,
delete_marker_version_id: parts.delete_marker_version_id,
version_id: parts.version_id,
replication_state: roi.replication_state.clone(),
delete_marker_mtime: roi.mod_time,
delete_marker: roi.delete_marker,
delete_marker: parts.delete_marker,
..Default::default()
},
bucket: roi.bucket.clone(),
@@ -399,6 +400,7 @@ mod tests {
let version_id = Uuid::new_v4();
let mut roi = replicate_object_info(ReplicationStatusType::Completed);
roi.version_id = Some(version_id);
roi.delete_marker = true;
roi.version_purge_status = VersionPurgeStatusType::Pending;
let action = replication_heal_queue_action(&mut roi);
@@ -408,6 +410,16 @@ mod tests {
};
assert_eq!(delete_info.delete_object.version_id, Some(version_id));
assert_eq!(delete_info.delete_object.delete_marker_version_id, None);
assert!(!delete_info.delete_object.delete_marker);
}
#[test]
fn heal_queue_action_skips_version_purge_without_version_id() {
let mut roi = replicate_object_info(ReplicationStatusType::Completed);
roi.delete_marker = true;
roi.version_purge_status = VersionPurgeStatusType::Pending;
assert!(matches!(replication_heal_queue_action(&mut roi), ReplicationHealQueueAction::Skip));
}
#[test]
+2
View File
@@ -1044,6 +1044,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
},
MrfReplicateEntry {
@@ -1055,6 +1056,7 @@ mod tests {
op: MrfOpKind::Object,
delete_marker_version_id: None,
delete_marker: false,
replica: false,
delete_marker_mtime: None,
},
],
+65 -21
View File
@@ -2637,18 +2637,24 @@ where
fn delete_replication_state_source<'a>(
opts: &ObjectOptions,
existing_object_info: Option<&'a ObjectInfo>,
deleted_object_info: &'a ObjectInfo,
deleted_object_source: &'a ObjectInfo,
delete_result: &'a ObjectInfo,
) -> &'a ObjectInfo {
let replication_source = if opts.replication_request {
deleted_object_source
} else {
delete_result
};
if should_use_existing_delete_replication_source(
opts.replication_request,
deleted_object_info.delete_marker,
replication_source.delete_marker,
existing_object_info.is_some(),
) && let Some(existing) = existing_object_info
{
return existing;
}
deleted_object_info
replication_source
}
const AMZ_SNOWBALL_EXTRACT_COMPAT: &str = "X-Amz-Snowball-Auto-Extract";
@@ -7190,14 +7196,14 @@ impl DefaultObjectUsecase {
let _delete_tail_guard = DeleteTailActivityGuard::new(DeleteTailStage::Tail);
let deleted_object_source = deleted_replication_info.unwrap_or(&obj_info);
let replication_state_source =
delete_replication_state_source(&opts, existing_object_info.as_ref(), deleted_object_source);
delete_replication_state_source(&opts, existing_object_info.as_ref(), deleted_object_source, &obj_info);
let deleted_delete_marker_version = deleted_replication_info.is_some_and(|info| info.delete_marker);
let delete_replication_version_id = delete_replication_version_id(deleted_object_source, deleted_delete_marker_version);
let schedule_delete_replication = if opts.replication_request && replica {
should_schedule_replica_delete_replication(&bucket, replication_state_source, delete_replication_version_id).await
} else {
should_schedule_delete_replication(&opts, deleted_object_source, deleted_delete_marker_version)
should_schedule_delete_replication(&opts, replication_state_source, deleted_delete_marker_version)
};
if schedule_delete_replication {
@@ -13388,7 +13394,7 @@ mod tests {
}
#[test]
fn should_schedule_delete_replication_keeps_delete_marker_version_purge_from_source() {
fn should_schedule_delete_replication_uses_pending_marker_version_purge_state() {
let opts = ObjectOptions {
replication_request: false,
version_id: Some(Uuid::new_v4().to_string()),
@@ -13396,18 +13402,18 @@ mod tests {
};
let replication_source = ObjectInfo {
delete_marker: true,
replication_status: ReplicationStatusType::Completed,
version_purge_status: VersionPurgeStatusType::Pending,
..Default::default()
};
assert!(
should_schedule_delete_replication(&opts, &replication_source, true),
"source-side delete-marker version purge still needs replication scheduling"
"source-side delete-marker version purge needs scheduling when the current delete rule marked it pending"
);
}
#[test]
fn should_schedule_delete_replication_keeps_object_version_purge_from_completed_source() {
fn should_schedule_delete_replication_uses_pending_object_version_purge_state() {
let opts = ObjectOptions {
replication_request: false,
version_id: Some(Uuid::new_v4().to_string()),
@@ -13416,12 +13422,13 @@ mod tests {
let replication_source = ObjectInfo {
delete_marker: false,
replication_status: ReplicationStatusType::Completed,
version_purge_status: VersionPurgeStatusType::Pending,
..Default::default()
};
assert!(
should_schedule_delete_replication(&opts, &replication_source, false),
"source-side object version purge must still enqueue delete replication after the original PUT completed"
"source-side object version purge needs scheduling when the current delete rule marked it pending"
);
}
@@ -13681,7 +13688,7 @@ mod tests {
}
#[test]
fn delete_replication_state_from_config_tracks_delete_marker_version_purges() {
fn delete_replication_state_from_config_skips_delete_marker_version_purges_when_delete_is_disabled() {
let arn = "arn:aws:s3:::target-bucket".to_string();
let config = ReplicationConfiguration {
role: arn.clone(),
@@ -13691,7 +13698,7 @@ mod tests {
}),
delete_replication: None,
destination: Destination {
bucket: arn.clone(),
bucket: arn,
..Default::default()
},
existing_object_replication: Some(ExistingObjectReplication {
@@ -13714,13 +13721,10 @@ mod tests {
};
let version_id = Some(Uuid::new_v4());
let state = delete_replication_state_from_config(&config, &obj_info, version_id, false)
.expect("delete-marker version purge should honor delete-marker replication rules");
let pending = format!("{arn}=PENDING;");
assert_eq!(state.version_purge_status_internal.as_deref(), Some(pending.as_str()));
assert_eq!(state.replicate_decision_str, format!("{arn}=true;false;{arn};"));
assert!(state.purge_targets.contains_key(&arn));
assert!(
delete_replication_state_from_config(&config, &obj_info, version_id, false).is_none(),
"delete-marker version purge must remain local when delete replication is disabled"
);
}
#[test]
@@ -13741,7 +13745,7 @@ mod tests {
..Default::default()
};
let source = delete_replication_state_source(&opts, Some(&existing), &deleted);
let source = delete_replication_state_source(&opts, Some(&existing), &deleted, &deleted);
assert_eq!(source.replication_status, ReplicationStatusType::Completed);
assert!(
@@ -13764,7 +13768,7 @@ mod tests {
..Default::default()
};
let source = delete_replication_state_source(&opts, Some(&existing), &deleted);
let source = delete_replication_state_source(&opts, Some(&existing), &deleted, &deleted);
assert!(
source.delete_marker,
@@ -13772,6 +13776,46 @@ mod tests {
);
}
#[test]
fn delete_replication_state_source_uses_current_state_for_source_version_purge() {
let opts = ObjectOptions {
version_id: Some(Uuid::new_v4().to_string()),
..Default::default()
};
let existing = ObjectInfo {
replication_status: ReplicationStatusType::Completed,
..Default::default()
};
let delete_result = ObjectInfo {
version_purge_status: VersionPurgeStatusType::Pending,
..Default::default()
};
let source = delete_replication_state_source(&opts, Some(&existing), &existing, &delete_result);
assert_eq!(source.version_purge_status, VersionPurgeStatusType::Pending);
}
#[test]
fn delete_replication_state_source_preserves_tags_for_replica_version_purge() {
let opts = ObjectOptions {
replication_request: true,
version_id: Some(Uuid::new_v4().to_string()),
..Default::default()
};
let existing = ObjectInfo {
user_tags: Arc::new("environment=production".to_string()),
replication_status: ReplicationStatusType::Replica,
..Default::default()
};
let delete_result = ObjectInfo::default();
let source = delete_replication_state_source(&opts, Some(&existing), &existing, &delete_result);
assert_eq!(source.user_tags.as_str(), "environment=production");
assert_eq!(source.replication_status, ReplicationStatusType::Replica);
}
#[test]
fn replica_delete_enrichment_must_not_reuse_upstream_targets() {
let upstream_state = ReplicationState {