diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index 8c6ff2a04..2aa12b296 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -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(()) } diff --git a/crates/ecstore/src/bucket/bucket_target_sys.rs b/crates/ecstore/src/bucket/bucket_target_sys.rs index cf2a3e1ec..f91956e82 100644 --- a/crates/ecstore/src/bucket/bucket_target_sys.rs +++ b/crates/ecstore/src/bucket/bucket_target_sys.rs @@ -1278,6 +1278,12 @@ fn build_remove_object_headers(version_id: Option<&str>, opts: &RemoveObjectOpti fn resolve_delete_api_version_id(version_id: Option, opts: &RemoveObjectOptions) -> Option { 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 diff --git a/crates/ecstore/src/bucket/replication/replication_config_boundary.rs b/crates/ecstore/src/bucket/replication/replication_config_boundary.rs index da3e52a9d..c1f4b4377 100644 --- a/crates/ecstore/src/bucket/replication/replication_config_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_config_boundary.rs @@ -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, }; diff --git a/crates/ecstore/src/bucket/replication/replication_object_config.rs b/crates/ecstore/src/bucket/replication/replication_object_config.rs index 2d23466d3..478bbba0b 100644 --- a/crates/ecstore/src/bucket/replication/replication_object_config.rs +++ b/crates/ecstore/src/bucket/replication/replication_object_config.rs @@ -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()); + } } diff --git a/crates/ecstore/src/bucket/replication/replication_object_decision_boundary.rs b/crates/ecstore/src/bucket/replication/replication_object_decision_boundary.rs index d6c4f1ceb..ab630cf75 100644 --- a/crates/ecstore/src/bucket/replication/replication_object_decision_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_object_decision_boundary.rs @@ -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, }; diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index f2972aae9..5da6998d7 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -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 ReplicationPool { 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 ReplicationPool { 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 ReplicationPool { 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"); diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index e8a2f3472..fbf66312a 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -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(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(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( } } -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 { + 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), + Stop, + Retry, +} + +fn delete_marker_purge_config( + result: std::result::Result, 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(dobj: &DeletedObjectReplicationInfo, storage: Arc) { @@ -1851,10 +1952,10 @@ async fn replicate_force_delete_to_targets(dobj: &Deleted } async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_client: Arc) -> 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(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!( diff --git a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs index a2992bc5d..3b74d6f25 100644 --- a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs @@ -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, version_purge: bool) -> Option { + 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, @@ -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()); diff --git a/crates/replication/src/config.rs b/crates/replication/src/config.rs index 9dea6d7eb..3aee7ff82 100644 --- a/crates/replication/src/config.rs +++ b/crates/replication/src/config.rs @@ -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; } +pub fn delete_replication_target_arns(config: &ReplicationConfiguration, object_name: &str, replica: bool) -> HashSet { + 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" + ); + } } diff --git a/crates/replication/src/delete.rs b/crates/replication/src/delete.rs index 11f248e3a..2a4f4c9d8 100644 --- a/crates/replication/src/delete.rs +++ b/crates/replication/src/delete.rs @@ -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) -> 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)); + } } diff --git a/crates/replication/src/filemeta.rs b/crates/replication/src/filemeta.rs index a802c05c2..367c94ca0 100644 --- a/crates/replication/src/filemeta.rs +++ b/crates/replication/src/filemeta.rs @@ -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, } } diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index e6759cc49..ca8e64838 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -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::{ diff --git a/crates/replication/src/mrf.rs b/crates/replication/src/mrf.rs index 08c7c7095..81c3c94f7 100644 --- a/crates/replication/src/mrf.rs +++ b/crates/replication/src/mrf.rs @@ -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, + pub delete_marker_version_id: Option, + pub delete_marker: bool, +} + +impl MrfReplicateEntry { + pub fn delete_parts_for_replay(&self) -> Option { + 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> { 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> { #[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); diff --git a/crates/replication/src/operation.rs b/crates/replication/src/operation.rs index d587b1bee..c64483a24 100644 --- a/crates/replication/src/operation.rs +++ b/crates/replication/src/operation.rs @@ -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, + pub delete_marker_version_id: Option, +} + +pub fn delete_replication_parts( + source_delete_marker: bool, + source_version_id: Option, + version_purge: bool, +) -> Option { + 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)); diff --git a/crates/replication/src/queue.rs b/crates/replication/src/queue.rs index 8d34a1d58..a762c436d 100644 --- a/crates/replication/src/queue.rs +++ b/crates/replication/src/queue.rs @@ -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] diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index bf1ff1f07..cdae6b836 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -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, }, ], diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 8d46feadd..bec995a2f 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -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 {