From 9a142fb1234860576f06a39b6651a96cb3760f71 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Thu, 2 Jul 2026 22:11:29 +0800 Subject: [PATCH] refactor(replication): move delete schedule decisions (#4199) --- crates/replication/src/lib.rs | 8 +- crates/replication/src/operation.rs | 127 +++++++++++++++++++++++++++- rustfs/src/app/bucket_usecase.rs | 2 +- rustfs/src/app/object_usecase.rs | 46 +++++----- 4 files changed, 149 insertions(+), 34 deletions(-) diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index 840b6cf1a..2cbb11b50 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -38,9 +38,11 @@ pub use object::{ replication_etags_match, target_is_newer_than_source_null_version, }; pub use operation::{ - MustReplicateOptions, ReplicationDeleteSource, ReplicationDeleteStateSource, ReplicationResyncTargetObject, - delete_replication_missing_source_decision, delete_replication_object_opts, delete_replication_state_from_config, - heal_uses_delete_replication_path, is_ssec_encrypted, resync_target_for_object, + 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, }; pub use queue::{ ReplicationHealQueueAction, ReplicationHealQueueResult, ReplicationHealResyncDeletes, ReplicationOperation, diff --git a/crates/replication/src/operation.rs b/crates/replication/src/operation.rs index c95bbcb8b..893b834fa 100644 --- a/crates/replication/src/operation.rs +++ b/crates/replication/src/operation.rs @@ -154,6 +154,62 @@ pub fn delete_replication_state_from_config( Some(state) } +#[derive(Debug, Clone, Copy)] +pub struct ReplicationDeleteScheduleInput<'a> { + pub replication_request: bool, + pub version_id_requested: bool, + pub source_delete_marker: bool, + pub source_replication_status: &'a ReplicationStatusType, + pub source_version_purge_status: &'a VersionPurgeStatusType, + 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 +} + +pub fn should_schedule_delete_replication(input: ReplicationDeleteScheduleInput<'_>) -> bool { + if input.replication_request { + 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); + } + + 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) +} + +pub fn delete_replication_version_id( + source_delete_marker: bool, + source_version_id: Option, + deleted_delete_marker_version: bool, +) -> Option { + if source_delete_marker && !deleted_delete_marker_version { + None + } else { + source_version_id + } +} + +pub fn should_use_existing_delete_replication_source( + replication_request: bool, + deleted_delete_marker: bool, + has_existing_source: bool, +) -> bool { + replication_request && deleted_delete_marker && has_existing_source +} + +pub fn should_use_existing_delete_replication_info(version_id_requested: bool, delete_marker_request: bool) -> bool { + version_id_requested && !delete_marker_request +} + pub fn heal_uses_delete_replication_path(delete_marker: bool, version_purge_status: &VersionPurgeStatusType) -> bool { delete_marker || !version_purge_status.is_empty() } @@ -246,9 +302,11 @@ pub fn resync_target_for_object( #[cfg(test)] mod tests { use super::{ - MustReplicateOptions, ReplicationDeleteSource, ReplicationDeleteStateSource, ReplicationResyncTargetObject, - delete_replication_missing_source_decision, delete_replication_object_opts, delete_replication_state_from_config, - heal_uses_delete_replication_path, is_ssec_encrypted, resync_target_for_object, + 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, }; use crate::{ReplicationStatusType, ReplicationType, VersionPurgeStatusType, target_reset_header}; use rustfs_storage_api::ObjectToDelete; @@ -444,6 +502,69 @@ mod tests { assert!(state.purge_targets.contains_key(arn)); } + #[test] + fn delete_replication_schedule_skips_replica_requests() { + assert!(!should_schedule_delete_replication(ReplicationDeleteScheduleInput { + replication_request: true, + version_id_requested: true, + source_delete_marker: true, + source_replication_status: &ReplicationStatusType::Completed, + source_version_purge_status: &VersionPurgeStatusType::Empty, + deleted_delete_marker_version: true, + })); + } + + #[test] + fn delete_replication_schedule_keeps_marker_and_version_purges() { + 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, + deleted_delete_marker_version: true, + })); + assert!(should_schedule_delete_replication(ReplicationDeleteScheduleInput { + replication_request: false, + version_id_requested: true, + source_delete_marker: false, + source_replication_status: &ReplicationStatusType::Completed, + source_version_purge_status: &VersionPurgeStatusType::Empty, + 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, + deleted_delete_marker_version: false, + })); + } + + #[test] + fn delete_replication_version_id_splits_marker_creation_and_purge() { + let version_id = Uuid::new_v4(); + + assert_eq!(delete_replication_version_id(true, Some(version_id), false), None); + assert_eq!(delete_replication_version_id(true, Some(version_id), true), Some(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)); + assert!(!should_use_existing_delete_replication_source(false, true, true)); + assert!(!should_use_existing_delete_replication_source(true, false, true)); + assert!(!should_use_existing_delete_replication_source(true, true, false)); + } + + #[test] + fn existing_delete_replication_info_is_limited_to_version_delete_requests() { + assert!(should_use_existing_delete_replication_info(true, false)); + assert!(!should_use_existing_delete_replication_info(true, true)); + assert!(!should_use_existing_delete_replication_info(false, false)); + } + #[test] fn resync_target_includes_object_at_reset_before_boundary() { let reset_before = OffsetDateTime::UNIX_EPOCH + Duration::seconds(30); diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index 0ccf3f32a..07ca12f86 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -261,7 +261,7 @@ fn validate_replication_config_targets(targets: &BucketTargets, config: &Replica } rustfs_replication::ReplicationTargetValidationError::StaleTarget => "replication config has a stale target", }; - Err(S3Error::with_message(S3ErrorCode::InvalidRequest, message.to_string())) + Err(S3Error::with_message(S3ErrorCode::InvalidRequest, message)) } } } diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 73b4e4837..20523223e 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -1498,24 +1498,14 @@ fn should_schedule_delete_replication( replication_source: &ObjectInfo, deleted_delete_marker_version: bool, ) -> bool { - if opts.replication_request { - return false; - } - - if opts.version_id.is_some() && !deleted_delete_marker_version && !replication_source.delete_marker { - return matches!( - replication_source.replication_status, - ReplicationStatusType::Replica - | ReplicationStatusType::Pending - | ReplicationStatusType::Completed - | ReplicationStatusType::Failed - ); - } - - replication_source.replication_status == ReplicationStatusType::Replica - || replication_source.replication_status == ReplicationStatusType::Pending - || replication_source.version_purge_status == VersionPurgeStatusType::Pending - || (deleted_delete_marker_version && replication_source.replication_status == ReplicationStatusType::Completed) + rustfs_replication::should_schedule_delete_replication(rustfs_replication::ReplicationDeleteScheduleInput { + replication_request: opts.replication_request, + version_id_requested: opts.version_id.is_some(), + source_delete_marker: replication_source.delete_marker, + source_replication_status: &replication_source.replication_status, + source_version_purge_status: &replication_source.version_purge_status, + deleted_delete_marker_version, + }) } async fn should_schedule_replica_delete_replication( @@ -1531,15 +1521,15 @@ async fn should_schedule_replica_delete_replication( } fn delete_replication_version_id(replication_source: &ObjectInfo, deleted_delete_marker_version: bool) -> Option { - if replication_source.delete_marker && !deleted_delete_marker_version { - None - } else { - replication_source.version_id - } + rustfs_replication::delete_replication_version_id( + replication_source.delete_marker, + replication_source.version_id, + deleted_delete_marker_version, + ) } fn should_use_existing_delete_replication_info(opts: &ObjectOptions) -> bool { - opts.version_id.is_some() && !opts.delete_marker + rustfs_replication::should_use_existing_delete_replication_info(opts.version_id.is_some(), opts.delete_marker) } fn internal_object_info_lookup_opts(mut opts: ObjectOptions) -> ObjectOptions { @@ -1576,9 +1566,11 @@ fn delete_replication_state_source<'a>( existing_object_info: Option<&'a ObjectInfo>, deleted_object_info: &'a ObjectInfo, ) -> &'a ObjectInfo { - if opts.replication_request - && deleted_object_info.delete_marker - && let Some(existing) = existing_object_info + if rustfs_replication::should_use_existing_delete_replication_source( + opts.replication_request, + deleted_object_info.delete_marker, + existing_object_info.is_some(), + ) && let Some(existing) = existing_object_info { return existing; }