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 3544ac27b..553a57e1b 100644 --- a/crates/ecstore/src/bucket/replication/replication_object_decision_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_object_decision_boundary.rs @@ -20,9 +20,9 @@ pub use rustfs_replication::{ pub(crate) use rustfs_replication::{ ReplicationDeleteSource, ReplicationMultipartPartInput, ReplicationResyncTargetObject, delete_marker_purge_mrf_entry, delete_marker_purge_version_id, delete_replication_creates_marker, delete_replication_missing_source_decision, - delete_replication_object_opts, heal_uses_delete_replication_path, is_object_lock_denied_delete, - is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome, replication_etags_match, - replication_multipart_complete_actual_size, replication_multipart_part_plan, replication_single_put_size_error, - resync_existing_delete_replication_info, resync_target_for_object, should_retry_delete_marker_purge, - single_part_replica_etag_mismatch, target_delete_version_id, + delete_replication_object_opts, delete_replication_target_version_id, heal_uses_delete_replication_path, + is_object_lock_denied_delete, is_retryable_delete_replication_head_error, is_version_delete_replication, + replicate_delete_outcome, replication_etags_match, replication_multipart_complete_actual_size, + replication_multipart_part_plan, replication_single_put_size_error, resync_existing_delete_replication_info, + resync_target_for_object, should_retry_delete_marker_purge, single_part_replica_etag_mismatch, }; diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index cf78dee70..64cbebeaa 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -32,11 +32,11 @@ use super::replication_msgp_boundary::ReplicationMsgpCodec; use super::replication_object_config::{ReplicationConfig, get_replication_config, must_replicate}; use super::replication_object_decision_boundary::{ MustReplicateOptions, ReplicationMultipartPartInput, delete_marker_purge_mrf_entry, delete_marker_purge_version_id, - delete_replication_creates_marker, heal_uses_delete_replication_path, is_object_lock_denied_delete, - is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome, replication_etags_match, - replication_multipart_complete_actual_size, replication_multipart_part_plan, replication_single_put_size_error, - resync_existing_delete_replication_info, should_retry_delete_marker_purge, single_part_replica_etag_mismatch, - target_delete_version_id, + delete_replication_creates_marker, delete_replication_target_version_id, heal_uses_delete_replication_path, + is_object_lock_denied_delete, is_retryable_delete_replication_head_error, is_version_delete_replication, + replicate_delete_outcome, replication_etags_match, replication_multipart_complete_actual_size, + replication_multipart_part_plan, replication_single_put_size_error, resync_existing_delete_replication_info, + should_retry_delete_marker_purge, single_part_replica_etag_mismatch, }; use super::replication_queue_boundary::{DeletedObjectReplicationInfo, ReplicationQueueAdmission}; use super::replication_resync_boundary::ResyncStatusType; @@ -2051,7 +2051,11 @@ pub(crate) async fn replicate_delete_with_outcome( let is_version_purge = is_version_delete_replication(&dobj.delete_object); - let requires_delayed_purge = should_retry_delete_marker_purge(&dobj.delete_object); + // The watcher exists to purge a replicated marker once the SOURCE marker + // vanishes. A version purge is that purge already (its failures reach the + // journal as a purge entry), so it must not spawn a second watcher that + // journals a duplicate intent (backlog#2290). + let requires_delayed_purge = should_retry_delete_marker_purge(&dobj.delete_object) && !is_version_purge; let (replication_status, prev_status) = if !is_version_purge { ( @@ -2761,12 +2765,6 @@ fn unavailable_delete_target_info(dobj: &DeletedObjectReplicationInfo, arn: &str } 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() - } else { - dobj.delete_object.version_id.unwrap_or_default() - }; - let mut rinfo = dobj .delete_object .replication_state @@ -2799,7 +2797,25 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli return rinfo; } - let version_id = target_delete_version_id(version_id, is_version_purge); + // Purging a replicated delete marker addresses the version the target + // assigned (recorded when the marker was created there); see + // `delete_replication_target_version_id`. A corrupt record is a failure, + // not a guess: the entry stays visible until the metadata is repaired. + let Some(version_id) = delete_replication_target_version_id(&dobj.delete_object, &tgt_client.arn) else { + warn!( + event = EVENT_DELETE_MARKER_PURGE_FAILED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, + bucket = tgt_client.bucket, + object = dobj.delete_object.object_name, + arn = %tgt_client.arn, + reason = "recorded_target_version_inconsistent", + "Replicated version purge refused: recorded target delete-marker version metadata is inconsistent" + ); + rinfo.version_purge_status = VersionPurgeStatusType::Failed; + rinfo.error = Some("recorded target delete-marker version metadata is inconsistent".to_string()); + return rinfo; + }; if dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() { match head_object_for_worker( diff --git a/crates/replication/src/delete.rs b/crates/replication/src/delete.rs index ff6fe1591..88a492533 100644 --- a/crates/replication/src/delete.rs +++ b/crates/replication/src/delete.rs @@ -253,6 +253,28 @@ pub fn delete_marker_purge_version_id( }) } +/// The version a delete replication addresses on `arn`, or `None` to refuse. +/// +/// A version purge whose purged version is a delete marker must address the +/// marker version the TARGET assigned — the recorded mapping, exactly as the +/// delayed-purge watcher does. The source-side `DELETE ?versionId=` +/// replicates as such a purge, and a generic S3 target answers a DELETE of an +/// unknown versionId with 204 while keeping its marker, so addressing it by +/// the source id reported success and left the marker behind (backlog#2290, +/// R6.1 on the VMs). Nothing recorded falls back to the source-derived id +/// (id-mirroring peers); a corrupt record refuses, as the watcher does. +pub fn delete_replication_target_version_id(dobj: &DeletedObject, arn: &str) -> Option> { + let is_version_purge = is_version_delete_replication(dobj); + if is_version_purge + && !dobj.delete_marker + && let Some(marker) = dobj.delete_marker_version_id + { + return delete_marker_purge_version_id(dobj.replication_state.as_ref(), arn, marker); + } + let source_version = dobj.delete_marker_version_id.or(dobj.version_id).unwrap_or_default(); + Some(target_delete_version_id(source_version, is_version_purge)) +} + /// Shape an exhausted purge intent as a marker-creation delete entry. Replay /// reconstructs it with `delete_marker: true`, finds the source marker gone, /// and funnels into the stale-marker branch of `replicate_delete_with_outcome` @@ -273,9 +295,9 @@ mod tests { use super::{ DeletedObjectReplicationInfo, delete_marker_purge_mrf_entry, delete_marker_purge_version_id, - delete_replication_creates_marker, is_object_lock_denied_delete, is_retryable_delete_replication_head_error, - is_version_delete_replication, replicate_delete_outcome, resync_existing_delete_replication_info, - should_retry_delete_marker_purge, target_delete_version_id, + delete_replication_creates_marker, delete_replication_target_version_id, is_object_lock_denied_delete, + is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome, + resync_existing_delete_replication_info, should_retry_delete_marker_purge, target_delete_version_id, }; use crate::storage_api::DeletedObject; use crate::{ @@ -741,4 +763,57 @@ mod tests { assert!(!is_object_lock_denied_delete(Some("InternalError"), Some("retention lookup failed"))); assert!(!is_object_lock_denied_delete(None, Some("legal hold"))); } + + fn purge_of_marker(marker: Uuid, state: Option) -> DeletedObject { + DeletedObject { + object_name: "obj".to_string(), + delete_marker: false, + delete_marker_version_id: Some(marker), + version_id: None, + replication_state: state, + ..Default::default() + } + } + + #[test] + fn delete_replication_target_version_id_addresses_recorded_marker_for_purges() { + let arn = "arn:minio:replication::generic:photos"; + let marker = Uuid::new_v4(); + let mut state = ReplicationState::default(); + state + .target_delete_marker_version_ids + .insert(arn.to_string(), "remote-marker".to_string()); + + // purge of a replicated marker: the target's own version + assert_eq!( + delete_replication_target_version_id(&purge_of_marker(marker, Some(state.clone())), arn), + Some(Some("remote-marker".to_string())) + ); + // nothing recorded for this arn: the source-derived id (id-mirroring peers) + assert_eq!( + delete_replication_target_version_id(&purge_of_marker(marker, None), arn), + Some(Some(marker.to_string())) + ); + // corrupt record: refuse instead of guessing + state.target_delete_marker_version_ids_corrupt = true; + assert_eq!(delete_replication_target_version_id(&purge_of_marker(marker, Some(state)), arn), None); + + // marker creation keeps the source id (the target mints its own on a + // versionless DELETE; the id only travels in the source header) + let creation = DeletedObject { + object_name: "obj".to_string(), + delete_marker: true, + delete_marker_version_id: Some(marker), + ..Default::default() + }; + assert_eq!(delete_replication_target_version_id(&creation, arn), Some(Some(marker.to_string()))); + // plain version purge: the source version id + let version = Uuid::new_v4(); + let purge = DeletedObject { + object_name: "obj".to_string(), + version_id: Some(version), + ..Default::default() + }; + assert_eq!(delete_replication_target_version_id(&purge, arn), Some(Some(version.to_string()))); + } } diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index 1aa20daca..295340026 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -41,9 +41,9 @@ pub use config::{ }; pub use delete::{ DeletedObjectReplicationInfo, delete_marker_purge_mrf_entry, delete_marker_purge_version_id, - delete_replication_creates_marker, is_object_lock_denied_delete, is_retryable_delete_replication_head_error, - is_version_delete_replication, replicate_delete_outcome, resync_existing_delete_replication_info, - should_retry_delete_marker_purge, target_delete_version_id, + delete_replication_creates_marker, delete_replication_target_version_id, is_object_lock_denied_delete, + is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome, + resync_existing_delete_replication_info, should_retry_delete_marker_purge, target_delete_version_id, }; pub use filemeta::{ NULL_VERSION_ID, REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, REPLICATE_HEAL, REPLICATE_HEAL_DELETE, REPLICATE_INCOMING,