From d960acb1707cd98561ad830db4c534361f249ac3 Mon Sep 17 00:00:00 2001 From: cxymds Date: Tue, 31 Mar 2026 18:12:56 +0800 Subject: [PATCH] fix(admin): reconcile site replication peer identity (#2356) --- rustfs/src/admin/handlers/site_replication.rs | 97 ++++++++++++++++++- 1 file changed, 92 insertions(+), 5 deletions(-) diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 4319859c2..ffcb5010f 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -102,6 +102,11 @@ struct SiteReplicationState { resync_status: BTreeMap, } +#[derive(Debug, Clone, Serialize, Deserialize, Default)] +struct SRPeerJoinResponse { + peer: PeerInfo, +} + const GO_GOB_SITE_NETPERF_SCHEMA: &[u8] = &[ 0x7d, 0x7f, 0x03, 0x01, 0x01, 0x15, 0x53, 0x69, 0x74, 0x65, 0x4e, 0x65, 0x74, 0x50, 0x65, 0x72, 0x66, 0x4e, 0x6f, 0x64, 0x65, 0x52, 0x65, 0x73, 0x75, 0x6c, 0x74, 0x01, 0xff, 0x80, 0x00, 0x01, 0x07, 0x01, 0x08, 0x45, 0x6e, 0x64, 0x70, 0x6f, 0x69, 0x6e, @@ -772,6 +777,15 @@ fn normalize_join_peers_for_local(local_peer: &PeerInfo, peers: BTreeMap SiteReplicationState { + let actual_peer = normalize_peer_info(actual_peer); + state + .peers + .retain(|_, peer| !same_endpoint(&peer.endpoint, &actual_peer.endpoint)); + state.peers.insert(actual_peer.deployment_id.clone(), actual_peer); + state +} + async fn ensure_site_replicator_service_account(parent_user: &str, state: &SiteReplicationState) -> S3Result<(String, String)> { let Some(iam_sys) = get_global_iam_sys() else { return Err(s3_error!(InvalidRequest, "iam not init")); @@ -1440,6 +1454,17 @@ fn update_peer(mut state: SiteReplicationState, incoming: PeerInfo, ilm_expiry_o state } +fn sync_state_name_for_local_peer( + mut state: SiteReplicationState, + local_peer: &PeerInfo, + incoming: &PeerInfo, +) -> SiteReplicationState { + if same_endpoint(&incoming.endpoint, &local_peer.endpoint) && !incoming.name.is_empty() { + state.name = incoming.name.clone(); + } + state +} + fn edit_state(mut state: SiteReplicationState, incoming: PeerInfo, ilm_expiry_override: Option) -> SiteReplicationState { if let Some(enabled) = ilm_expiry_override { for peer in state.peers.values_mut() { @@ -1947,7 +1972,7 @@ impl Operation for SiteReplicationAddHandler { let sites: Vec = read_site_replication_json(req, &cred.secret_key, true).await?; let (service_account_access_key, service_account_secret_key) = ensure_site_replicator_service_account(&cred.access_key, ¤t_state).await?; - let state = merge_add_sites( + let mut state = merge_add_sites( current_state, local_peer.clone(), sites.clone(), @@ -1973,7 +1998,7 @@ impl Operation for SiteReplicationAddHandler { let mut peer_join_req = join_req.clone(); peer_join_req.svc_acct_parent = site.access_key.clone(); - send_peer_admin_request( + let body = send_peer_admin_request( &site.endpoint, SITE_REPLICATION_PEER_JOIN_PATH, &site.access_key, @@ -1981,6 +2006,14 @@ impl Operation for SiteReplicationAddHandler { &peer_join_req, ) .await?; + + let join_response: SRPeerJoinResponse = serde_json::from_slice(&body).map_err(|e| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("parse peer join response from {} failed: {e}", site.endpoint), + ) + })?; + state = reconcile_peer_with_actual_identity(state, join_response.peer); } persist_site_replication_state(&state).await?; @@ -2180,9 +2213,11 @@ impl Operation for SRPeerJoinHandler { .get(&local_peer.deployment_id) .map(|peer| peer.name.clone()) .filter(|name| !name.is_empty()) - .unwrap_or(local_peer.name); + .unwrap_or_else(|| local_peer.name.clone()); persist_site_replication_state(&state).await?; - Ok(empty_response(StatusCode::OK)) + json_response(&SRPeerJoinResponse { + peer: state.peers.get(&local_peer.deployment_id).cloned().unwrap_or(local_peer), + }) } } @@ -2416,7 +2451,8 @@ impl Operation for SRPeerEditHandler { incoming.name = local_peer.name.clone(); } } - let state = update_peer(state, incoming, ilm_expiry_override); + let state = + sync_state_name_for_local_peer(update_peer(state, incoming.clone(), ilm_expiry_override), &local_peer, &incoming); save_site_replication_state(&state).await?; Ok(empty_response(StatusCode::OK)) } @@ -2655,6 +2691,57 @@ mod tests { assert!(normalized.contains_key("hash-remote")); } + #[test] + fn test_reconcile_peer_with_actual_identity_replaces_endpoint_hash_key() { + let mut state = SiteReplicationState::default(); + state.peers.insert( + "local".to_string(), + PeerInfo { + deployment_id: "local".to_string(), + ..peer("local", "https://local.example.com") + }, + ); + state.peers.insert( + "hash-remote".to_string(), + PeerInfo { + deployment_id: "hash-remote".to_string(), + ..peer("remote", "https://remote.example.com") + }, + ); + + let reconciled = reconcile_peer_with_actual_identity( + state, + PeerInfo { + deployment_id: "real-remote".to_string(), + ..peer("remote", "https://remote.example.com/") + }, + ); + + assert!(reconciled.peers.contains_key("local")); + assert!(reconciled.peers.contains_key("real-remote")); + assert!(!reconciled.peers.contains_key("hash-remote")); + } + + #[test] + fn test_sync_state_name_for_local_peer_updates_top_level_name() { + let mut state = SiteReplicationState { + name: "old-local".to_string(), + ..Default::default() + }; + let local_peer = PeerInfo { + deployment_id: "local".to_string(), + ..peer("old-local", "https://local.example.com") + }; + let incoming = PeerInfo { + deployment_id: "local".to_string(), + ..peer("new-local", "https://local.example.com/") + }; + + state = sync_state_name_for_local_peer(state, &local_peer, &incoming); + + assert_eq!(state.name, "new-local"); + } + #[test] fn test_site_replication_state_requires_remote_peer_to_be_enabled() { let mut state = SiteReplicationState::default();