diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 9f962acd7..c8dcf5b9e 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -196,9 +196,9 @@ pub mod bucket { ReplicationOperation, ReplicationPoolTrait, ReplicationPriority, ReplicationQueueAdmission, ReplicationScannerBridge, ReplicationState, ReplicationStats, ReplicationStatusType, ReplicationStorage, ReplicationTargetValidationError, ReplicationType, ResyncOpts, ResyncStatusType, RuntimeReplicationTargetBacklog, TargetReplicationResyncStatus, - VersionPurgeStatusType, XferStats, commit_force_delete_intent, complete_force_delete_intent, - delete_replication_state_from_config, delete_replication_version_id, get_global_replication_pool, - get_global_replication_stats, get_proxy_targets, init_background_replication, + VersionPurgeStatusType, XferStats, assign_site_replication_rule_priorities, commit_force_delete_intent, + complete_force_delete_intent, delete_replication_state_from_config, delete_replication_version_id, + get_global_replication_pool, get_global_replication_stats, get_proxy_targets, init_background_replication, invalid_replication_config_status_field, is_site_replication_rule, merge_incoming_replication_config, merge_user_replication_config, persist_force_delete_intent, read_durable_mrf_backlog, replication_state_to_filemeta, replication_status_to_filemeta, replication_statuses_map, replication_target_arn_deployment_id, diff --git a/crates/ecstore/src/bucket/replication/mod.rs b/crates/ecstore/src/bucket/replication/mod.rs index 2469e508c..7755f1136 100644 --- a/crates/ecstore/src/bucket/replication/mod.rs +++ b/crates/ecstore/src/bucket/replication/mod.rs @@ -47,10 +47,10 @@ pub use replication_config_boundary::{ ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS, REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS, ReplicationConfigStructureError, ReplicationConfigurationExt, ReplicationTargetValidationError, - invalid_replication_config_status_field, is_site_replication_rule, merge_incoming_replication_config, - merge_user_replication_config, replication_target_arn_deployment_id, replication_target_arns, - should_remove_replication_target, unsupported_replication_config_field, validate_replication_config_structure, - validate_replication_config_target_arns, + assign_site_replication_rule_priorities, invalid_replication_config_status_field, is_site_replication_rule, + merge_incoming_replication_config, merge_user_replication_config, replication_target_arn_deployment_id, + replication_target_arns, should_remove_replication_target, unsupported_replication_config_field, + validate_replication_config_structure, validate_replication_config_target_arns, }; pub(crate) use replication_filemeta_boundary::version_purge_statuses_map; pub use replication_filemeta_boundary::{ diff --git a/crates/ecstore/src/bucket/replication/replication_config_boundary.rs b/crates/ecstore/src/bucket/replication/replication_config_boundary.rs index c476e385a..11142e2ac 100644 --- a/crates/ecstore/src/bucket/replication/replication_config_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_config_boundary.rs @@ -16,8 +16,8 @@ pub use rustfs_replication::{ ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS, REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS, ReplicationConfigStructureError, ReplicationConfigurationExt, ReplicationRuleExt, ReplicationTargetValidationError, - invalid_replication_config_status_field, is_site_replication_rule, merge_incoming_replication_config, - merge_user_replication_config, replication_target_arn_deployment_id, replication_target_arns, - should_remove_replication_target, unsupported_replication_config_field, validate_replication_config_structure, - validate_replication_config_target_arns, + assign_site_replication_rule_priorities, invalid_replication_config_status_field, is_site_replication_rule, + merge_incoming_replication_config, merge_user_replication_config, replication_target_arn_deployment_id, + replication_target_arns, should_remove_replication_target, unsupported_replication_config_field, + validate_replication_config_structure, validate_replication_config_target_arns, }; diff --git a/crates/replication/src/config.rs b/crates/replication/src/config.rs index 47b8b0047..eb007b72b 100644 --- a/crates/replication/src/config.rs +++ b/crates/replication/src/config.rs @@ -382,9 +382,7 @@ fn merge_replication_config_keeping_site_rules( return None; } - for (index, rule) in rules.iter_mut().enumerate() { - rule.priority = Some(i32::try_from(index + 1).unwrap_or(i32::MAX)); - } + assign_site_replication_rule_priorities(&mut rules, &is_site_rule); // A site-replication ARN in `role` is the sender's, and the reconciler's // per-peer target lookup reads it — carrying it over would pin the @@ -397,6 +395,30 @@ fn merge_replication_config_keeping_site_rules( Some(ReplicationConfiguration { role, rules }) } +/// Give the site rules in `rules` the lowest priorities no operator rule uses, +/// in rule order, leaving every operator rule's priority untouched. Operator +/// priorities decide which rule wins per target, so they are part of the +/// submitted policy; site rules are derived state and only need to be unique +/// (`validate_replication_config_structure` rejects duplicates). The result +/// is a pure function of the rule list, so the site-replication reconciler, +/// the peer ingestion merge and the S3 edit merge all converge on the same +/// bytes and the reconciler's no-op check holds. +pub fn assign_site_replication_rule_priorities(rules: &mut [ReplicationRule], is_site_rule: impl Fn(&ReplicationRule) -> bool) { + let taken: HashSet = rules + .iter() + .filter(|rule| !is_site_rule(rule)) + .map(|rule| rule.priority.unwrap_or(0)) + .collect(); + let mut next = 1; + for rule in rules.iter_mut().filter(|rule| is_site_rule(rule)) { + while taken.contains(&next) { + next += 1; + } + rule.priority = Some(next); + next = next.saturating_add(1); + } +} + pub fn replication_target_arns(config: &ReplicationConfiguration) -> HashSet { let role = config.role.trim(); if !role.is_empty() { @@ -1694,4 +1716,80 @@ mod tests { let removed_peer = replication_rule("site-repl-gone-dep", "arn:rustfs:replication::gone-dep:bucket"); assert!(!is_reconciler_owned_site_replication_rule(&removed_peer, &peers)); } + + // The merge must not rewrite the operator's priorities: with the + // priority-5 rule listed first and renumbered 1 then 2, the priority-1 + // delete-marker-disabled rule would win the replication decision. + #[test] + fn merge_keeps_operator_priorities_and_replication_decision() { + let user_arn = "arn:minio:replication:us-east-1:2f1c-remote:bucket"; + let peer_arn = "arn:rustfs:replication::peer-dep:bucket"; + let incoming = ReplicationConfiguration { + role: String::new(), + rules: vec![ + delete_marker_rule("dm-enabled", user_arn, "logs/", 5, true), + delete_marker_rule("dm-disabled", user_arn, "logs/2026/", 1, false), + ], + }; + let mut site_rule = delete_marker_rule("site-repl-peer-dep", peer_arn, "", 7, true); + site_rule.prefix = None; + let local = structure_config(vec![site_rule]); + let opts = ObjectOpts { + name: "logs/2026/app.log".to_string(), + op_type: ReplicationType::Delete, + delete_marker: true, + version_id: None, + ..Default::default() + }; + let submitted: Vec<_> = incoming.filter_target_replication_decisions(&opts); + + let peers = HashSet::from(["peer-dep".to_string()]); + let merged = merge_user_replication_config(Some(incoming.clone()), Some(local.clone()), &peers).expect("rules"); + + let priorities: Vec<_> = merged + .rules + .iter() + .map(|rule| (rule.id.as_deref().unwrap(), rule.priority)) + .collect(); + assert_eq!( + priorities, + vec![ + ("dm-enabled", Some(5)), + ("dm-disabled", Some(1)), + ("site-repl-peer-dep", Some(2)) + ], + "operator priorities are kept verbatim; the site rule takes the lowest free slot" + ); + assert!(validate_replication_config_structure(&merged).is_ok()); + let mut decisions = merged.filter_target_replication_decisions(&opts); + decisions.retain(|(arn, _)| arn == user_arn); + assert_eq!(decisions, submitted, "the merged config must replicate exactly as the operator submitted"); + assert_eq!(decisions, vec![(user_arn.to_string(), true)]); + + // The peer ingestion merge follows the same rule. + let merged = merge_incoming_replication_config(Some(incoming), Some(local)).expect("rules"); + let priorities: Vec<_> = merged.rules.iter().map(|rule| rule.priority).collect(); + assert_eq!(priorities, vec![Some(5), Some(1), Some(2)]); + } + + #[test] + fn site_rule_priorities_skip_every_operator_priority() { + let mut rules = vec![ + delete_marker_rule("a", "arn:a", "", 2, true), + delete_marker_rule("site-repl-x", "arn:rustfs:replication::x:b", "", 9, true), + delete_marker_rule("b", "arn:a", "", 1, true), + delete_marker_rule("site-repl-y", "arn:rustfs:replication::y:b", "", 9, true), + delete_marker_rule("c", "arn:a", "", 4, true), + ]; + assign_site_replication_rule_priorities(&mut rules, is_site_replication_rule); + let priorities: Vec<_> = rules.iter().map(|rule| rule.priority).collect(); + assert_eq!(priorities, vec![Some(2), Some(3), Some(1), Some(5), Some(4)]); + assert!(validate_replication_config_structure(&structure_config(rules.clone())).is_ok()); + + // Idempotent, so the reconciler's pass over an already-merged config + // is a byte-stable no-op rather than a rewrite every period. + let settled = rules.clone(); + assign_site_replication_rule_priorities(&mut rules, is_site_replication_rule); + assert_eq!(rules, settled); + } } diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index b80e8525d..b25d63575 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -32,11 +32,11 @@ pub use config::{ ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS, REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS, ReplicationConfigStructureError, ReplicationConfigurationExt, ReplicationTargetValidationError, - active_replication_rule_destination_arns, invalid_replication_config_status_field, is_reconciler_owned_site_replication_rule, - is_site_replication_rule, merge_incoming_replication_config, merge_user_replication_config, - replication_target_arn_deployment_id, replication_target_arns, should_remove_replication_target, - site_replication_rule_deployment_id, unsupported_replication_config_field, validate_replication_config_structure, - validate_replication_config_target_arns, + active_replication_rule_destination_arns, assign_site_replication_rule_priorities, invalid_replication_config_status_field, + is_reconciler_owned_site_replication_rule, is_site_replication_rule, merge_incoming_replication_config, + merge_user_replication_config, replication_target_arn_deployment_id, replication_target_arns, + should_remove_replication_target, site_replication_rule_deployment_id, unsupported_replication_config_field, + validate_replication_config_structure, validate_replication_config_target_arns, }; pub use delete::{ DeletedObjectReplicationInfo, delete_marker_purge_mrf_entry, delete_marker_purge_version_id, diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 60b0bdc89..502f7efc0 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -32,7 +32,8 @@ use crate::admin::storage_api::bucket::metadata_sys; use crate::admin::storage_api::bucket::quota::BucketQuota; use crate::admin::storage_api::bucket::replication; use crate::admin::storage_api::bucket::replication::{ - is_site_replication_rule, merge_incoming_replication_config, replication_target_arn_deployment_id, + assign_site_replication_rule_priorities, is_site_replication_rule, merge_incoming_replication_config, + replication_target_arn_deployment_id, }; use crate::admin::storage_api::bucket::target::{ARN, BucketTarget, BucketTargetType, BucketTargets, Credentials}; use crate::admin::storage_api::bucket::target_sys::BucketTargetSys; @@ -8168,9 +8169,7 @@ fn prune_removed_site_replication_rules( return (None, removed); } - for (index, rule) in config.rules.iter_mut().enumerate() { - rule.priority = Some(i32::try_from(index + 1).unwrap_or(i32::MAX)); - } + assign_site_replication_rule_priorities(&mut config.rules, is_site_replication_rule); (Some(config), removed) } @@ -8344,9 +8343,10 @@ async fn ensure_site_replication_bucket_replication_config_with_runtime( .cloned() .collect(); rules.extend(desired.rules); - for (index, rule) in rules.iter_mut().enumerate() { - rule.priority = Some(i32::try_from(index + 1).unwrap_or(i32::MAX)); - } + // Operator priorities are the operator's policy; only the derived rules + // take free slots, by the same function as the config merges so a merged + // write and this pass agree byte for byte. + assign_site_replication_rule_priorities(&mut rules, is_site_replication_rule); // Only a site-replication ARN in `role` is ours to drop — an operator-authored role is // part of the bucket's S3-visible configuration, and repairing a reverse rule must not @@ -16969,7 +16969,7 @@ mod tests { } #[test] - fn test_prune_removed_site_replication_rules_removes_site_rule_and_reorders_priorities() { + fn test_prune_removed_site_replication_rules_removes_site_rule_and_keeps_operator_priority() { let removed_deployment_ids = HashSet::from(["removed-dep".to_string()]); let kept_rule = build_site_replication_rule("arn:rustfs:replication::kept-dep:photos", 3, "site-repl-kept-dep"); let removed_rule = build_site_replication_rule("arn:rustfs:replication::removed-dep:photos", 1, "site-repl-removed-dep"); @@ -16986,9 +16986,9 @@ mod tests { assert!(updated.role.is_empty()); assert_eq!(updated.rules.len(), 2); assert_eq!(updated.rules[0].id.as_deref(), Some("user-managed-rule")); - assert_eq!(updated.rules[0].priority, Some(1)); + assert_eq!(updated.rules[0].priority, Some(9), "the operator's priority is policy and stays"); assert_eq!(updated.rules[1].id.as_deref(), Some("site-repl-kept-dep")); - assert_eq!(updated.rules[1].priority, Some(2)); + assert_eq!(updated.rules[1].priority, Some(1), "the derived rule moves to the lowest free slot"); } #[test] diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index 3538d6329..cef2c71f7 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -443,7 +443,8 @@ pub(crate) mod replication { pub(crate) use super::ecstore_bucket::replication::{ REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS, REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS, - is_site_replication_rule, merge_incoming_replication_config, replication_target_arn_deployment_id, + assign_site_replication_rule_priorities, is_site_replication_rule, merge_incoming_replication_config, + replication_target_arn_deployment_id, }; pub(crate) type BucketReplicationResyncStatus = super::ecstore_bucket::replication::BucketReplicationResyncStatus; pub(crate) type BucketStats = super::ecstore_bucket::replication::BucketStats;