From ec67884f8df1c40871ce086d16892d2224e9d36d Mon Sep 17 00:00:00 2001 From: cxymds Date: Sun, 2 Aug 2026 23:53:53 +0800 Subject: [PATCH] fix(replication): preserve durable MRF delete admission (#5643) --- .../bucket/replication/replication_pool.rs | 63 +++++++++----- .../replication/replication_resyncer.rs | 23 +++-- crates/replication/src/delete.rs | 87 +++++++++++++++++-- crates/replication/src/filemeta.rs | 7 ++ crates/replication/src/mrf.rs | 34 ++++++++ rustfs/src/admin/handlers/replication.rs | 3 + rustfs/src/app/object_usecase.rs | 26 ++++-- 7 files changed, 196 insertions(+), 47 deletions(-) diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 3da1392d7..472f3f9b5 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -297,6 +297,7 @@ where retry_count: 0, size: entry_size, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -837,11 +838,7 @@ impl ReplicationPool { /// Queues a replica delete task pub async fn queue_replica_delete_task(&self, doi: DeletedObjectReplicationInfo) -> ReplicationQueueAdmission { - let target_arns = if doi.target_arn.is_empty() { - Vec::new() - } else { - vec![doi.target_arn.clone()] - }; + let target_arns = doi.admitted_target_arns(); let ch = self .worker_queue_channel(&doi.op_type, &doi.bucket, &doi.delete_object.object_name, 0) .await; @@ -1027,6 +1024,7 @@ impl ReplicationPool { delete_marker_version_id: entry.delete_marker_version_id, delete_marker: entry.delete_marker, delete_marker_mtime, + force_delete: entry.force_delete, replication_state: Some(rstate), ..Default::default() }, @@ -1627,18 +1625,13 @@ impl ReplicationBacklogGuard { } fn for_delete(stats: Arc, delete: &DeletedObjectReplicationInfo) -> Self { - let target_arns = if delete.target_arn.is_empty() { - Vec::new() - } else { - vec![delete.target_arn.clone()] - }; Self { stats, bucket: delete.bucket.clone(), size: 0, is_delete_marker: true, op_type: delete.op_type, - target_arns, + target_arns: delete.admitted_target_arns(), } } } @@ -1958,18 +1951,24 @@ pub(crate) async fn schedule_replication_delete(dv: DeletedObjectReplicationInfo let _ = pool.queue_replica_delete_task(dv.clone()).await; } - if let (Some(rs), Some(stats)) = (dv.delete_object.replication_state, runtime_sources::replication_stats()) { - for k in rs.targets.keys() { - let ri = ReplicatedTargetInfo { - arn: k.clone(), - size: 0, - duration: Duration::default(), - op_type: ReplicationType::Delete, - ..Default::default() - }; - stats - .update(&dv.bucket, &ri, ReplicationStatusType::Pending, ReplicationStatusType::Empty) - .await; + if let Some(stats) = runtime_sources::replication_stats() { + let target_arns = dv.admitted_target_arns(); + if let Some(rs) = dv.delete_object.replication_state.as_ref() { + for k in target_arns + .iter() + .filter(|target_arn| rs.targets.contains_key(*target_arn) || rs.purge_targets.contains_key(*target_arn)) + { + let ri = ReplicatedTargetInfo { + arn: k.clone(), + size: 0, + duration: Duration::default(), + op_type: ReplicationType::Delete, + ..Default::default() + }; + stats + .update(&dv.bucket, &ri, ReplicationStatusType::Pending, ReplicationStatusType::Empty) + .await; + } } } } @@ -3089,6 +3088,7 @@ mod tests { retry_count: 1, size: 1, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3186,6 +3186,7 @@ mod tests { retry_count: 1, size: 1, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3217,6 +3218,7 @@ mod tests { retry_count: 1, size: 2048, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3271,6 +3273,7 @@ mod tests { retry_count: 1, size: 1024, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3293,6 +3296,7 @@ mod tests { retry_count: 1, size: 1024, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3405,6 +3409,7 @@ mod tests { retry_count: 3, size: 1024, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3439,6 +3444,7 @@ mod tests { retry_count: 0, size: 0, op: MrfOpKind::Delete, + force_delete: false, delete_marker_version_id: Some(dm_vid), delete_marker: true, delete_marker_mtime: Some(mtime_nanos), @@ -3473,6 +3479,7 @@ mod tests { retry_count: 0, size: 0, op: MrfOpKind::Delete, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3502,6 +3509,7 @@ mod tests { retry_count: 1, size: 512, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3514,6 +3522,7 @@ mod tests { retry_count: 0, size: 0, op: MrfOpKind::Delete, + force_delete: false, delete_marker_version_id: Some(del_dm_vid), delete_marker: true, delete_marker_mtime: None, @@ -3544,6 +3553,7 @@ mod tests { retry_count: 0, size: 0, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3559,6 +3569,7 @@ mod tests { retry_count: 0, size: 0, op: MrfOpKind::Delete, + force_delete: false, delete_marker_version_id: Some(Uuid::new_v4()), delete_marker: true, delete_marker_mtime: None, @@ -3575,6 +3586,7 @@ mod tests { retry_count: 0, size: 0, op: MrfOpKind::default(), + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3640,6 +3652,7 @@ mod tests { retry_count: 1, size: 512, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3687,6 +3700,7 @@ mod tests { retry_count: 0, size: 1024, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3699,6 +3713,7 @@ mod tests { retry_count: 0, size: 512, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3711,6 +3726,7 @@ mod tests { retry_count: 0, size: 256, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -3765,6 +3781,7 @@ mod tests { retry_count: 0, size: -1, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index c65160e7b..fc5eea204 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -1383,6 +1383,7 @@ pub async fn replicate_delete(dobj: DeletedObjectReplicat let mut join_set = JoinSet::new(); // Process each target + let target_arns = dobj.admitted_target_arns(); for tgt_entry in dsc.targets_map.values() { // Skip targets that should not be replicated if !tgt_entry.replicate { @@ -1390,7 +1391,7 @@ pub async fn replicate_delete(dobj: DeletedObjectReplicat } // If dobj.TargetArn is not empty string, this is a case of specific target being re-synced. - if !dobj.target_arn.is_empty() && dobj.target_arn != tgt_entry.arn { + if !target_arns.is_empty() && !target_arns.iter().any(|arn| arn == &tgt_entry.arn) { continue; } @@ -1617,7 +1618,8 @@ async fn replicate_delete_marker_purge_to_targets(bucket: &str, dobj: &DeletedOb if !tgt_entry.replicate { continue; } - if !dobj.target_arn.is_empty() && dobj.target_arn != tgt_entry.arn { + let target_arns = dobj.admitted_target_arns(); + if !target_arns.is_empty() && !target_arns.iter().any(|arn| arn == &tgt_entry.arn) { continue; } let Some(tgt_client) = ReplicationTargetStore::remote_target_client(bucket, &tgt_entry.arn).await else { @@ -1747,13 +1749,16 @@ async fn replicate_force_delete_to_targets(dobj: &Deleted } }; - let tgt_arns = if !dobj.target_arn.is_empty() { - vec![dobj.target_arn.clone()] - } else { - rcfg.filter_target_arns(&ObjectOpts { - name: object_name.clone(), - ..Default::default() - }) + let tgt_arns = { + let admitted = dobj.admitted_target_arns(); + if admitted.is_empty() { + rcfg.filter_target_arns(&ObjectOpts { + name: object_name.clone(), + ..Default::default() + }) + } else { + admitted + } }; let mut join_set = JoinSet::new(); diff --git a/crates/replication/src/delete.rs b/crates/replication/src/delete.rs index 36d6c6f9c..e4fbfd7ba 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, ReplicationState, ReplicationType, ReplicationWorkerOperation}; #[derive(Debug, Clone, Default)] pub struct DeletedObjectReplicationInfo { @@ -27,6 +27,24 @@ pub struct DeletedObjectReplicationInfo { pub target_arn: String, } +impl DeletedObjectReplicationInfo { + pub fn admitted_target_arns(&self) -> Vec { + if !self.target_arn.is_empty() { + return vec![self.target_arn.clone()]; + } + + let mut target_arns = self + .delete_object + .replication_state + .as_ref() + .map(admitted_target_arns_from_replication_state) + .unwrap_or_default(); + target_arns.sort(); + target_arns.dedup(); + target_arns + } +} + impl ReplicationWorkerOperation for DeletedObjectReplicationInfo { fn as_any(&self) -> &dyn Any { self @@ -40,6 +58,7 @@ impl ReplicationWorkerOperation for DeletedObjectReplicationInfo { retry_count: 0, size: 0, op: MrfOpKind::Delete, + force_delete: self.delete_object.force_delete, delete_marker_version_id: self.delete_object.delete_marker_version_id, delete_marker: self.delete_object.delete_marker, // Persist the original delete-marker mtime as Unix nanoseconds so replay after a @@ -49,11 +68,7 @@ impl ReplicationWorkerOperation for DeletedObjectReplicationInfo { .delete_object .delete_marker_mtime .and_then(|t| i64::try_from(t.unix_timestamp_nanos()).ok()), - target_arns: if self.target_arn.is_empty() { - Vec::new() - } else { - vec![self.target_arn.clone()] - }, + target_arns: self.admitted_target_arns(), } } @@ -86,18 +101,28 @@ pub fn should_retry_delete_marker_purge(dobj: &DeletedObject) -> bool { dobj.delete_marker_version_id.is_some() } +fn admitted_target_arns_from_replication_state(state: &ReplicationState) -> Vec { + let mut target_arns = state.targets.keys().cloned().collect::>(); + target_arns.extend(state.purge_targets.keys().cloned()); + target_arns +} + pub fn is_retryable_delete_replication_head_error(is_not_found: bool, code: Option<&str>) -> bool { !(is_not_found || matches!(code, Some("MethodNotAllowed" | "405"))) } #[cfg(test)] mod tests { + use std::collections::HashMap; + use super::{ DeletedObjectReplicationInfo, is_retryable_delete_replication_head_error, is_version_delete_replication, should_retry_delete_marker_purge, }; use crate::storage_api::DeletedObject; - use crate::{MrfOpKind, ReplicationType, ReplicationWorkerOperation}; + use crate::{ + MrfOpKind, ReplicationState, ReplicationStatusType, ReplicationType, ReplicationWorkerOperation, VersionPurgeStatusType, + }; use uuid::Uuid; #[test] @@ -125,6 +150,7 @@ mod tests { assert_eq!(entry.bucket, "bucket"); assert_eq!(entry.object, "object"); assert_eq!(entry.version_id, Some(version_id)); + assert!(!entry.force_delete); assert_eq!(entry.delete_marker_version_id, Some(delete_marker_version_id)); assert_eq!(entry.op, MrfOpKind::Delete); assert!(entry.delete_marker); @@ -156,6 +182,53 @@ mod tests { assert_eq!(info.to_mrf_entry().delete_marker_mtime, None); assert!(info.to_mrf_entry().target_arns.is_empty()); + assert!(!info.to_mrf_entry().force_delete); + } + + #[test] + fn deleted_object_replication_info_uses_explicit_target_over_replay_state() { + let info = DeletedObjectReplicationInfo { + bucket: "bucket".to_string(), + delete_object: DeletedObject { + object_name: "object".to_string(), + replication_state: Some(ReplicationState { + targets: HashMap::from([("arn:target-state".to_string(), ReplicationStatusType::Pending)]), + purge_targets: HashMap::from([("arn:purge-state".to_string(), VersionPurgeStatusType::Pending)]), + ..Default::default() + }), + ..Default::default() + }, + target_arn: "arn:target-explicit".to_string(), + ..Default::default() + }; + + assert_eq!(info.to_mrf_entry().target_arns, vec!["arn:target-explicit".to_string()]); + } + + #[test] + fn deleted_object_replication_info_serializes_replay_targets_when_target_is_implicit() { + let info = DeletedObjectReplicationInfo { + bucket: "bucket".to_string(), + delete_object: DeletedObject { + object_name: "object".to_string(), + replication_state: Some(ReplicationState { + targets: HashMap::from([("arn:target-b".to_string(), ReplicationStatusType::Pending)]), + purge_targets: HashMap::from([ + ("arn:target-a".to_string(), VersionPurgeStatusType::Pending), + ("arn:target-b".to_string(), VersionPurgeStatusType::Complete), + ]), + ..Default::default() + }), + ..Default::default() + }, + ..Default::default() + }; + + assert_eq!( + info.to_mrf_entry().target_arns, + vec!["arn:target-a".to_string(), "arn:target-b".to_string()], + "MRF deletes must preserve the admitted target identities in stable order" + ); } #[test] diff --git a/crates/replication/src/filemeta.rs b/crates/replication/src/filemeta.rs index 20fc30883..eabb843a1 100644 --- a/crates/replication/src/filemeta.rs +++ b/crates/replication/src/filemeta.rs @@ -579,6 +579,11 @@ pub struct MrfReplicateEntry { #[serde(rename = "op", default)] pub op: MrfOpKind, + // For delete entries: whether the source operation was a force-delete. Old files lack + // this key; default=false preserves the pre-existing replay contract. + #[serde(rename = "forceDelete", default)] + pub force_delete: bool, + // For delete entries: the delete-marker version id (distinct from version_id, which is // the version being purged). Old files lack this; default=None is correct. #[serde(rename = "deleteMarkerVersionID", skip_serializing_if = "Option::is_none", default)] @@ -838,6 +843,7 @@ impl ReplicationWorkerOperation for ReplicateObjectInfo { } else { MrfOpKind::Object }, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -897,6 +903,7 @@ impl ReplicateObjectInfo { } else { MrfOpKind::Object }, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, diff --git a/crates/replication/src/mrf.rs b/crates/replication/src/mrf.rs index 45faaf301..dc7ecefa0 100644 --- a/crates/replication/src/mrf.rs +++ b/crates/replication/src/mrf.rs @@ -68,6 +68,7 @@ mod tests { retry_count: 1, size: 1024, op: MrfOpKind::Metadata, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -80,6 +81,7 @@ mod tests { retry_count: 2, size: 1024, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -92,6 +94,7 @@ mod tests { retry_count: 0, size: 0, op: MrfOpKind::Delete, + force_delete: true, delete_marker_version_id: Some(del_vid), delete_marker: true, delete_marker_mtime: Some(1_705_312_200_123_456_789), @@ -107,11 +110,14 @@ mod tests { assert_eq!(decoded[0].op, MrfOpKind::Metadata); assert_eq!(decoded[0].target_arns, vec!["arn:target-a".to_string()]); assert_eq!(decoded[0].delete_marker_mtime, None); + assert!(!decoded[0].force_delete); assert_eq!(decoded[1].op, MrfOpKind::Object); + assert!(!decoded[1].force_delete); assert_eq!(decoded[1].target_arns, vec!["arn:target-a".to_string(), "arn:target-b".to_string()]); assert_eq!(decoded[1].delete_marker_mtime, None); assert_eq!(decoded[2].delete_marker_version_id, Some(del_vid)); assert_eq!(decoded[2].op, MrfOpKind::Delete); + assert!(decoded[2].force_delete); assert_eq!(decoded[2].target_arns, vec!["arn:target-a".to_string()]); assert!(decoded[2].delete_marker); assert_eq!( @@ -149,11 +155,39 @@ mod tests { assert_eq!(decoded[0].size, 100); assert_eq!(decoded[0].op, MrfOpKind::Object); assert!(decoded[0].target_arns.is_empty()); + assert!(!decoded[0].force_delete); // 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); } + #[test] + fn mrf_file_defaults_missing_force_delete_to_false() { + let mut payload = Vec::new(); + rmp::encode::write_array_len(&mut payload, 1).expect("array len should encode"); + rmp::encode::write_map_len(&mut payload, 5).expect("map len should encode"); + rmp::encode::write_str(&mut payload, "bucket").expect("bucket key should encode"); + rmp::encode::write_str(&mut payload, "old-bucket").expect("bucket value should encode"); + rmp::encode::write_str(&mut payload, "object").expect("object key should encode"); + rmp::encode::write_str(&mut payload, "old-key").expect("object value should encode"); + rmp::encode::write_str(&mut payload, "retryCount").expect("retry key should encode"); + rmp::encode::write_i32(&mut payload, 1).expect("retry value should encode"); + rmp::encode::write_str(&mut payload, "size").expect("size key should encode"); + rmp::encode::write_i64(&mut payload, 0).expect("size value should encode"); + rmp::encode::write_str(&mut payload, "op").expect("op key should encode"); + rmp::encode::write_str(&mut payload, "delete").expect("op value should encode"); + + let mut data = Vec::with_capacity(4 + payload.len()); + data.extend_from_slice(&MRF_META_FORMAT.to_le_bytes()); + data.extend_from_slice(&MRF_META_VERSION.to_le_bytes()); + data.extend_from_slice(&payload); + + let decoded = decode_mrf_file(&data).expect("MRF payload should decode"); + + assert_eq!(decoded[0].op, MrfOpKind::Delete); + assert!(!decoded[0].force_delete); + } + #[test] fn mrf_file_rejects_invalid_header() { let mut data = Vec::new(); diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index 3cc890b91..b5136762f 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -1120,6 +1120,7 @@ mod tests { retry_count: 0, size: 250, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -1132,6 +1133,7 @@ mod tests { retry_count: 0, size: 999, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, @@ -1190,6 +1192,7 @@ mod tests { retry_count: 0, size: 250, op: MrfOpKind::Object, + force_delete: false, delete_marker_version_id: None, delete_marker: false, delete_marker_mtime: None, diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 59abd3fde..9f9f88481 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -7288,16 +7288,26 @@ impl DefaultObjectUsecase { if obj_info.name.is_empty() { if replicate_force_delete { - schedule_replication_delete( - StorageDeletedObject { - object_name: key.clone(), - force_delete: true, + let mut delete_object = StorageDeletedObject { + object_name: key.clone(), + force_delete: true, + ..Default::default() + }; + if let Some(replication_state) = delete_replication_state_from_config( + delete_config_snapshot + .replication_config() + .unwrap_or_else(|| unreachable!("force-delete requires a replication config")), + &ObjectInfo { + bucket: bucket.clone(), + name: key.clone(), ..Default::default() }, - bucket.clone(), - REPLICATE_INCOMING_DELETE.to_string(), - ) - .await; + None, + false, + ) { + set_deleted_object_replication_state(&mut delete_object, &replication_state); + } + schedule_replication_delete(delete_object, bucket.clone(), REPLICATE_INCOMING_DELETE.to_string()).await; } // Prefix/force-delete returns empty ObjectInfo; still emit bucket notification so webhooks match S3 DELETE. helper = helper