fix(replication): correct peer joins and remote-state reporting (#7650)

* fix(replication): propagate verified peer deployment identities

* fix(replication): report actual remote peer state

* fix(replication): defer initial sync until all peers join

* test(replication): shut down TLS fixtures cleanly

---------

Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
Jason Kossis
2026-09-11 03:12:04 -04:00
committed by GitHub
parent f02bc947cd
commit 853ae63b6a
3 changed files with 543 additions and 80 deletions
@@ -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?;
+486 -65
View File
@@ -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<String, SRInfo>,
reachable_peers: &HashSet<String>,
) -> BTreeMap<String, SRStateInfo> {
// 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<SRStatusInfo> {
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<PeerSite>,
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<bool>) -> 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<SRPe
return Ok(SRPeerJoinResponse {
peer: fallback_peer,
initial_sync_error_message: String::new(),
initial_sync_deferred: false,
applied: None,
});
}
@@ -6801,6 +6843,7 @@ impl Operation for SiteReplicationAddHandler {
current_state,
local_peer.clone(),
sites.clone(),
&preflight_infos,
service_account_access_key.clone(),
admin_access_key,
replicate_ilm_expiry,
@@ -6815,66 +6858,86 @@ impl Operation for SiteReplicationAddHandler {
updated_at: state.updated_at,
},
defer_sync_state_enable: true,
defer_initial_sync: true,
};
let peer_join_path = with_site_replication_bootstrap_token(
SITE_REPLICATION_PEER_JOIN_PATH,
&add_in_progress_guard.token.to_string(),
);
let mut joined_endpoints = HashSet::new();
// Install every site's service account before any receiver probes or
// backfills to a peer that may not have joined yet. Reuse the join
// snapshot for the second pass without repeating IAM/topology writes.
// Only peers acknowledging deferral get a second request; older
// receivers retain their one-pass behavior and reported errors.
let initial_sync_path = format!("{peer_join_path}&initial-sync=true");
let mut initial_sync_errors = SiteReplicationErrorSummary::default();
for (site, preflight) in sites.iter().zip(preflight_infos.iter()) {
if same_identity_endpoint(&site.endpoint, &local_peer.endpoint)
|| !joined_endpoints.insert(site_identity_key(&site.endpoint))
{
continue;
}
let mut deferred_endpoints = HashSet::new();
for (path, defer_initial_sync) in [(&peer_join_path, true), (&initial_sync_path, false)] {
let mut joined_endpoints = HashSet::new();
for (site, preflight) in sites.iter().zip(preflight_infos.iter()) {
if same_identity_endpoint(&site.endpoint, &local_peer.endpoint)
|| (!defer_initial_sync && !deferred_endpoints.contains(&site_identity_key(&site.endpoint)))
|| !joined_endpoints.insert(site_identity_key(&site.endpoint))
{
continue;
}
let mut peer_join_req = join_req.clone();
peer_join_req.request.svc_acct_parent = site.access_key.clone();
let connection = PeerConnection::try_from(site)?;
let body = PeerAdminRequest::put(&connection, &peer_join_path, &site.access_key)
.send(&site.secret_key, &peer_join_req)
.await?;
let mut peer_join_req = join_req.clone();
peer_join_req.defer_initial_sync = defer_initial_sync;
peer_join_req.request.svc_acct_parent = site.access_key.clone();
let connection = PeerConnection::try_from(site)?;
let body = PeerAdminRequest::put(&connection, path, &site.access_key)
.send(&site.secret_key, &peer_join_req)
.await?;
let mut fallback_peer = existing_peer_for_endpoint(&state, &site.endpoint)
.unwrap_or_else(|| normalize_peer_site(site.clone(), replicate_ilm_expiry));
fallback_peer.deployment_id = preflight.deployment_id.clone();
let join_response = parse_peer_join_response(&body, fallback_peer).map_err(|e| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("parse peer join response from {} failed: {e}", site.endpoint),
)
})?;
if !join_response.initial_sync_error_message.is_empty() {
initial_sync_errors.push(format!("{}: {}", site.endpoint, join_response.initial_sync_error_message));
let mut fallback_peer = existing_peer_for_endpoint(&state, &site.endpoint)
.unwrap_or_else(|| normalize_peer_site(site.clone(), replicate_ilm_expiry));
fallback_peer.deployment_id = preflight.deployment_id.clone();
let join_response = parse_peer_join_response(&body, fallback_peer).map_err(|e| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("parse peer join response from {} failed: {e}", site.endpoint),
)
})?;
if join_response.initial_sync_deferred {
if defer_initial_sync {
deferred_endpoints.insert(site_identity_key(&site.endpoint));
} else {
initial_sync_errors.push(format!("{}: peer did not complete initial sync", site.endpoint));
}
}
if !join_response.initial_sync_error_message.is_empty() {
initial_sync_errors.push(format!("{}: {}", site.endpoint, join_response.initial_sync_error_message));
}
// An explicit no-op join. The peer answered 200 but wrote nothing —
// its persisted state is already newer than the snapshot it was
// sent — so the add is only PARTIALLY configured and saying
// "configured successfully" would be a lie (rustfs/rustfs#5963).
// `None` (a MinIO peer, or one older than the field) is not a
// no-op signal and is deliberately not reported.
if join_response.applied == Some(false) {
let phase = if defer_initial_sync { "join" } else { "initial sync" };
initial_sync_errors.push(format!(
"{}: peer did not apply the {phase} (its site replication state is newer than the snapshot it was sent); \
the site is not configured against this peer",
site.endpoint
));
}
state = reconcile_peer_with_actual_identity(state, join_response.peer);
let reconciled_peer = existing_peer_for_endpoint(&state, &site.endpoint).ok_or_else(|| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("peer join response from {} did not identify the requested site", site.endpoint),
)
})?;
validate_proposed_peer(&reconciled_peer).map_err(|err| {
S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!("invalid peer join response from {}: {err}", site.endpoint),
)
})?;
}
// An explicit no-op join. The peer answered 200 but wrote nothing —
// its persisted state is already newer than the snapshot it was
// sent — so the add is only PARTIALLY configured and saying
// "configured successfully" would be a lie (rustfs/rustfs#5963).
// `None` (a MinIO peer, or one older than the field) is not a
// no-op signal and is deliberately not reported.
if join_response.applied == Some(false) {
initial_sync_errors.push(format!(
"{}: peer did not apply the join (its site replication state is newer than the snapshot it was sent); \
the site is not configured against this peer",
site.endpoint
));
}
state = reconcile_peer_with_actual_identity(state, join_response.peer);
let reconciled_peer = existing_peer_for_endpoint(&state, &site.endpoint).ok_or_else(|| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("peer join response from {} did not identify the requested site", site.endpoint),
)
})?;
validate_proposed_peer(&reconciled_peer).map_err(|err| {
S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!("invalid peer join response from {}: {err}", site.endpoint),
)
})?;
}
mark_unknown_peer_sync_enabled(&mut state.peers);
@@ -7148,6 +7211,24 @@ impl Operation for SiteReplicationNetPerfHandler {
pub struct SRPeerJoinHandler {}
fn ensure_initial_sync_join_current(state: &SiteReplicationState, join_req: &SRPeerJoinReq) -> 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<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
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::<Vec<_>>();
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::<Vec<_>>(), state.peers.keys().collect::<Vec<_>>());
}
}
#[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");
+36 -1
View File
@@ -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");