From f31cd4b71652c9d3eddc26cf398b0693116d0696 Mon Sep 17 00:00:00 2001 From: LeonWang0735 Date: Sat, 21 Feb 2026 20:12:05 +0800 Subject: [PATCH] fix(replication): replicate delete all versions to targets (#1898) Co-authored-by: loverustfs --- .../replication/replication_resyncer.rs | 188 ++++++++++++++++++ crates/ecstore/src/store_api.rs | 1 + rustfs/src/storage/ecfs.rs | 30 ++- 3 files changed, 217 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index e21db787d..ec434de26 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -1193,6 +1193,11 @@ pub async fn must_replicate(bucket: &str, object: &str, mopts: MustReplicateOpti } pub async fn replicate_delete(dobj: DeletedObjectReplicationInfo, storage: Arc) { + if dobj.delete_object.force_delete { + replicate_force_delete_to_targets(&dobj, storage).await; + return; + } + let bucket = dobj.bucket.clone(); let version_id = if let Some(version_id) = &dobj.delete_object.delete_marker_version_id { Some(version_id.to_owned()) @@ -1481,6 +1486,189 @@ pub async fn replicate_delete(dobj: DeletedObjectReplicationInfo, } } +async fn replicate_force_delete_to_targets(dobj: &DeletedObjectReplicationInfo, storage: Arc) { + let bucket = &dobj.bucket; + let object_name = &dobj.delete_object.object_name; + + let rcfg = match get_replication_config(bucket).await { + Ok(Some(config)) => config, + Ok(None) => { + warn!("replicate force-delete: no replication config for bucket:{}", bucket); + send_event(EventArgs { + event_name: EventName::ObjectReplicationNotTracked.as_ref().to_string(), + bucket_name: bucket.clone(), + object: ObjectInfo { + bucket: bucket.clone(), + name: object_name.clone(), + ..Default::default() + }, + user_agent: "Internal: [Replication]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + return; + } + Err(err) => { + warn!("replicate force-delete: replication config error bucket:{} error:{}", bucket, err); + send_event(EventArgs { + event_name: EventName::ObjectReplicationNotTracked.as_ref().to_string(), + bucket_name: bucket.clone(), + object: ObjectInfo { + bucket: bucket.clone(), + name: object_name.clone(), + ..Default::default() + }, + user_agent: "Internal: [Replication]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + return; + } + }; + + let ns_lock = match storage + .new_ns_lock(bucket, format!("/[replicate]/{}", object_name).as_str()) + .await + { + Ok(ns_lock) => ns_lock, + Err(e) => { + warn!( + "replicate force-delete: failed to get ns lock bucket:{} object:{} error:{}", + bucket, object_name, e + ); + send_event(EventArgs { + event_name: EventName::ObjectReplicationNotTracked.as_ref().to_string(), + bucket_name: bucket.clone(), + object: ObjectInfo { + bucket: bucket.clone(), + name: object_name.clone(), + ..Default::default() + }, + user_agent: "Internal: [Replication]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + return; + } + }; + + let _lock_guard = match ns_lock.get_write_lock(get_lock_acquire_timeout()).await { + Ok(guard) => guard, + Err(e) => { + warn!( + "replicate force-delete: failed to get write lock bucket:{} object:{} error:{}", + bucket, object_name, e + ); + send_event(EventArgs { + event_name: EventName::ObjectReplicationNotTracked.as_ref().to_string(), + bucket_name: bucket.clone(), + object: ObjectInfo { + bucket: bucket.clone(), + name: object_name.clone(), + ..Default::default() + }, + user_agent: "Internal: [Replication]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + return; + } + }; + + 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 mut join_set = JoinSet::new(); + + for arn in tgt_arns { + let Some(tgt_client) = BucketTargetSys::get().get_remote_target_client(bucket, &arn).await else { + warn!("replicate force-delete: failed to get target client bucket:{} arn:{}", bucket, arn); + send_event(EventArgs { + event_name: EventName::ObjectReplicationNotTracked.as_ref().to_string(), + bucket_name: bucket.clone(), + object: ObjectInfo { + bucket: bucket.clone(), + name: object_name.clone(), + ..Default::default() + }, + user_agent: "Internal: [Replication]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + continue; + }; + + let bucket = bucket.clone(); + let object_name = object_name.clone(); + + join_set.spawn(async move { + if BucketTargetSys::get().is_offline(&tgt_client.to_url()).await { + error!("replicate force-delete: target offline bucket:{} arn:{}", bucket, tgt_client.arn); + send_event(EventArgs { + event_name: EventName::ObjectReplicationFailed.as_ref().to_string(), + bucket_name: bucket.clone(), + object: ObjectInfo { + bucket: bucket.clone(), + name: object_name.clone(), + ..Default::default() + }, + user_agent: "Internal: [Replication]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + return; + } + + if let Err(e) = tgt_client + .remove_object( + &tgt_client.bucket, + &object_name, + None, + RemoveObjectOptions { + force_delete: true, + governance_bypass: false, + replication_delete_marker: false, + replication_mtime: None, + replication_status: ReplicationStatusType::Replica, + replication_request: true, + replication_validity_check: false, + }, + ) + .await + { + error!( + "replicate force-delete failed bucket:{} object:{} arn:{} error:{}", + bucket, object_name, tgt_client.arn, e + ); + send_event(EventArgs { + event_name: EventName::ObjectReplicationFailed.as_ref().to_string(), + bucket_name: bucket.clone(), + object: ObjectInfo { + bucket: bucket.clone(), + name: object_name.clone(), + ..Default::default() + }, + user_agent: "Internal: [Replication]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + } + }); + } + + while let Some(result) = join_set.join_next().await { + if let Err(e) = result { + error!("replicate force-delete task panicked: {}", e); + } + } +} + 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() diff --git a/crates/ecstore/src/store_api.rs b/crates/ecstore/src/store_api.rs index cbb696b49..dad8d4a92 100644 --- a/crates/ecstore/src/store_api.rs +++ b/crates/ecstore/src/store_api.rs @@ -1287,6 +1287,7 @@ pub struct DeletedObject { // to support delete marker replication pub replication_state: Option, pub found: bool, + pub force_delete: bool, } impl DeletedObject { diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 5c8633d4c..8b83d1053 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -78,8 +78,8 @@ use rustfs_ecstore::{ policy_sys::PolicySys, quota::QuotaOperation, replication::{ - DeletedObjectReplicationInfo, check_replicate_delete, get_must_replicate_options, must_replicate, - schedule_replication, schedule_replication_delete, + DeletedObjectReplicationInfo, ObjectOpts, ReplicationConfigurationExt, check_replicate_delete, + get_must_replicate_options, must_replicate, schedule_replication, schedule_replication_delete, }, tagging::{decode_tags, encode_tags}, utils::serialize, @@ -1437,6 +1437,8 @@ impl S3 for FS { // } } + let is_force_delete = opts.delete_prefix; + let Some(store) = new_object_layer_fn() else { return Err(not_initialized_error()); }; @@ -1501,6 +1503,30 @@ impl S3 for FS { .await; }); + if is_force_delete + && !replica + && let Ok((rcfg, _)) = metadata_sys::get_replication_config(&bucket).await + { + let tgt_arns = rcfg.filter_target_arns(&ObjectOpts { + name: key.clone(), + ..Default::default() + }); + for arn in tgt_arns { + schedule_replication_delete(DeletedObjectReplicationInfo { + delete_object: rustfs_ecstore::store_api::DeletedObject { + object_name: key.clone(), + force_delete: true, + ..Default::default() + }, + bucket: bucket.clone(), + target_arn: arn, + event_type: REPLICATE_INCOMING_DELETE.to_string(), + ..Default::default() + }) + .await; + } + } + if obj_info.name.is_empty() { return Ok(S3Response::with_status(DeleteObjectOutput::default(), StatusCode::NO_CONTENT)); }