From b57820a486a27ccd263f2b4a8e329d9c37615a8c Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Thu, 2 Jul 2026 10:17:48 +0800 Subject: [PATCH] refactor(replication): move config contracts into crate (#4162) --- Cargo.lock | 2 + .../ecstore/src/bucket/replication/config.rs | 384 +---------------- crates/ecstore/src/bucket/replication/mod.rs | 1 - .../replication_tagging_boundary.rs | 24 +- crates/replication/Cargo.toml | 6 +- crates/replication/src/config.rs | 392 ++++++++++++++++++ crates/replication/src/lib.rs | 6 + .../replication => replication/src}/rule.rs | 2 +- crates/replication/src/tagging.rs | 60 +++ scripts/check_architecture_migration_rules.sh | 16 + 10 files changed, 482 insertions(+), 411 deletions(-) create mode 100644 crates/replication/src/config.rs rename crates/{ecstore/src/bucket/replication => replication/src}/rule.rs (98%) create mode 100644 crates/replication/src/tagging.rs diff --git a/Cargo.lock b/Cargo.lock index cec0095d2..446dbcc44 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9897,8 +9897,10 @@ dependencies = [ "rmp", "rmp-serde", "rustfs-filemeta", + "s3s", "serde", "time", + "url", "uuid", ] diff --git a/crates/ecstore/src/bucket/replication/config.rs b/crates/ecstore/src/bucket/replication/config.rs index 807a25141..2003f970d 100644 --- a/crates/ecstore/src/bucket/replication/config.rs +++ b/crates/ecstore/src/bucket/replication/config.rs @@ -12,386 +12,4 @@ // See the License for the specific language governing permissions and // limitations under the License. -use super::replication_filemeta_boundary::ReplicationType; -use super::replication_tagging_boundary::ReplicationTagFilter; -use super::rule::ReplicationRuleExt as _; -use s3s::dto::DeleteMarkerReplicationStatus; -use s3s::dto::DeleteReplicationStatus; -use s3s::dto::Destination; -use s3s::dto::{ExistingObjectReplicationStatus, ReplicationConfiguration, ReplicationRuleStatus, ReplicationRules}; -use serde::{Deserialize, Serialize}; -use std::collections::HashSet; -use uuid::Uuid; - -#[derive(Debug, Clone, Serialize, Deserialize, Default)] -pub struct ObjectOpts { - pub name: String, - pub user_tags: String, - pub version_id: Option, - pub delete_marker: bool, - pub ssec: bool, - pub op_type: ReplicationType, - pub replica: bool, - pub existing_object: bool, - pub target_arn: String, -} - -pub trait ReplicationConfigurationExt { - fn replicate(&self, opts: &ObjectOpts) -> bool; - fn has_existing_object_replication(&self, arn: &str) -> (bool, bool); - fn filter_actionable_rules(&self, obj: &ObjectOpts) -> ReplicationRules; - fn get_destination(&self) -> Destination; - fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool; - fn filter_target_arns(&self, obj: &ObjectOpts) -> Vec; -} - -impl ReplicationConfigurationExt for ReplicationConfiguration { - /// Check whether any object-replication rules exist - fn has_existing_object_replication(&self, arn: &str) -> (bool, bool) { - let mut has_arn = false; - let arn = arn.trim(); - - for rule in &self.rules { - if rule.destination.bucket.trim() == arn || self.role.trim() == arn { - if !has_arn { - has_arn = true; - } - if let Some(status) = &rule.existing_object_replication - && status.status == ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::ENABLED) - { - return (true, true); - } - } - } - (has_arn, false) - } - - fn filter_actionable_rules(&self, obj: &ObjectOpts) -> ReplicationRules { - if obj.name.is_empty() && obj.op_type != ReplicationType::Resync && obj.op_type != ReplicationType::All { - return vec![]; - } - - let mut rules = ReplicationRules::default(); - - for rule in &self.rules { - if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) { - continue; - } - - if !obj.target_arn.is_empty() - && rule.destination.bucket.trim() != obj.target_arn.trim() - && self.role.trim() != obj.target_arn.trim() - { - continue; - } - - if obj.op_type == ReplicationType::Resync || obj.op_type == ReplicationType::All { - rules.push(rule.clone()); - continue; - } - - if let Some(status) = &rule.existing_object_replication - && obj.existing_object - && status.status == ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::DISABLED) - { - continue; - } - - if !obj.name.starts_with(rule.prefix()) { - continue; - } - - if let Some(filter) = &rule.filter { - let object_tags = ReplicationTagFilter::decode_tags_to_map(&obj.user_tags); - if filter.test_tags(&object_tags) { - rules.push(rule.clone()); - } - } else { - rules.push(rule.clone()); - } - } - - rules.sort_by(|a, b| { - if a.destination == b.destination { - a.priority.cmp(&b.priority) - } else { - std::cmp::Ordering::Equal - } - }); - - rules - } - - /// Retrieve the destination configuration - fn get_destination(&self) -> Destination { - if !self.rules.is_empty() { - self.rules[0].destination.clone() - } else { - Destination { - account: None, - bucket: "".to_string(), - encryption_configuration: None, - metrics: None, - replication_time: None, - access_control_translation: None, - storage_class: None, - } - } - } - - /// Determine whether an object should be replicated - fn replicate(&self, obj: &ObjectOpts) -> bool { - let rules = self.filter_actionable_rules(obj); - - for rule in rules.iter() { - if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) { - 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); - } - false - } - - /// Check for an active rule - /// Optionally accept a prefix - /// When recursive is true, return true if any level under the prefix has an active rule - /// Without a prefix, recursive behaves as true - fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool { - if self.rules.is_empty() { - return false; - } - - for rule in &self.rules { - if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) { - continue; - } - - if let Some(filter) = &rule.filter - && let Some(filter_prefix) = &filter.prefix - { - if !prefix.is_empty() && !filter_prefix.is_empty() { - // The provided prefix must fall within the rule prefix - if !recursive && !prefix.starts_with(filter_prefix) { - continue; - } - } - - // When recursive, skip this rule if it does not match the test prefix or hierarchy - if recursive && !rule.prefix().starts_with(prefix) && !prefix.starts_with(rule.prefix()) { - continue; - } - } - return true; - } - false - } - - /// Filter target ARNs and return a slice of the distinct values in the config - fn filter_target_arns(&self, obj: &ObjectOpts) -> Vec { - let role = self.role.trim(); - if !role.is_empty() { - return vec![role.to_string()]; - } - - let mut arns = Vec::new(); - let mut targets_map: HashSet = HashSet::new(); - let rules = self.filter_actionable_rules(obj); - - for rule in rules { - if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) { - continue; - } - - let arn = rule.destination.bucket.trim(); - if !arn.is_empty() && !targets_map.contains(arn) { - targets_map.insert(arn.to_string()); - } - } - - for arn in targets_map { - arns.push(arn); - } - arns - } -} - -#[cfg(test)] -mod tests { - use super::*; - use s3s::dto::{DeleteMarkerReplication, Destination, ExistingObjectReplication, ReplicationRule}; - - fn replication_rule(id: &str, arn: &str) -> ReplicationRule { - ReplicationRule { - delete_marker_replication: Some(DeleteMarkerReplication::default()), - delete_replication: None, - destination: Destination { - bucket: arn.to_string(), - ..Default::default() - }, - existing_object_replication: Some(ExistingObjectReplication { - status: ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::ENABLED), - }), - filter: None, - id: Some(id.to_string()), - prefix: Some(String::new()), - priority: Some(1), - source_selection_criteria: None, - status: ReplicationRuleStatus::from_static(ReplicationRuleStatus::ENABLED), - } - } - - #[test] - fn filter_target_arns_uses_role_when_role_is_present() { - let config = ReplicationConfiguration { - role: " arn:legacy:target ".to_string(), - rules: vec![ - replication_rule("rule-1", "arn:target:a"), - replication_rule("rule-2", "arn:target:b"), - ], - }; - - let arns = config.filter_target_arns(&ObjectOpts { - name: "object".to_string(), - op_type: ReplicationType::Object, - ..Default::default() - }); - - assert_eq!(arns, vec!["arn:legacy:target".to_string()]); - } - - #[test] - fn filter_target_arns_falls_back_to_role_when_destination_is_empty() { - let config = ReplicationConfiguration { - role: "arn:legacy:target".to_string(), - rules: vec![ReplicationRule { - delete_marker_replication: Some(DeleteMarkerReplication::default()), - delete_replication: None, - destination: Destination { - bucket: String::new(), - ..Default::default() - }, - existing_object_replication: Some(ExistingObjectReplication { - status: ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::ENABLED), - }), - filter: None, - id: Some("rule-1".to_string()), - prefix: Some(String::new()), - priority: Some(1), - source_selection_criteria: None, - status: ReplicationRuleStatus::from_static(ReplicationRuleStatus::ENABLED), - }], - }; - - let arns = config.filter_target_arns(&ObjectOpts { - name: "object".to_string(), - op_type: ReplicationType::Object, - ..Default::default() - }); - - assert_eq!(arns, vec!["arn:legacy:target".to_string()]); - } - - fn replication_rule_existing_object_disabled(id: &str, arn: &str) -> ReplicationRule { - ReplicationRule { - delete_marker_replication: Some(DeleteMarkerReplication::default()), - delete_replication: None, - destination: Destination { - bucket: arn.to_string(), - ..Default::default() - }, - existing_object_replication: Some(ExistingObjectReplication { - status: ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::DISABLED), - }), - filter: None, - id: Some(id.to_string()), - prefix: Some(String::new()), - priority: Some(1), - source_selection_criteria: None, - status: ReplicationRuleStatus::from_static(ReplicationRuleStatus::ENABLED), - } - } - - // Regression test for BUG-3: replicate_object was calling filter_target_arns with - // existing_object:false regardless of op_type, letting ExistingObject resync operations - // fan out to targets whose rule has ExistingObjectReplicationStatus::DISABLED. - #[test] - fn filter_target_arns_excludes_disabled_existing_object_target_for_existing_object_op() { - let config = ReplicationConfiguration { - role: String::new(), - rules: vec![ - replication_rule("rule-enabled", "arn:target:enabled"), - replication_rule_existing_object_disabled("rule-disabled", "arn:target:disabled"), - ], - }; - - let arns = config.filter_target_arns(&ObjectOpts { - name: "object".to_string(), - op_type: ReplicationType::ExistingObject, - existing_object: true, - ..Default::default() - }); - - assert_eq!(arns.len(), 1, "only the ENABLED target should be returned for ExistingObject ops"); - assert!(arns.contains(&"arn:target:enabled".to_string())); - assert!(!arns.contains(&"arn:target:disabled".to_string())); - } - - // Heal operations intentionally bypass ExistingObjectReplicationStatus — healing a past - // failure is not subject to the existing-object opt-out. - #[test] - fn filter_target_arns_includes_disabled_existing_object_target_for_heal_op() { - let config = ReplicationConfiguration { - role: String::new(), - rules: vec![ - replication_rule("rule-enabled", "arn:target:enabled"), - replication_rule_existing_object_disabled("rule-disabled", "arn:target:disabled"), - ], - }; - - let arns = config.filter_target_arns(&ObjectOpts { - name: "object".to_string(), - op_type: ReplicationType::Heal, - existing_object: false, - ..Default::default() - }); - - assert_eq!( - arns.len(), - 2, - "Heal ops must reach all targets regardless of existing_object_replication setting" - ); - assert!(arns.contains(&"arn:target:enabled".to_string())); - assert!(arns.contains(&"arn:target:disabled".to_string())); - } -} +pub use rustfs_replication::{ObjectOpts, ReplicationConfigurationExt}; diff --git a/crates/ecstore/src/bucket/replication/mod.rs b/crates/ecstore/src/bucket/replication/mod.rs index 1cf8f6a6d..61af5ea7b 100644 --- a/crates/ecstore/src/bucket/replication/mod.rs +++ b/crates/ecstore/src/bucket/replication/mod.rs @@ -34,7 +34,6 @@ mod replication_tagging_boundary; mod replication_target_boundary; mod replication_target_config_bridge; mod replication_versioning_boundary; -mod rule; mod runtime_boundary; pub use config::{ObjectOpts, ReplicationConfigurationExt}; diff --git a/crates/ecstore/src/bucket/replication/replication_tagging_boundary.rs b/crates/ecstore/src/bucket/replication/replication_tagging_boundary.rs index 69b81383f..73733598e 100644 --- a/crates/ecstore/src/bucket/replication/replication_tagging_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_tagging_boundary.rs @@ -12,26 +12,4 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::collections::HashMap; - -pub(crate) struct ReplicationTagFilter; - -impl ReplicationTagFilter { - pub(crate) fn decode_tags_to_map(tags: &str) -> HashMap { - crate::bucket::tagging::decode_tags_to_map(tags) - } -} - -#[cfg(test)] -mod tests { - use super::ReplicationTagFilter; - - #[test] - fn decode_tags_to_map_preserves_bucket_tagging_parser_behavior() { - let tags = ReplicationTagFilter::decode_tags_to_map("env=prod&encoded=a%2Fb&=ignored"); - - assert_eq!(tags.get("env").map(String::as_str), Some("prod")); - assert_eq!(tags.get("encoded").map(String::as_str), Some("a/b")); - assert!(!tags.contains_key("")); - } -} +pub(crate) use rustfs_replication::ReplicationTagFilter; diff --git a/crates/replication/Cargo.toml b/crates/replication/Cargo.toml index 4e49b89bb..9e6bf9dff 100644 --- a/crates/replication/Cargo.toml +++ b/crates/replication/Cargo.toml @@ -30,11 +30,11 @@ byteorder.workspace = true rmp.workspace = true rmp-serde.workspace = true rustfs-filemeta.workspace = true +s3s.workspace = true serde.workspace = true time.workspace = true - -[dev-dependencies] -uuid = { workspace = true, features = ["v4"] } +url.workspace = true +uuid = { workspace = true, features = ["v4", "serde"] } [lints] workspace = true diff --git a/crates/replication/src/config.rs b/crates/replication/src/config.rs new file mode 100644 index 000000000..4a3171717 --- /dev/null +++ b/crates/replication/src/config.rs @@ -0,0 +1,392 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use crate::ReplicationTagFilter; +use crate::ReplicationType; +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 serde::{Deserialize, Serialize}; +use std::collections::HashSet; +use uuid::Uuid; + +#[derive(Debug, Clone, Serialize, Deserialize, Default)] +pub struct ObjectOpts { + pub name: String, + pub user_tags: String, + pub version_id: Option, + pub delete_marker: bool, + pub ssec: bool, + pub op_type: ReplicationType, + pub replica: bool, + pub existing_object: bool, + pub target_arn: String, +} + +pub trait ReplicationConfigurationExt { + fn replicate(&self, opts: &ObjectOpts) -> bool; + fn has_existing_object_replication(&self, arn: &str) -> (bool, bool); + fn filter_actionable_rules(&self, obj: &ObjectOpts) -> ReplicationRules; + fn get_destination(&self) -> Destination; + fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool; + fn filter_target_arns(&self, obj: &ObjectOpts) -> Vec; +} + +impl ReplicationConfigurationExt for ReplicationConfiguration { + /// Check whether any object-replication rules exist + fn has_existing_object_replication(&self, arn: &str) -> (bool, bool) { + let mut has_arn = false; + let arn = arn.trim(); + + for rule in &self.rules { + if rule.destination.bucket.trim() == arn || self.role.trim() == arn { + if !has_arn { + has_arn = true; + } + if let Some(status) = &rule.existing_object_replication + && status.status == ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::ENABLED) + { + return (true, true); + } + } + } + (has_arn, false) + } + + fn filter_actionable_rules(&self, obj: &ObjectOpts) -> ReplicationRules { + if obj.name.is_empty() && obj.op_type != ReplicationType::Resync && obj.op_type != ReplicationType::All { + return vec![]; + } + + let mut rules = ReplicationRules::default(); + + for rule in &self.rules { + if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) { + continue; + } + + if !obj.target_arn.is_empty() + && rule.destination.bucket.trim() != obj.target_arn.trim() + && self.role.trim() != obj.target_arn.trim() + { + continue; + } + + if obj.op_type == ReplicationType::Resync || obj.op_type == ReplicationType::All { + rules.push(rule.clone()); + continue; + } + + if let Some(status) = &rule.existing_object_replication + && obj.existing_object + && status.status == ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::DISABLED) + { + continue; + } + + if !obj.name.starts_with(rule.prefix()) { + continue; + } + + if let Some(filter) = &rule.filter { + let object_tags = ReplicationTagFilter::decode_tags_to_map(&obj.user_tags); + if filter.test_tags(&object_tags) { + rules.push(rule.clone()); + } + } else { + rules.push(rule.clone()); + } + } + + rules.sort_by(|a, b| { + if a.destination == b.destination { + a.priority.cmp(&b.priority) + } else { + std::cmp::Ordering::Equal + } + }); + + rules + } + + /// Retrieve the destination configuration + fn get_destination(&self) -> Destination { + if !self.rules.is_empty() { + self.rules[0].destination.clone() + } else { + Destination { + bucket: String::new(), + ..Default::default() + } + } + } + + /// Determine whether an object should be replicated + fn replicate(&self, obj: &ObjectOpts) -> bool { + let rules = self.filter_actionable_rules(obj); + + for rule in rules.iter() { + if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) { + 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); + } + false + } + + /// Check for an active rule + /// Optionally accept a prefix + /// When recursive is true, return true if any level under the prefix has an active rule + /// Without a prefix, recursive behaves as true + fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool { + if self.rules.is_empty() { + return false; + } + + for rule in &self.rules { + if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) { + continue; + } + + if let Some(filter) = &rule.filter + && let Some(filter_prefix) = &filter.prefix + { + if !prefix.is_empty() && !filter_prefix.is_empty() { + // The provided prefix must fall within the rule prefix + if !recursive && !prefix.starts_with(filter_prefix) { + continue; + } + } + + // When recursive, skip this rule if it does not match the test prefix or hierarchy + if recursive && !rule.prefix().starts_with(prefix) && !prefix.starts_with(rule.prefix()) { + continue; + } + } + return true; + } + false + } + + /// Filter target ARNs and return a slice of the distinct values in the config + fn filter_target_arns(&self, obj: &ObjectOpts) -> Vec { + let role = self.role.trim(); + if !role.is_empty() { + return vec![role.to_string()]; + } + + let mut arns = Vec::new(); + let mut targets_map: HashSet = HashSet::new(); + let rules = self.filter_actionable_rules(obj); + + for rule in rules { + if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) { + continue; + } + + let arn = rule.destination.bucket.trim(); + if !arn.is_empty() && !targets_map.contains(arn) { + targets_map.insert(arn.to_string()); + } + } + + for arn in targets_map { + arns.push(arn); + } + arns + } +} + +#[cfg(test)] +mod tests { + use super::*; + use s3s::dto::{DeleteMarkerReplication, Destination, ExistingObjectReplication, ReplicationRule}; + + fn replication_rule(id: &str, arn: &str) -> ReplicationRule { + ReplicationRule { + delete_marker_replication: Some(DeleteMarkerReplication::default()), + delete_replication: None, + destination: Destination { + bucket: arn.to_string(), + ..Default::default() + }, + existing_object_replication: Some(ExistingObjectReplication { + status: ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::ENABLED), + }), + filter: None, + id: Some(id.to_string()), + prefix: Some(String::new()), + priority: Some(1), + source_selection_criteria: None, + status: ReplicationRuleStatus::from_static(ReplicationRuleStatus::ENABLED), + } + } + + #[test] + fn filter_target_arns_uses_role_when_role_is_present() { + let config = ReplicationConfiguration { + role: " arn:legacy:target ".to_string(), + rules: vec![ + replication_rule("rule-1", "arn:target:a"), + replication_rule("rule-2", "arn:target:b"), + ], + }; + + let arns = config.filter_target_arns(&ObjectOpts { + name: "object".to_string(), + op_type: ReplicationType::Object, + ..Default::default() + }); + + assert_eq!(arns, vec!["arn:legacy:target".to_string()]); + } + + #[test] + fn filter_target_arns_falls_back_to_role_when_destination_is_empty() { + let config = ReplicationConfiguration { + role: "arn:legacy:target".to_string(), + rules: vec![ReplicationRule { + delete_marker_replication: Some(DeleteMarkerReplication::default()), + delete_replication: None, + destination: Destination { + bucket: String::new(), + ..Default::default() + }, + existing_object_replication: Some(ExistingObjectReplication { + status: ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::ENABLED), + }), + filter: None, + id: Some("rule-1".to_string()), + prefix: Some(String::new()), + priority: Some(1), + source_selection_criteria: None, + status: ReplicationRuleStatus::from_static(ReplicationRuleStatus::ENABLED), + }], + }; + + let arns = config.filter_target_arns(&ObjectOpts { + name: "object".to_string(), + op_type: ReplicationType::Object, + ..Default::default() + }); + + assert_eq!(arns, vec!["arn:legacy:target".to_string()]); + } + + fn replication_rule_existing_object_disabled(id: &str, arn: &str) -> ReplicationRule { + ReplicationRule { + delete_marker_replication: Some(DeleteMarkerReplication::default()), + delete_replication: None, + destination: Destination { + bucket: arn.to_string(), + ..Default::default() + }, + existing_object_replication: Some(ExistingObjectReplication { + status: ExistingObjectReplicationStatus::from_static(ExistingObjectReplicationStatus::DISABLED), + }), + filter: None, + id: Some(id.to_string()), + prefix: Some(String::new()), + priority: Some(1), + source_selection_criteria: None, + status: ReplicationRuleStatus::from_static(ReplicationRuleStatus::ENABLED), + } + } + + // Regression test for BUG-3: replicate_object was calling filter_target_arns with + // existing_object:false regardless of op_type, letting ExistingObject resync operations + // fan out to targets whose rule has ExistingObjectReplicationStatus::DISABLED. + #[test] + fn filter_target_arns_excludes_disabled_existing_object_target_for_existing_object_op() { + let config = ReplicationConfiguration { + role: String::new(), + rules: vec![ + replication_rule("rule-enabled", "arn:target:enabled"), + replication_rule_existing_object_disabled("rule-disabled", "arn:target:disabled"), + ], + }; + + let arns = config.filter_target_arns(&ObjectOpts { + name: "object".to_string(), + op_type: ReplicationType::ExistingObject, + existing_object: true, + ..Default::default() + }); + + assert_eq!(arns.len(), 1, "only the ENABLED target should be returned for ExistingObject ops"); + assert!(arns.contains(&"arn:target:enabled".to_string())); + assert!(!arns.contains(&"arn:target:disabled".to_string())); + } + + // Heal operations intentionally bypass ExistingObjectReplicationStatus — healing a past + // failure is not subject to the existing-object opt-out. + #[test] + fn filter_target_arns_includes_disabled_existing_object_target_for_heal_op() { + let config = ReplicationConfiguration { + role: String::new(), + rules: vec![ + replication_rule("rule-enabled", "arn:target:enabled"), + replication_rule_existing_object_disabled("rule-disabled", "arn:target:disabled"), + ], + }; + + let arns = config.filter_target_arns(&ObjectOpts { + name: "object".to_string(), + op_type: ReplicationType::Heal, + existing_object: false, + ..Default::default() + }); + + assert_eq!( + arns.len(), + 2, + "Heal ops must reach all targets regardless of existing_object_replication setting" + ); + assert!(arns.contains(&"arn:target:enabled".to_string())); + assert!(arns.contains(&"arn:target:disabled".to_string())); + } +} diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index 2d809bc10..b138c5ed2 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -12,14 +12,19 @@ // See the License for the specific language governing permissions and // limitations under the License. +pub mod config; pub mod mrf; pub mod resync; +pub mod rule; +pub mod tagging; +pub use config::{ObjectOpts, ReplicationConfigurationExt}; pub use mrf::{MrfOpKind, MrfReplicateEntry, decode_mrf_file, encode_mrf_file}; pub use resync::{ BucketReplicationResyncStatus, Error, Result, ResyncOpts, ResyncStatusType, TargetReplicationResyncStatus, decode_resync_file, encode_resync_file, }; +pub use rule::ReplicationRuleExt; pub use rustfs_filemeta::{ REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, REPLICATE_HEAL, REPLICATE_HEAL_DELETE, REPLICATE_INCOMING_DELETE, ReplicateDecision, ReplicateObjectInfo, ReplicateTargetDecision, ReplicatedInfos, ReplicatedTargetInfo, ReplicationAction, @@ -27,3 +32,4 @@ pub use rustfs_filemeta::{ VersionPurgeStatusType, get_replication_state, parse_replicate_decision, replication_statuses_map, target_reset_header, version_purge_statuses_map, }; +pub use tagging::{ReplicationTagFilter, decode_tags_to_map}; diff --git a/crates/ecstore/src/bucket/replication/rule.rs b/crates/replication/src/rule.rs similarity index 98% rename from crates/ecstore/src/bucket/replication/rule.rs rename to crates/replication/src/rule.rs index 9944109b3..4ca1871a0 100644 --- a/crates/ecstore/src/bucket/replication/rule.rs +++ b/crates/replication/src/rule.rs @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -use super::config::ObjectOpts; +use crate::config::ObjectOpts; use s3s::dto::ReplicaModificationsStatus; use s3s::dto::ReplicationRule; diff --git a/crates/replication/src/tagging.rs b/crates/replication/src/tagging.rs new file mode 100644 index 000000000..23953be02 --- /dev/null +++ b/crates/replication/src/tagging.rs @@ -0,0 +1,60 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::collections::HashMap; + +use url::form_urlencoded; + +pub struct ReplicationTagFilter; + +impl ReplicationTagFilter { + pub fn decode_tags_to_map(tags: &str) -> HashMap { + decode_tags_to_map(tags) + } +} + +pub fn decode_tags_to_map(tags: &str) -> HashMap { + let mut list = HashMap::new(); + + for (k, v) in form_urlencoded::parse(tags.as_bytes()) { + if k.is_empty() { + continue; + } + + list.insert(k.to_string(), v.to_string()); + } + + list +} + +#[cfg(test)] +mod tests { + use super::{ReplicationTagFilter, decode_tags_to_map}; + + #[test] + fn decode_tags_to_map_preserves_bucket_tagging_parser_behavior() { + let tags = decode_tags_to_map("env=prod&encoded=a%2Fb&=ignored"); + + assert_eq!(tags.get("env").map(String::as_str), Some("prod")); + assert_eq!(tags.get("encoded").map(String::as_str), Some("a/b")); + assert!(!tags.contains_key("")); + } + + #[test] + fn replication_tag_filter_uses_shared_parser() { + let tags = ReplicationTagFilter::decode_tags_to_map("env=prod"); + + assert_eq!(tags.get("env").map(String::as_str), Some("prod")); + } +} diff --git a/scripts/check_architecture_migration_rules.sh b/scripts/check_architecture_migration_rules.sh index ad7f41c66..3c3f12166 100755 --- a/scripts/check_architecture_migration_rules.sh +++ b/scripts/check_architecture_migration_rules.sh @@ -199,6 +199,7 @@ FUZZ_ECSTORE_COMPAT_BYPASS_HITS_FILE="${TMP_DIR}/fuzz_ecstore_compat_bypass_hits EXTERNAL_ECSTORE_API_BOUNDARY_HITS_FILE="${TMP_DIR}/external_ecstore_api_boundary_hits.txt" REPLICATION_FACADE_BYPASS_HITS_FILE="${TMP_DIR}/replication_facade_bypass_hits.txt" REPLICATION_FACADE_WILDCARD_EXPORT_HITS_FILE="${TMP_DIR}/replication_facade_wildcard_export_hits.txt" +REPLICATION_CONFIG_RULE_CONTRACT_BACKSLIDE_HITS_FILE="${TMP_DIR}/replication_config_rule_contract_backslide_hits.txt" REPLICATION_RESYNC_CONTRACT_BACKSLIDE_HITS_FILE="${TMP_DIR}/replication_resync_contract_backslide_hits.txt" REPLICATION_MRF_WIRE_FORMAT_BACKSLIDE_HITS_FILE="${TMP_DIR}/replication_mrf_wire_format_backslide_hits.txt" STORAGE_REPLICATION_HANDLE_BOUNDARY_BYPASS_HITS_FILE="${TMP_DIR}/storage_replication_handle_boundary_bypass_hits.txt" @@ -2543,6 +2544,21 @@ if [[ -s "$REPLICATION_FACADE_WILDCARD_EXPORT_HITS_FILE" ]]; then report_failure "replication facade must use explicit compatibility exports instead of wildcard re-exports: $(paste -sd '; ' "$REPLICATION_FACADE_WILDCARD_EXPORT_HITS_FILE")" fi +( + cd "$ROOT_DIR" + replication_config_rule_status=0 + rg -n --with-filename '^\s*(?:pub(?:\([^)]*\))?\s+)?(?:struct\s+ObjectOpts|trait\s+ReplicationConfigurationExt|trait\s+ReplicationRuleExt|struct\s+ReplicationTagFilter|fn\s+decode_tags_to_map)\b' \ + crates/ecstore/src/bucket/replication \ + --glob '*.rs' >"$REPLICATION_CONFIG_RULE_CONTRACT_BACKSLIDE_HITS_FILE" || replication_config_rule_status=$? + if [[ "$replication_config_rule_status" -ne 0 && "$replication_config_rule_status" -ne 1 ]]; then + exit "$replication_config_rule_status" + fi +) + +if [[ -s "$REPLICATION_CONFIG_RULE_CONTRACT_BACKSLIDE_HITS_FILE" ]]; then + report_failure "replication config/rule/tag contracts must stay in crates/replication: $(paste -sd '; ' "$REPLICATION_CONFIG_RULE_CONTRACT_BACKSLIDE_HITS_FILE")" +fi + ( cd "$ROOT_DIR" rg -n --with-filename 'pub\s+(struct|enum)\s+(ResyncOpts|TargetReplicationResyncStatus|BucketReplicationResyncStatus|ResyncStatusType)\b' \