From cb37a64a81bb857ea22efd2ebb015fbd3b64bef0 Mon Sep 17 00:00:00 2001 From: abdullahnah92 <157595835+abdullahnah92@users.noreply.github.com> Date: Sat, 13 Jun 2026 18:34:39 +0300 Subject: [PATCH] fix(site-replication): bidirectional sync, pre-existing data back-fill, and related fixes (#3401) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(site-repl): clamp replication_cfg_mismatch to owning deployments and add resync status branch Fix 3: replication_cfg_mismatch was set on ALL deployments for a bucket when any replication-config discrepancy existed, even on deployments that simply have no config. mc computes "in sync" as max_buckets minus per-deployment mismatch entries; with N deployments all flagged for 1 bucket, the result underflows to 1-N = -1. Now only deployments that own a replication config are flagged. Fix 4: SiteReplicationResyncOpHandler accepted "start" and "cancel" but returned an empty body (and later a parse error on the client) for any other operation string. Add a SITE_REPL_RESYNC_STATUS="status" arm that returns the stored SRResyncOpStatus or an explicit {"status":"not-found"} object so mc never sees an unexpected end of JSON input. Co-Authored-By: Claude Sonnet 4.6 * fix(site-repl): derive sync_state from reachability and replication rule completeness PeerInfo.sync_state was set to SyncStatus::Unknown at every construction site and never updated from real signals. build_status_info now tracks which peers were reachable during the metainfo fetch phase and, after all bucket stats are merged, derives sync_state as: - Enable : reachable AND no replication_cfg_mismatch for any bucket - Disable : reachable BUT at least one bucket has incomplete/missing rules - Unknown : unreachable (fetch failed) The "Sync" column in mc admin replicate status will now reflect actual state rather than always showing Unknown. Co-Authored-By: Claude Sonnet 4.6 * fix(site-repl): add SRRotateServiceAccountHandler for split-brain svc-acct repair When site-replicator-0 gets desynced (different secret on different peers), calls via the old service account return 403, blocking remove and replicate operations with no in-band recovery path. SiteReplicationRemoveHandler already purges local state unconditionally before sending peer notifications, so force-remove works locally even when peers reject the 403. The new POST /v3/site-replication/rotate-svc-acct endpoint provides the missing recovery path: it generates a fresh service-account secret, applies it locally, and pushes a peer/join to every member. Each peer's SRPeerJoinHandler accepts the join idempotently (update if exists, create otherwise), repairing the desynced credential without a full teardown + restart. Co-Authored-By: Claude Sonnet 4.6 * fix(site-repl): back-fill pre-existing buckets and objects on replicate add SiteReplicationAddHandler and SRPeerJoinHandler both called persist_site_replication_state and returned success without doing anything for buckets that existed before the sites were linked. Only buckets created after the link (via site_replication_make_bucket_hook) ever replicated. Introduce backfill_existing_buckets_after_add which, after persisting the new state, iterates every local bucket and: 1. Ensures versioning is enabled (required by replication). 2. Reconciles bucket targets so a target entry exists for each remote peer. 3. Reconciles the replication config so a rule pointing to each peer is present. 4. Broadcasts a make-bucket-hook to peers (idempotent) so they create the bucket. 5. Kicks start_site_bucket_resync toward every remote peer so pre-existing objects travel across. Errors per bucket are logged but never abort the overall add — manual resync remains available as a fallback. Both the initiator and the receiving (join) side run the backfill so convergence happens regardless of which side held data first. Co-Authored-By: Claude Sonnet 4.6 * fix(site-repl): reconcile replication config to restore bidirectional replication Root cause: ensure_site_replication_bucket_replication_config bailed with Ok(()) the moment any replication config was found on the bucket. When bucket B was propagated to site-2 via the make-with-versioning bucket-op, site-2's configure-replication step loaded the freshly-written config and immediately returned, never adding the reverse-direction rule pointing back to site-1. Result: objects uploaded to site-2 failed with "replication head_object fallback failed ... service error" because no rule targeted the originating site. Fix: drop the early-return. Instead load the existing rules, build the full desired config via build_site_replication_config, and MERGE — adding only the rules that are absent (identified by their "site-repl-" id). Re- number priorities after each merge to avoid conflicts. Existing non-site-repl rules are preserved. The write is skipped entirely when all desired rules are already present, so repeat calls remain cheap. Together with Fix 1 (backfill), any file written to ANY member now replicates to ALL members. The offline-and-recover case also benefits: when a peer returns, the per-bucket rules are complete and the resyncer can catch up. Also register the new rotate-svc-acct route in route_policy and route_registration_test so the route inventory assertions stay green. Co-Authored-By: Claude Sonnet 4.6 * fix(replication): remove broken resync worker-signal gate The resync_bucket task opened a new broadcast receiver via worker_rx.resubscribe(), which positions the receiver at the current write-head of the ring buffer — past all 10 bootstrap signals written in ReplicationResyncer::new(). Every spawned resync task therefore blocked on recv() forever, making `mc admin replicate resync start` report "started" while no objects ever moved. Remove the dead wait entirely. Each resync_bucket call is already spawned on-demand (tokio::spawn in start_bucket_resync / load_resync), so no additional gate is needed. The per-object concurrency limit is already enforced by the inner mpsc worker channels (line ~877). Also remove the now-dead worker_tx/worker_rx fields, bootstrap loop, and signal-send in resync_bucket_mark_status. Co-Authored-By: Claude Sonnet 4.6 * docs(site-repl): document rotate-svc-acct idempotency and partial-failure behavior * style: cargo fmt and remove redundant clones * fix(site-repl): preserve object-lock state when back-filling existing buckets --------- Co-authored-by: Claude Sonnet 4.6 Co-authored-by: loverustfs Co-authored-by: houseme --- .../replication/replication_resyncer.rs | 52 +- rustfs/src/admin/handlers/site_replication.rs | 493 +++++++++++++++++- rustfs/src/admin/route_policy.rs | 6 + rustfs/src/admin/route_registration_test.rs | 2 + 4 files changed, 508 insertions(+), 45 deletions(-) diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index 308e4e57f..54b3d7e53 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -98,7 +98,6 @@ const EVENT_RESYNC_OBJECT_PROCESSED: &str = "replication_resync_object_processed const EVENT_RESYNC_RUNTIME_SKIPPED: &str = "replication_resync_runtime_skipped"; const EVENT_REPLICATION_DELETE_SKIPPED: &str = "replication_delete_skipped"; const EVENT_REPLICATION_FORCE_DELETE_SKIPPED: &str = "replication_force_delete_skipped"; -const EVENT_RESYNC_WORKER_SIGNAL_FAILED: &str = "replication_resync_worker_signal_failed"; const EVENT_RESYNC_TASK_FAILED: &str = "replication_resync_task_failed"; const EVENT_RESYNC_TARGET_OPERATION_FAILED: &str = "replication_resync_target_operation_failed"; const EVENT_RESYNC_RUNTIME_CHANNEL_FAILED: &str = "replication_resync_runtime_channel_failed"; @@ -459,33 +458,14 @@ pub struct ReplicationResyncer { pub status_map: Arc>>, pub worker_size: usize, pub cancel_tokens: Arc>>, - pub worker_tx: tokio::sync::broadcast::Sender<()>, - pub worker_rx: tokio::sync::broadcast::Receiver<()>, } impl ReplicationResyncer { pub async fn new() -> Self { - let (worker_tx, worker_rx) = tokio::sync::broadcast::channel(RESYNC_WORKER_COUNT); - - for _ in 0..RESYNC_WORKER_COUNT { - if let Err(err) = worker_tx.send(()) { - error!( - event = EVENT_RESYNC_WORKER_SIGNAL_FAILED, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, - reason = "worker_bootstrap_signal_failed", - error = %err, - "Failed to signal replication resync worker" - ); - } - } - Self { status_map: Arc::new(RwLock::new(HashMap::new())), worker_size: RESYNC_WORKER_COUNT, cancel_tokens: Arc::new(RwLock::new(HashMap::new())), - worker_tx, - worker_rx, } } @@ -686,18 +666,6 @@ impl ReplicationResyncer { "Failed to update resync status" ); } - if let Err(err) = self.worker_tx.send(()) { - error!( - event = EVENT_RESYNC_WORKER_SIGNAL_FAILED, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, - bucket = %opts.bucket, - arn = %opts.arn, - reason = "worker_release_signal_failed", - error = %err, - "Failed to signal replication resync worker" - ); - } // TODO: Metrics } @@ -709,14 +677,18 @@ impl ReplicationResyncer { heal: bool, opts: ResyncOpts, ) { - let mut worker_rx = self.worker_rx.resubscribe(); - - tokio::select! { - _ = cancellation_token.cancelled() => { - return; - } - - _ = worker_rx.recv() => {} + // Check cancellation before starting the scan. + // NOTE: the previous design waited here on `worker_rx.resubscribe().recv()` to + // throttle concurrent resyncs, but `resubscribe()` positions the new receiver at + // the current write-head of the broadcast ring buffer, so all pre-sent bootstrap + // signals (written in `ReplicationResyncer::new`) are invisible to it. Every + // spawned task therefore blocked forever, which is why `resync start` reported + // "started" yet objects never moved. Throttling at this level is also incorrect + // for broadcast channels (one send unblocks ALL receivers). The inner + // per-object worker pool (mpsc channels, line ~877) already provides the right + // concurrency limit. + if cancellation_token.is_cancelled() { + return; } let cfg = match get_replication_config(&opts.bucket).await { diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 6a70db574..ed0f00045 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -99,6 +99,7 @@ const SITE_REPL_EDIT_SUCCESS: &str = "Requested site was updated successfully."; const SITE_REPL_REMOVE_SUCCESS: &str = "Requested site(s) were removed from cluster replication successfully."; const SITE_REPL_RESYNC_START: &str = "start"; const SITE_REPL_RESYNC_CANCEL: &str = "cancel"; +const SITE_REPL_RESYNC_STATUS: &str = "status"; const SITE_REPL_MIN_NETPERF_DURATION: Duration = Duration::from_secs(1); const SITE_REPLICATION_PEER_REQUEST_TIMEOUT: Duration = Duration::from_secs(10); const SITE_REPLICATION_PEER_CONNECT_TIMEOUT: Duration = Duration::from_secs(3); @@ -332,6 +333,11 @@ pub fn register_site_replication_route(r: &mut S3Router) -> std: AdminOperation(&SiteReplicationResyncOpHandler {}), ), (Method::PUT, "/v3/site-replication/state/edit", AdminOperation(&SRStateEditHandler {})), + ( + Method::POST, + "/v3/site-replication/rotate-svc-acct", + AdminOperation(&SRRotateServiceAccountHandler {}), + ), ] { r.insert(method, format!("{ADMIN_PREFIX}{path}").as_str(), operation)?; } @@ -1795,7 +1801,7 @@ fn merge_bucket_status_info(status: &mut SRStatusInfo, site_infos: &BTreeMap { site_infos.insert(deployment_id.clone(), filter_sr_info(peer_info, &opts)); + reachable_peers.insert(deployment_id.clone()); } Err(err) => { warn!(peer = %peer.endpoint, error = ?err, "site replication peer metainfo fetch failed"); @@ -2031,6 +2040,32 @@ async fn build_status_info(state: &SiteReplicationState, local_peer: &PeerInfo, } prune_in_sync_status_details(&mut status, &opts); + // Fix 2: derive sync_state from real signals — reachability + replication rule completeness + // instead of always returning SyncStatus::Unknown as stored in the persisted peer map. + { + let peer_has_replication_issue: HashMap = status + .sites + .keys() + .map(|dep_id| { + let has_issue = status + .bucket_stats + .values() + .any(|by_dep| by_dep.get(dep_id.as_str()).is_some_and(|s| s.replication_cfg_mismatch)); + (dep_id.clone(), has_issue) + }) + .collect(); + + for (deployment_id, peer) in status.sites.iter_mut() { + if !reachable_peers.contains(deployment_id) { + peer.sync_state = SyncStatus::Unknown; + } else if peer_has_replication_issue.get(deployment_id).copied().unwrap_or(false) { + peer.sync_state = SyncStatus::Disable; + } else { + peer.sync_state = SyncStatus::Enable; + } + } + } + if metrics_requested { status.metrics = build_metrics_summary(local_peer).await; } @@ -2429,14 +2464,51 @@ async fn ensure_site_replication_bucket_replication_config( state: &SiteReplicationState, local_peer: &PeerInfo, ) -> S3Result<()> { - match metadata_sys::get_replication_config(bucket).await { - Ok(_) => return Ok(()), - Err(StorageError::ConfigNotFound) => {} + // Fix 6: reconcile rather than early-returning when any config already exists. + // The old code bailed on Ok(_), so the second site joined a replicated bucket and + // ended up with NO rule pointing back to the first site. Objects on that site could + // never travel back, producing the "one-directional" replication symptom. + let Some(desired) = build_site_replication_config(bucket, state, local_peer)? else { + return Ok(()); + }; + + // Load the existing rules (may be empty if never configured). + let mut existing_rules = match metadata_sys::get_replication_config(bucket).await { + Ok((existing, _)) => existing.rules, + Err(StorageError::ConfigNotFound) => Vec::new(), Err(err) => return Err(ApiError::from(err).into()), + }; + + // Collect the IDs of existing site-repl-* rules so we don't duplicate them. + let existing_rule_ids: HashSet = existing_rules + .iter() + .filter_map(|r| r.id.as_deref()) + .filter(|id| id.starts_with("site-repl-")) + .map(String::from) + .collect(); + + let mut added = false; + for rule in desired.rules { + let rule_id = rule.id.as_deref().unwrap_or(""); + if !existing_rule_ids.contains(rule_id) { + existing_rules.push(rule); + added = true; + } } - let Some(config) = build_site_replication_config(bucket, state, local_peer)? else { + if !added { + // All desired rules are already present — nothing to write. return Ok(()); + } + + // Re-assign contiguous priorities to avoid conflicts with any preserved rules. + for (i, rule) in existing_rules.iter_mut().enumerate() { + rule.priority = Some((i + 1) as i32); + } + + let config = ReplicationConfiguration { + role: String::new(), + rules: existing_rules, }; let data = serialize(&config) @@ -2453,6 +2525,63 @@ pub async fn site_replication_peer_deployment_id_for_endpoint(endpoint: &str) -> peer_deployment_id_for_endpoint(&state, endpoint) } +/// Fix 1: after persisting a new site-replication state (add or join), enumerate every bucket +/// that already exists locally, wire up versioning + targets + replication config for each, and +/// kick a resync toward every remote peer so pre-existing objects back-fill. Errors are logged +/// but never abort the caller — the admin can run a manual resync if needed. +async fn backfill_existing_buckets_after_add(state: &SiteReplicationState, local_peer: &PeerInfo) { + let Some(store) = new_object_layer_fn() else { + return; + }; + let buckets = match store.list_bucket(&BucketOptions::default()).await { + Ok(b) => b, + Err(err) => { + warn!(error = ?err, "site replication backfill: failed to list buckets"); + return; + } + }; + + let resync_id = Uuid::new_v4().to_string(); + for bucket in &buckets { + let name = &bucket.name; + + if let Err(err) = ensure_site_replication_bucket_versioning(name).await { + warn!(bucket = %name, error = ?err, "site replication backfill: versioning setup failed"); + continue; + } + if let Err(err) = ensure_site_replication_bucket_targets(name, state, local_peer, None).await { + warn!(bucket = %name, error = ?err, "site replication backfill: targets setup failed"); + } + if let Err(err) = ensure_site_replication_bucket_replication_config(name, state, local_peer).await { + warn!(bucket = %name, error = ?err, "site replication backfill: replication config setup failed"); + } + // Broadcast the bucket to peers so they create it too (idempotent on the peer side). + // Read the real lock_enabled flag so peers recreate the bucket with the same object-lock + // setting — object lock cannot be added after bucket creation. + let lock_enabled = match metadata_sys::get(name).await { + Ok(bm) => bm.lock_enabled, + Err(err) => { + warn!(bucket = %name, error = ?err, "site replication backfill: failed to read bucket metadata, assuming lock_enabled=false"); + false + } + }; + if let Err(err) = site_replication_make_bucket_hook(name, lock_enabled).await { + warn!(bucket = %name, error = ?err, "site replication backfill: make-bucket broadcast failed"); + } + // Kick a resync toward every remote peer so existing objects travel across. + for peer in state.peers.values() { + if peer.deployment_id == local_peer.deployment_id || same_identity_endpoint(&peer.endpoint, &local_peer.endpoint) { + continue; + } + let result = start_site_bucket_resync(name, peer, &resync_id).await; + if result.status == "failed" { + warn!(bucket = %name, peer = %peer.endpoint, detail = %result.err_detail, + "site replication backfill: resync kick failed"); + } + } + } +} + async fn start_site_bucket_resync(bucket: &str, peer: &PeerInfo, resync_id: &str) -> ResyncBucketStatus { let mut bucket_status = ResyncBucketStatus { bucket: bucket.to_string(), @@ -2996,6 +3125,12 @@ impl Operation for SiteReplicationAddHandler { } persist_site_replication_state(&state).await?; + + // Fix 1: back-fill pre-existing buckets so objects created before `replicate add` + // are not silently left out of replication. Failures are logged but do not abort + // the overall add operation — the admin can trigger a manual resync if needed. + backfill_existing_buckets_after_add(&state, &local_peer).await; + json_response(&ReplicateAddStatus { success: true, status: SITE_REPL_ADD_SUCCESS.to_string(), @@ -3209,6 +3344,9 @@ impl Operation for SRPeerJoinHandler { .filter(|name| !name.is_empty()) .unwrap_or_else(|| local_peer.name.clone()); persist_site_replication_state(&state).await?; + // Fix 1 (receiving side): ensure the joining peer also sets up replication for any + // buckets it already owns so the reverse direction works from the start. + backfill_existing_buckets_after_add(&state, &local_peer).await; json_response(&SRPeerJoinResponse { peer: state.peers.get(&local_peer.deployment_id).cloned().unwrap_or(local_peer), }) @@ -3542,6 +3680,18 @@ impl Operation for SiteReplicationResyncOpHandler { state.resync_status.insert(peer.deployment_id.clone(), status.clone()); status } + SITE_REPL_RESYNC_STATUS => { + let status = state + .resync_status + .get(&peer.deployment_id) + .cloned() + .unwrap_or_else(|| SRResyncOpStatus { + op_type: SITE_REPL_RESYNC_STATUS.to_string(), + status: "not-found".to_string(), + ..Default::default() + }); + return json_response(&status); + } _ => return Err(s3_error!(InvalidRequest, "unsupported resync operation")), }; save_site_replication_state(&state).await?; @@ -3562,6 +3712,77 @@ impl Operation for SRStateEditHandler { } } +/// Repairs a split-brained `site-replicator-0` service account. +/// +/// When the internal service account is desynced (e.g. after a failed `rm` left stale state on +/// one peer), admin calls to that peer return 403. This handler recovers the cluster without a +/// full teardown: +/// +/// 1. Generates a fresh service-account secret locally. +/// 2. Applies it to the local node and persists state. +/// 3. Pushes `peer/join` with the new credentials to every remote peer. +/// A peer whose secret is already correct accepts the update idempotently. +/// A peer whose secret was stale is repaired. +/// +/// **Partial failure**: if one or more peers are unreachable the local node is still updated and +/// `status="Partial"` is returned with `err_detail` listing each failed endpoint and its error. +/// The call is **idempotent** — re-run it until `status="Success"` to repair all peers. +pub struct SRRotateServiceAccountHandler {} + +#[async_trait::async_trait] +impl Operation for SRRotateServiceAccountHandler { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?; + let mut state = load_site_replication_state().await?; + if !state.enabled() { + return Err(s3_error!(InvalidRequest, "site replication is not configured")); + } + let local_peer = current_local_peer(&req, &state); + + // Force generation of a new secret by passing a state with empty secret key. + let rotation_state = SiteReplicationState { + service_account_access_key: SITE_REPLICATOR_SERVICE_ACCOUNT.to_string(), + service_account_secret_key: String::new(), + ..state.clone() + }; + let (svc_ak, new_sk) = ensure_site_replicator_service_account(&cred.access_key, &rotation_state).await?; + + state.service_account_access_key = svc_ak.clone(); + state.service_account_secret_key = new_sk.clone(); + + let join_req = SRPeerJoinReq { + svc_acct_access_key: svc_ak.clone(), + svc_acct_secret_key: new_sk.clone(), + svc_acct_parent: cred.access_key.clone(), + peers: state.peers.clone(), + updated_at: state.updated_at, + }; + + let mut peer_errors = Vec::new(); + for peer in state.peers.values() { + if same_identity_endpoint(&peer.endpoint, &local_peer.endpoint) { + continue; + } + if let Err(err) = + send_peer_admin_request(&peer.endpoint, SITE_REPLICATION_PEER_JOIN_PATH, &svc_ak, &new_sk, &join_req).await + { + let detail = summarize_peer_error_detail(&format!("{}: {err}", peer.endpoint)); + warn!(peer = %peer.endpoint, error = %detail, "site replication service account rotation failed for peer"); + peer_errors.push(detail); + } + } + + persist_site_replication_state(&state).await?; + + json_response(&ReplicateEditStatus { + success: peer_errors.is_empty(), + status: if peer_errors.is_empty() { "Success" } else { "Partial" }.to_string(), + err_detail: peer_errors.join("; "), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + }) + } +} + #[cfg(test)] mod tests { use super::*; @@ -4658,4 +4879,266 @@ mod tests { assert!(group_info_requires_upsert(&update)); } + + // Fix 3: replication_cfg_mismatch must not be set for deployments that simply have no + // replication config. Setting it globally caused mc to count N mismatch entries for a + // single bucket (one per deployment), while max_buckets=1, producing -1/N in sync. + #[test] + fn test_replication_cfg_mismatch_only_set_for_deployments_with_config() { + use rustfs_madmin::{SRBucketInfo, SRInfo}; + + let repl_xml = { + let config = ReplicationConfiguration { + role: String::new(), + rules: vec![build_site_replication_rule( + "arn:rustfs:replication::site-b:photos", + 1, + "site-repl-site-b", + )], + }; + String::from_utf8(serialize(&config).unwrap()).unwrap() + }; + + let mut site_a_info = SRInfo::default(); + site_a_info.buckets.insert( + "photos".to_string(), + SRBucketInfo { + bucket: "photos".to_string(), + replication_config: Some(repl_xml), + ..Default::default() + }, + ); + + // Site B has the bucket but NO replication config yet (partial setup) + let mut site_b_info = SRInfo::default(); + site_b_info.buckets.insert( + "photos".to_string(), + SRBucketInfo { + bucket: "photos".to_string(), + replication_config: None, + ..Default::default() + }, + ); + + let site_infos: BTreeMap = [("dep-a".to_string(), site_a_info), ("dep-b".to_string(), site_b_info)] + .into_iter() + .collect(); + + let mut status = SRStatusInfo { + sites: site_infos + .keys() + .map(|k| { + ( + k.clone(), + PeerInfo { + deployment_id: k.clone(), + ..Default::default() + }, + ) + }) + .collect(), + ..Default::default() + }; + for k in site_infos.keys() { + status.stats_summary.insert( + k.clone(), + SRSiteSummary { + api_version: Some(SITE_REPL_API_VERSION.to_string()), + ..Default::default() + }, + ); + } + + let opts = SRStatusOptions { + buckets: true, + ..Default::default() + }; + merge_bucket_status_info(&mut status, &site_infos, &opts); + + let bucket_stats = status.bucket_stats.get("photos").expect("photos bucket stats"); + let dep_a = bucket_stats.get("dep-a").expect("dep-a stats"); + let dep_b = bucket_stats.get("dep-b").expect("dep-b stats"); + + // dep-a has a config but it doesn't cover all peers → mismatch + assert!(dep_a.replication_cfg_mismatch, "dep-a has config but it is incomplete"); + // dep-b has NO config → must NOT be flagged as mismatch (only has_replication_cfg=false) + assert!( + !dep_b.replication_cfg_mismatch, + "dep-b has no config, mismatch must not be set to avoid -1 in mc output" + ); + assert!(!dep_b.has_replication_cfg, "dep-b should show has_replication_cfg=false"); + } + + // Fix 4: status operation must return a well-formed SRResyncOpStatus (not an empty body) + #[test] + fn test_resync_status_returns_not_found_when_no_resync_in_progress() { + let state = SiteReplicationState::default(); + let status = state + .resync_status + .get("nonexistent-peer") + .cloned() + .unwrap_or_else(|| SRResyncOpStatus { + op_type: SITE_REPL_RESYNC_STATUS.to_string(), + status: "not-found".to_string(), + ..Default::default() + }); + assert_eq!(status.status, "not-found"); + assert_eq!(status.op_type, SITE_REPL_RESYNC_STATUS); + } + + #[test] + fn test_resync_status_returns_existing_status_for_known_peer() { + let mut state = SiteReplicationState::default(); + state.resync_status.insert( + "peer-dep".to_string(), + SRResyncOpStatus { + op_type: SITE_REPL_RESYNC_START.to_string(), + resync_id: "abc-123".to_string(), + status: "success".to_string(), + ..Default::default() + }, + ); + let status = state + .resync_status + .get("peer-dep") + .cloned() + .unwrap_or_else(|| SRResyncOpStatus { + op_type: SITE_REPL_RESYNC_STATUS.to_string(), + status: "not-found".to_string(), + ..Default::default() + }); + assert_eq!(status.op_type, SITE_REPL_RESYNC_START); + assert_eq!(status.resync_id, "abc-123"); + } + + // Fix 2: sync_state must derive from real health signals, not always Unknown + #[test] + fn test_derive_sync_state_from_replication_completeness() { + // A peer that is reachable and has complete replication rules for all other peers + // should be Enable; one that is reachable but has an incomplete config should be Disable. + let repl_complete_xml = { + let config = ReplicationConfiguration { + role: String::new(), + rules: vec![build_site_replication_rule( + "arn:rustfs:replication::dep-b:bucket", + 1, + "site-repl-dep-b", + )], + }; + String::from_utf8(serialize(&config).unwrap()).unwrap() + }; + + // Peer that has complete config for 2-site setup + assert!(site_replication_rule_complete(&build_site_replication_rule( + "arn:rustfs:replication::dep-b:bucket", + 1, + "site-repl-dep-b" + ))); + assert_eq!( + site_replication_config_mismatch(vec![Some(&repl_complete_xml), Some(&repl_complete_xml)].into_iter(), 2), + (2, false), + "complete rules on both sites → no mismatch" + ); + assert_eq!( + site_replication_config_mismatch(vec![Some(&repl_complete_xml)].into_iter(), 2), + (1, true), + "config only on one of two sites → mismatch" + ); + } + + // Fix 5: remove --all must purge local state unconditionally even when peer errors occur + #[test] + fn test_remove_all_purges_local_state_unconditionally() { + let mut state = SiteReplicationState { + name: "local".to_string(), + service_account_access_key: "site-replicator-0".to_string(), + service_account_secret_key: "some-secret".to_string(), + ..Default::default() + }; + state.peers.insert( + "local-dep".to_string(), + PeerInfo { + deployment_id: "local-dep".to_string(), + ..peer("local", "https://local.example.com") + }, + ); + state.peers.insert( + "remote-dep".to_string(), + PeerInfo { + deployment_id: "remote-dep".to_string(), + ..peer("remote", "https://remote.example.com") + }, + ); + state.resync_status.insert( + "remote-dep".to_string(), + SRResyncOpStatus { + resync_id: "r1".to_string(), + status: "success".to_string(), + ..Default::default() + }, + ); + + // Simulate remove --all + let state = remove_sites( + state, + SRRemoveReq { + remove_all: true, + ..Default::default() + }, + ); + + // Local state must be cleared regardless of whether peer notifications succeed + assert!(state.peers.is_empty(), "peers must be cleared on remove --all"); + assert!(state.resync_status.is_empty(), "resync_status must be cleared on remove --all"); + + // Even if peers returned 403 (desynced account), status still reports success + let status = + site_replication_remove_status(&["https://remote.example.com: peer/remove returned 403 Forbidden".to_string()]); + assert_eq!( + status.status, SITE_REPL_REMOVE_SUCCESS, + "local remove reports success even when peer notifications fail" + ); + assert!( + status.err_detail.contains("403 Forbidden"), + "peer errors are included in err_detail for diagnostics" + ); + } + + // Fix 6: ensure_site_replication_bucket_replication_config must reconcile rather than + // early-return so that a bucket propagated to the second site gets a rule back to the first. + #[test] + fn test_reconcile_adds_missing_peer_rules_to_existing_config() { + // Start with a config that has only rule for dep-b (first site's initial config) + let rule_b = build_site_replication_rule("arn:rustfs:replication::dep-b:bucket", 1, "site-repl-dep-b"); + let rule_c = build_site_replication_rule("arn:rustfs:replication::dep-c:bucket", 2, "site-repl-dep-c"); + + let mut existing_rules = vec![rule_b.clone()]; + + // Desired config has rules for both dep-b and dep-c (3-site setup) + let desired_rules = vec![rule_b, rule_c]; + + // Simulate the reconcile: collect existing site-repl rule IDs + let existing_ids: std::collections::HashSet = existing_rules + .iter() + .filter_map(|r| r.id.as_deref()) + .filter(|id| id.starts_with("site-repl-")) + .map(String::from) + .collect(); + + let mut added = false; + for rule in &desired_rules { + let rid = rule.id.as_deref().unwrap_or(""); + if !existing_ids.contains(rid) { + existing_rules.push(rule.clone()); + added = true; + } + } + + assert!(added, "missing rule should have been added"); + assert_eq!(existing_rules.len(), 2, "should now have rules for both peers"); + + let rule_ids: Vec<&str> = existing_rules.iter().filter_map(|r| r.id.as_deref()).collect(); + assert!(rule_ids.contains(&"site-repl-dep-b")); + assert!(rule_ids.contains(&"site-repl-dep-c")); + } } diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index 4d06ae929..ec1429909 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -437,6 +437,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ SITE_REPLICATION_INFO, RouteRiskLevel::High, ), + admin( + HttpMethod::Post, + "/rustfs/admin/v3/site-replication/rotate-svc-acct", + SITE_REPLICATION_OPERATION, + RouteRiskLevel::Sensitive, + ), admin( HttpMethod::Put, "/rustfs/admin/v3/site-replication/peer/join", diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index a8d7b5c50..18fe08a3b 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -241,6 +241,7 @@ fn expected_admin_route_matrix() -> Vec { admin_route(Method::GET, "/v3/site-replication/status"), admin_route(Method::POST, "/v3/site-replication/devnull"), admin_route(Method::POST, "/v3/site-replication/netperf"), + admin_route(Method::POST, "/v3/site-replication/rotate-svc-acct"), admin_route(Method::PUT, "/v3/site-replication/peer/join"), admin_route(Method::PUT, "/v3/site-replication/peer/bucket-ops"), admin_route(Method::PUT, "/v3/site-replication/peer/iam-item"), @@ -769,6 +770,7 @@ fn test_register_routes_cover_representative_admin_paths() { assert_route(&router, Method::GET, &admin_path("/v3/site-replication/status")); assert_route(&router, Method::POST, &admin_path("/v3/site-replication/devnull")); assert_route(&router, Method::POST, &admin_path("/v3/site-replication/netperf")); + assert_route(&router, Method::POST, &admin_path("/v3/site-replication/rotate-svc-acct")); assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/peer/join")); assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/peer/bucket-ops")); assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/peer/iam-item"));