diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index 48dfd8cf5..bebecff9c 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -6383,6 +6383,18 @@ async fn test_site_replication_edit_and_status_peer_state_real_three_node() -> R let relayed_key = "after-edit-from-relay.txt"; let relayed_payload = b"site replication after endpoint edit from relay".to_vec(); + // The first joining receiver owns data before the third site has the + // shared account. Initial probes and backfill must wait for every join. + target_client.create_bucket().bucket(bucket).send().await?; + enable_bucket_versioning(&target_env, bucket).await?; + target_client + .put_object() + .bucket(bucket) + .key(baseline_key) + .body(ByteStream::from(baseline_payload.clone())) + .send() + .await?; + let add_status = site_replication_add( &source_env, &[ @@ -6410,7 +6422,10 @@ async fn test_site_replication_edit_and_status_peer_state_real_three_node() -> R ], ) .await?; - assert!(add_status.success, "unexpected site add result: {:?}", add_status); + assert!( + add_status.success && add_status.err_detail.is_empty() && add_status.initial_sync_error_message.is_empty(), + "unexpected site add result: {add_status:?}" + ); let source_info = wait_for_site_replication_enabled(&source_env, 3).await?; let _target_info = wait_for_site_replication_enabled(&target_env, 3).await?; @@ -6421,19 +6436,11 @@ async fn test_site_replication_edit_and_status_peer_state_real_three_node() -> R .find(|peer| peer.endpoint == target_env.url) .ok_or("target peer missing from source site replication info")?; - source_client.create_bucket().bucket(bucket).send().await?; - enable_bucket_versioning(&source_env, bucket).await?; - wait_for_bucket_on_target(&target_client, bucket).await?; - wait_for_bucket_on_target(&relay_client, bucket).await?; - source_client - .put_object() - .bucket(bucket) - .key(baseline_key) - .body(ByteStream::from(baseline_payload.clone())) - .send() - .await?; - let replicated_baseline = wait_for_object_on_target(&target_client, bucket, baseline_key).await?; - assert_eq!(replicated_baseline, baseline_payload); + for client in [&source_client, &relay_client] { + wait_for_bucket_on_target(client, bucket).await?; + let backfilled = wait_for_object_on_target(client, bucket, baseline_key).await?; + assert_eq!(backfilled, baseline_payload); + } let old_target_address = target_env.address.clone(); let new_target_port = RustFSTestEnvironment::find_available_port().await?; diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index c288dc047..40f008983 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -44,7 +44,8 @@ use crate::admin::utils::{empty_response, json_response, read_compatible_admin_b use crate::error::ApiError; use crate::server::ADMIN_PREFIX; use crate::site_replication::identity::{ - canonical_endpoint, is_https_endpoint, mark_unknown_peer_sync_enabled, same_identity_endpoint, site_identity_key, + canonical_endpoint, deployment_id_for_endpoint, is_https_endpoint, mark_unknown_peer_sync_enabled, same_identity_endpoint, + site_identity_key, }; use crate::storage::storage_api::{lock_bucket_targets_metadata, with_config_object_write_lock}; use base64_simd::URL_SAFE_NO_PAD; @@ -312,6 +313,8 @@ struct SRPeerJoinResponse { peer: PeerInfo, #[serde(rename = "initialSyncErrorMessage", default, skip_serializing_if = "String::is_empty")] initial_sync_error_message: String, + #[serde(rename = "initialSyncDeferred", default, skip_serializing_if = "std::ops::Not::not")] + initial_sync_deferred: bool, /// Whether the receiving site actually applied this join. /// /// Three-valued on purpose. `None` means the peer did not report — MinIO @@ -330,6 +333,8 @@ struct SRPeerJoinEnvelope { request: SRPeerJoinReq, #[serde(rename = "deferSyncStateEnable", default, skip_serializing_if = "std::ops::Not::not")] defer_sync_state_enable: bool, + #[serde(rename = "deferInitialSync", default, skip_serializing_if = "std::ops::Not::not")] + defer_initial_sync: bool, } #[derive(Debug, Default)] @@ -2792,6 +2797,19 @@ fn prune_in_sync_status_details(status: &mut SRStatusInfo, opts: &SRStatusOption } } +fn peer_states_from_infos( + site_infos: BTreeMap, + reachable_peers: &HashSet, +) -> BTreeMap { + // Failed metainfo fetches leave default entries in site_infos for comparison; + // they must not become fabricated peer state. PeerErrors describes the failure. + site_infos + .into_iter() + .filter(|(deployment_id, _)| reachable_peers.contains(deployment_id)) + .map(|(deployment_id, info)| (deployment_id, info.state)) + .collect() +} + async fn build_status_info(state: &SiteReplicationState, local_peer: &PeerInfo, uri: &Uri) -> S3Result { let opts = sr_status_options(uri); let mut local_info = Some(filter_sr_info(build_sr_info(state, local_peer).await?, &opts)); @@ -2920,17 +2938,7 @@ async fn build_status_info(state: &SiteReplicationState, local_peer: &PeerInfo, } if opts.peer_state { - for (deployment_id, peer) in &state.peers { - status.peer_states.insert( - deployment_id.clone(), - SRStateInfo { - name: peer.name.clone(), - peers: state.peers.clone(), - updated_at: state.updated_at, - api_version: Some(SITE_REPL_API_VERSION.to_string()), - }, - ); - } + status.peer_states = peer_states_from_infos(site_infos, &reachable_peers); } Ok(status) @@ -2940,6 +2948,7 @@ fn merge_add_sites( mut state: SiteReplicationState, local_peer: PeerInfo, sites: Vec, + preflight_infos: &[SiteReplicationAddPreflightInfo], service_account_access_key: String, service_account_parent: String, replicate_ilm_expiry: bool, @@ -2949,11 +2958,36 @@ fn merge_add_sites( state.service_account_parent = service_account_parent; state.updated_at = Some(OffsetDateTime::now_utc()); state.peers = build_join_peers(&state, &local_peer, sites, replicate_ilm_expiry); + // Every join must carry the verified identities, including peers that + // have not joined yet. Fixing only the coordinator after each reply + // leaves the other sites holding endpoint-derived placeholders. + for info in preflight_infos { + if let Some(mut peer) = existing_peer_for_endpoint(&state, &info.endpoint) { + peer.deployment_id = info.deployment_id.clone(); + state = reconcile_peer_with_actual_identity(state, peer); + } + } state } fn update_peer(mut state: SiteReplicationState, incoming: PeerInfo, ilm_expiry_override: Option) -> SiteReplicationState { let mut peer = normalize_peer_info(incoming); + // An older sender may still hold a placeholder after this site has + // learned the real ID. Do not let that delivery downgrade the identity. + if peer.deployment_id == deployment_id_for_endpoint(&peer.endpoint) + && let Some(existing) = state.peers.values().find(|existing| { + same_identity_endpoint(&existing.endpoint, &peer.endpoint) + && existing.deployment_id != deployment_id_for_endpoint(&existing.endpoint) + }) + { + peer.deployment_id = existing.deployment_id.clone(); + } + // Remove the placeholder before persistence normalizes duplicate + // endpoints; otherwise map ordering can discard the real identity. + state.peers.retain(|_, existing| { + !same_identity_endpoint(&existing.endpoint, &peer.endpoint) + || existing.deployment_id != deployment_id_for_endpoint(&existing.endpoint) + }); if let Some(enabled) = ilm_expiry_override { peer.replicate_ilm_expiry = enabled; } @@ -3539,6 +3573,13 @@ fn align_peer_edit_deployment_id(state: &SiteReplicationState, incoming: &mut Pe return; }; if matches.next().is_none() { + if same_identity_endpoint(&peer.endpoint, &incoming.endpoint) + && peer.deployment_id == deployment_id_for_endpoint(&peer.endpoint) + && !incoming.deployment_id.is_empty() + && incoming.deployment_id != deployment_id_for_endpoint(&incoming.endpoint) + { + return; + } incoming.deployment_id = peer.deployment_id.clone(); } } @@ -6710,6 +6751,7 @@ fn parse_peer_join_response(body: &[u8], fallback_peer: PeerInfo) -> Result S3Result<()> { + if !state.enabled() + || join_req.updated_at.is_none() + || state.updated_at != join_req.updated_at + || state.service_account_access_key.is_empty() + || state.service_account_access_key != join_req.svc_acct_access_key + || state.pending_remove.is_some() + || state.pending_rotation.is_some() + || pending_endpoint_refresh(state).is_some() + { + return Err(s3_error!( + InvalidRequest, + "site replication changed before initial sync; re-run replicate add" + )); + } + Ok(()) +} + /// What the join admission decided about an incoming peer join. The verdict — /// and the committed state the back-fill afterwards needs — travel out of /// [`admit_peer_join`] instead of being answered where they are decided. @@ -7310,6 +7391,7 @@ fn superseded_join_response(peer: PeerInfo) -> SRPeerJoinResponse { SRPeerJoinResponse { peer, initial_sync_error_message: String::new(), + initial_sync_deferred: false, applied: Some(false), } } @@ -7319,6 +7401,7 @@ fn applied_join_response(peer: PeerInfo, initial_sync_error_message: String) -> SRPeerJoinResponse { peer, initial_sync_error_message, + initial_sync_deferred: false, applied: Some(true), } } @@ -7328,17 +7411,30 @@ impl Operation for SRPeerJoinHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationAddAction).await?; let bootstrap_token = site_replication_bootstrap_token(&req.uri); + let initial_sync_only = query_pairs(&req.uri).get("initial-sync").is_some_and(|value| value == "true"); let local_endpoint = site_replication_local_endpoint(&req.uri, &req.headers); // The body is fully read before the admission takes the lifecycle // guard: a sender that stalls mid-body must not block this node's // add/remove/rotate/reconciler. let join_envelope: SRPeerJoinEnvelope = read_site_replication_json(req, &cred.secret_key, true).await?; let defer_sync_state_enable = join_envelope.defer_sync_state_enable; + let defer_initial_sync = join_envelope.defer_initial_sync; let join_req = join_envelope.request; validate_join_peer_snapshot(&join_req.peers)?; - let committed = - admit_peer_join(local_endpoint, join_req, defer_sync_state_enable, apply_peer_join_service_account).await?; + let _initial_sync_guard = if initial_sync_only { + Some(SiteReplicationLifecycleGuard::acquire().await?) + } else { + None + }; + let committed = if initial_sync_only { + let state = load_site_replication_state().await?; + ensure_initial_sync_join_current(&state, &join_req)?; + let local_peer = local_peer_at_endpoint(local_endpoint, &state); + PeerJoinOutcome::Applied(Box::new(state), local_peer) + } else { + admit_peer_join(local_endpoint, join_req, defer_sync_state_enable, apply_peer_join_service_account).await? + }; // Committed; the reverse-reachability probe and the bucket back-fill // run outside the transaction — their transport helpers' retry-event // bookkeeping re-enters it (P1-15). @@ -7355,6 +7451,12 @@ impl Operation for SRPeerJoinHandler { return json_response(StatusCode::OK, &superseded_join_response(peer)); } }; + if defer_initial_sync && !initial_sync_only { + let mut response = + applied_join_response(state.peers.get(&local_peer.deployment_id).cloned().unwrap_or(local_peer), String::new()); + response.initial_sync_deferred = true; + return json_response(StatusCode::OK, &response); + } // 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. Per-bucket // failures are logged (BUG2) so a reverse-direction back-fill gap is observable. @@ -8482,9 +8584,92 @@ impl Operation for SRRotateServiceAccountHandler { #[cfg(test)] mod tests { use super::*; - use crate::site_replication::identity::deployment_id_for_endpoint; use rustfs_madmin::SRSessionPolicy; + #[test] + fn peer_states_preserve_each_sites_actual_membership_and_metadata() { + let local = SRStateInfo { + name: "local".to_string(), + peers: BTreeMap::from([( + "actual-remote".to_string(), + PeerInfo { + deployment_id: "actual-remote".to_string(), + endpoint: "http://remote:9000".to_string(), + ..Default::default() + }, + )]), + updated_at: Some(OffsetDateTime::UNIX_EPOCH), + api_version: Some("1".to_string()), + }; + let remote = SRStateInfo { + name: "remote-reported-name".to_string(), + peers: BTreeMap::from([( + "legacy-placeholder".to_string(), + PeerInfo { + deployment_id: "legacy-placeholder".to_string(), + endpoint: "http://local:9000".to_string(), + ..Default::default() + }, + )]), + updated_at: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(10)), + api_version: None, + }; + let infos = BTreeMap::from([ + ( + "local".to_string(), + SRInfo { + state: local.clone(), + ..Default::default() + }, + ), + ( + "remote".to_string(), + SRInfo { + state: remote.clone(), + ..Default::default() + }, + ), + ]); + let states = peer_states_from_infos(infos, &HashSet::from(["local".to_string(), "remote".to_string()])); + assert_eq!(states.len(), 2); + assert_eq!(serde_json::to_value(&states["local"]).unwrap(), serde_json::to_value(local).unwrap()); + assert_eq!(serde_json::to_value(&states["remote"]).unwrap(), serde_json::to_value(remote).unwrap()); + } + + #[test] + fn peer_states_omit_unreachable_peers_instead_of_defaulting_them() { + let infos = BTreeMap::from([ + ("local".to_string(), SRInfo::default()), + ("offline".to_string(), SRInfo::default()), + ]); + let states = peer_states_from_infos(infos, &HashSet::from(["local".to_string()])); + assert_eq!(states.len(), 1); + assert!(states.contains_key("local")); + assert!(!states.contains_key("offline")); + } + + #[test] + fn peer_states_preserve_a_reachable_peers_empty_membership() { + let states = peer_states_from_infos( + BTreeMap::from([( + "remote".to_string(), + SRInfo { + enabled: false, + state: SRStateInfo { + name: "remote".to_string(), + ..Default::default() + }, + ..Default::default() + }, + )]), + &HashSet::from(["remote".to_string()]), + ); + assert_eq!(states["remote"].name, "remote"); + assert!(states["remote"].peers.is_empty()); + assert!(states["remote"].updated_at.is_none()); + assert!(states["remote"].api_version.is_none()); + } + /// A peer the status probe could not reach must render as offline. /// /// Regression: `build_metrics_summary` used to hardcode `online: true` and @@ -10916,6 +11101,7 @@ mod tests { secret_key: "remote-sk".to_string(), ..PeerSite::default() }], + &[], "svc-ak".to_string(), "root".to_string(), true, @@ -10949,6 +11135,7 @@ mod tests { ..PeerSite::default() }, ], + &[], "svc-ak".to_string(), "root".to_string(), true, @@ -12155,6 +12342,173 @@ mod tests { assert!(normalized.contains_key("hash-remote")); } + #[test] + fn test_peer_identity_join_snapshot_uses_verified_ids() { + let actual = ["site-a", "site-b", "site-c"].map(|name| PeerInfo { + deployment_id: format!("{name}-deployment"), + ..peer(name, &format!("https://{name}.example.com:9000")) + }); + let preflight = actual + .iter() + .map(|peer| preflight_site("reported-name", &peer.endpoint, &peer.deployment_id, 0)) + .collect::>(); + let sites = actual + .iter() + .map(|peer| PeerSite { + name: peer.name.clone(), + endpoint: peer.endpoint.clone(), + skip_tls_verify: true, + ca_cert_pem: "requested-ca".to_string(), + ..Default::default() + }) + .collect(); + let state = merge_add_sites( + SiteReplicationState::default(), + actual[0].clone(), + sites, + &preflight, + "svc-ak".to_string(), + "root".to_string(), + false, + ); + assert_eq!(state.peers.len(), actual.len()); + for expected in &actual { + let stored = state + .peers + .get(&expected.deployment_id) + .expect("verified ID in initial join map"); + assert_eq!(stored.name, expected.name); + assert_eq!(stored.endpoint, expected.endpoint); + assert!(stored.skip_tls_verify); + assert_eq!(stored.ca_cert_pem, "requested-ca"); + } + for local in &actual[1..] { + let mut joined = SiteReplicationState::default(); + apply_peer_join( + &mut joined, + local, + SRPeerJoinReq { + peers: state.peers.clone(), + ..Default::default() + }, + true, + ); + assert_eq!(joined.peers.keys().collect::>(), state.peers.keys().collect::>()); + } + } + + #[test] + fn test_peer_identity_legacy_edit_does_not_restore_placeholder() { + let actual = PeerInfo { + deployment_id: "actual-remote".to_string(), + ..peer("remote", "http://remote.example.com:9000") + }; + for name in ["remote", ""] { + let state = SiteReplicationState { + peers: BTreeMap::from([(actual.deployment_id.clone(), actual.clone())]), + ..Default::default() + }; + let mut incoming = PeerInfo { + deployment_id: deployment_id_for_endpoint("https://REMOTE.example.com:9000/"), + sync_state: SyncStatus::Enable, + ..peer(name, "https://REMOTE.example.com:9000/") + }; + align_peer_edit_deployment_id(&state, &mut incoming); + let state = update_peer(state, incoming, None); + assert_eq!(state.peers.len(), 1); + assert_eq!(state.peers[&actual.deployment_id].sync_state, SyncStatus::Enable); + } + } + + #[test] + fn test_peer_identity_finalization_repairs_legacy_three_site_join() { + let actual = ["site-a", "site-b", "site-c"].map(|name| PeerInfo { + deployment_id: format!("{name}-deployment"), + ..peer(name, &format!("http://{name}.example.com:9000")) + }); + let sites = actual + .iter() + .map(|peer| PeerSite { + name: peer.name.clone(), + endpoint: peer.endpoint.clone(), + ..Default::default() + }) + .collect(); + let mut coordinator = merge_add_sites( + SiteReplicationState::default(), + actual[0].clone(), + sites, + &[], + "svc-ak".to_string(), + "root".to_string(), + true, + ); + let join = SRPeerJoinReq { + peers: coordinator.peers.clone(), + ..Default::default() + }; + for remote in &actual[1..] { + coordinator = reconcile_peer_with_actual_identity(coordinator, remote.clone()); + } + mark_unknown_peer_sync_enabled(&mut coordinator.peers); + + for local in &actual[1..] { + let mut state = SiteReplicationState::default(); + apply_peer_join(&mut state, local, join.clone(), true); + for mut incoming in coordinator.peers.values().cloned() { + align_peer_edit_deployment_id(&state, &mut incoming); + state = apply_internal_peer_edit(state, local, incoming, None).expect("finalize peer identity"); + } + assert_eq!(state.peers.len(), actual.len(), "finalization must not retain placeholder peers"); + for expected in &actual { + let stored = existing_peer_for_endpoint(&state, &expected.endpoint).expect("peer remains configured"); + assert_eq!(stored.deployment_id, expected.deployment_id, "observer: {}", local.name); + assert_eq!(stored.sync_state, SyncStatus::Enable); + } + } + } + + #[test] + fn test_peer_identity_edit_replaces_placeholder_for_canonical_endpoint() { + let local = PeerInfo { + deployment_id: "local-deployment".to_string(), + ..peer("local", "https://local.example.com:9000") + }; + let endpoint = "http://remote.example.com:9000"; + let placeholder = PeerInfo { + deployment_id: deployment_id_for_endpoint(endpoint), + ..peer("remote", endpoint) + }; + for deployment_id in ["00000000-0000-4000-8000-000000000001", "ffffffff-ffff-4fff-bfff-ffffffffffff"] { + for name in ["remote", ""] { + for already_present in [false, true] { + let mut incoming = PeerInfo { + deployment_id: deployment_id.to_string(), + sync_state: SyncStatus::Enable, + ..peer(name, "https://REMOTE.example.com:9000/") + }; + let mut state = SiteReplicationState { + peers: BTreeMap::from([ + (local.deployment_id.clone(), local.clone()), + (placeholder.deployment_id.clone(), placeholder.clone()), + ]), + ..Default::default() + }; + if already_present { + state.peers.insert(incoming.deployment_id.clone(), incoming.clone()); + } + align_peer_edit_deployment_id(&state, &mut incoming); + let state = apply_internal_peer_edit(state, &local, incoming, None).expect("repair peer identity"); + assert_eq!(state.peers.len(), 2, "repair must replace, not duplicate, the placeholder"); + assert!(!state.peers.contains_key(&placeholder.deployment_id)); + assert!(state.peers.contains_key(deployment_id)); + let normalized = normalize_peer_map_by_identity(state.peers); + assert!(normalized.contains_key(deployment_id), "normalization must retain the actual ID"); + } + } + } + } + #[test] fn test_reconcile_peer_with_actual_identity_replaces_endpoint_hash_key() { let mut state = SiteReplicationState::default(); @@ -13477,6 +13831,7 @@ mod tests { assert_eq!(response.peer.deployment_id, "remote-deployment"); assert_eq!(response.peer.endpoint, "https://remote.example.com"); assert!(response.initial_sync_error_message.is_empty()); + assert!(!response.initial_sync_deferred); assert_eq!( response.applied, None, "a MinIO empty-body success reports nothing; it must not read as a no-op join" @@ -13486,6 +13841,7 @@ mod tests { let json = serde_json::to_vec(&SRPeerJoinResponse { peer: peer("actual", "https://actual.example.com"), initial_sync_error_message: "sync failed".to_string(), + initial_sync_deferred: false, applied: Some(true), }) .expect("serialize join response"); @@ -14338,6 +14694,69 @@ mod tests { assert_eq!(value.get("deferSyncStateEnable"), Some(&Value::Bool(true))); } + #[test] + fn test_initial_sync_join_requires_the_committed_snapshot() { + let now = OffsetDateTime::now_utc(); + let mut state = SiteReplicationState { + updated_at: Some(now), + service_account_access_key: "replicator".to_string(), + peers: BTreeMap::from([ + ("a".to_string(), peer("a", "https://a.example.com")), + ("b".to_string(), peer("b", "https://b.example.com")), + ]), + ..Default::default() + }; + let mut request = SRPeerJoinReq { + updated_at: Some(now), + svc_acct_access_key: "replicator".to_string(), + ..Default::default() + }; + ensure_initial_sync_join_current(&state, &request).expect("same committed join"); + for timestamp in [None, Some(now - time::Duration::SECOND), Some(now + time::Duration::SECOND)] { + request.updated_at = timestamp; + assert!(ensure_initial_sync_join_current(&state, &request).is_err()); + } + request.updated_at = Some(now); + request.svc_acct_access_key = "another-replicator".to_string(); + assert!(ensure_initial_sync_join_current(&state, &request).is_err()); + request.svc_acct_access_key.clone_from(&state.service_account_access_key); + state.pending_remove = Some(PendingRemove::default()); + assert!(ensure_initial_sync_join_current(&state, &request).is_err()); + state.pending_remove = None; + state.pending_rotation = Some(PendingRotation::default()); + assert!(ensure_initial_sync_join_current(&state, &request).is_err()); + state.pending_rotation = None; + state.pending_endpoint_refresh = Some(PendingEndpointRefresh::default()); + assert!(ensure_initial_sync_join_current(&state, &request).is_err()); + state.pending_endpoint_refresh = None; + state.service_account_access_key.clear(); + request.svc_acct_access_key.clear(); + assert!(ensure_initial_sync_join_current(&state, &request).is_err()); + state.service_account_access_key = "replicator".to_string(); + request.svc_acct_access_key.clone_from(&state.service_account_access_key); + state.peers.clear(); + assert!(ensure_initial_sync_join_current(&state, &request).is_err()); + } + + #[test] + fn test_join_initial_sync_deferral_preserves_legacy_requests() { + let legacy: SRPeerJoinEnvelope = serde_json::from_str("{}").expect("legacy join"); + assert!(!legacy.defer_initial_sync); + assert!(serde_json::to_value(&legacy).unwrap().get("deferInitialSync").is_none()); + let deferred: SRPeerJoinEnvelope = serde_json::from_str(r#"{"deferInitialSync":true}"#).expect("deferred join"); + assert!(deferred.defer_initial_sync); + assert_eq!(serde_json::to_value(deferred).unwrap()["deferInitialSync"], true); + let mut response = applied_join_response(peer("b", "https://b.example.com"), String::new()); + assert!(serde_json::to_value(&response).unwrap().get("initialSyncDeferred").is_none()); + response.initial_sync_deferred = true; + let wire = serde_json::to_vec(&response).unwrap(); + assert!( + parse_peer_join_response(&wire, PeerInfo::default()) + .unwrap() + .initial_sync_deferred + ); + } + // BUG2: pre-existing-bucket back-fill failures must be surfaced in the add response's // initial_sync_error_message, not swallowed behind an unqualified success. #[test] @@ -14390,6 +14809,7 @@ mod tests { let value = serde_json::to_value(SRPeerJoinResponse { peer: peer("remote", "https://remote.example.com"), initial_sync_error_message: "bucket setup failed".to_string(), + initial_sync_deferred: false, applied: Some(true), }) .expect("serialize peer join response"); @@ -14401,6 +14821,7 @@ mod tests { let value = serde_json::to_value(SRPeerJoinResponse { peer: peer("remote", "https://remote.example.com"), initial_sync_error_message: String::new(), + initial_sync_deferred: false, applied: None, }) .expect("serialize peer join response"); diff --git a/rustfs/src/site_replication/tests.rs b/rustfs/src/site_replication/tests.rs index 42568392c..7701c233f 100644 --- a/rustfs/src/site_replication/tests.rs +++ b/rustfs/src/site_replication/tests.rs @@ -98,11 +98,46 @@ async fn spawn_test_tls_server_with_response(response: &'static [u8]) -> (String break; } } - stream.write_all(response).await.is_ok() + // Flush buffered TLS records and send close_notify before dropping the socket. + stream.write_all(response).await.is_ok() && stream.shutdown().await.is_ok() }); (endpoint, ca_pem, task) } +#[tokio::test] +async fn tls_test_server_delivers_response_and_closes_cleanly() { + use rustls_pki_types::pem::PemObject; + + let (endpoint, ca_pem, server) = spawn_test_tls_server().await; + let mut roots = rustls::RootCertStore::empty(); + roots + .add(rustls_pki_types::CertificateDer::from_pem_slice(ca_pem.as_bytes()).expect("parse test CA")) + .expect("trust test CA"); + let config = rustls::ClientConfig::builder() + .with_root_certificates(roots) + .with_no_client_auth(); + let connector = tokio_rustls::TlsConnector::from(Arc::new(config)); + let socket = tokio::net::TcpStream::connect(endpoint.strip_prefix("https://").expect("TLS endpoint")) + .await + .expect("connect to TLS test server"); + let mut stream = connector + .connect(rustls_pki_types::ServerName::try_from("127.0.0.1").expect("test server name"), socket) + .await + .expect("trust TLS test server"); + stream + .write_all(b"GET / HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n") + .await + .expect("write test request"); + stream.flush().await.expect("flush test request"); + let mut response = Vec::new(); + tokio::time::timeout(Duration::from_secs(5), stream.read_to_end(&mut response)) + .await + .expect("TLS response must finish") + .expect("TLS test server must send close_notify before closing"); + assert!(response.ends_with(b"\r\n\r\nok")); + assert!(server.await.expect("TLS test server task")); +} + #[test] fn peer_connection_validation_accepts_supported_combinations() { let ca = valid_test_ca_pem("peer.example.com");