From 618858adb4f116f03dcb78c11f5dba4631cc015a Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 5 Sep 2026 10:23:50 +0800 Subject: [PATCH] fix(site-replication): serialize retry replay state --- rustfs/src/admin/handlers/site_replication.rs | 4 ++ rustfs/src/site_replication/hooks.rs | 32 +++++++++- rustfs/src/site_replication/retry.rs | 59 ++++++++++++++----- 3 files changed, 78 insertions(+), 17 deletions(-) diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 8ebf3c02a..253920ae1 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -9569,6 +9569,10 @@ mod tests { fn test_peer_timeout_constants_bound_unreachable_peer_probes() { assert_eq!(SITE_REPLICATION_PEER_REQUEST_TIMEOUT, Duration::from_secs(10)); assert_eq!(SITE_REPLICATION_PEER_CONNECT_TIMEOUT, Duration::from_secs(3)); + assert!( + SITE_REPLICATION_RETRY_PROBE_BATCH_TIMEOUT < SITE_REPLICATION_LIFECYCLE_LOCK_TIMEOUT, + "the entire retry probe batch must finish before lifecycle waiters time out" + ); assert!( SITE_REPLICATION_LIFECYCLE_LOCK_TIMEOUT >= SITE_REPLICATION_PEER_REQUEST_TIMEOUT, "a waiter must not give up before the holder's single wedged peer probe can finish" diff --git a/rustfs/src/site_replication/hooks.rs b/rustfs/src/site_replication/hooks.rs index cd62704ee..57b3e29c9 100644 --- a/rustfs/src/site_replication/hooks.rs +++ b/rustfs/src/site_replication/hooks.rs @@ -406,7 +406,37 @@ pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: boo .append_pair("operation", operation) .finish() ); - broadcast_site_replication_json(&path, &serde_json::json!({})).await + let Some(runtime) = runtime_site_replication_targets().await? else { + return Ok(()); + }; + let store = + current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?; + let retry_peers = runtime + .state + .peers + .values() + .filter(|peer| { + peer.deployment_id != runtime.local_peer.deployment_id + && !same_identity_endpoint(&peer.endpoint, &runtime.local_peer.endpoint) + }) + .cloned() + .collect::>(); + let retry_path = path.clone(); + let result = + with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move { + broadcast_site_replication_json_with_runtime(&runtime, &path, &serde_json::json!({})).await + }) + .await; + match result { + Ok(result) => result, + Err(err) => { + let err: S3Error = ApiError::from(err).into(); + for peer in &retry_peers { + enqueue_site_replication_retry_event(peer, &retry_path, &err).await; + } + Err(err) + } + } } pub async fn site_replication_bucket_meta_hook(mut item: SRBucketMeta) -> S3Result<()> { diff --git a/rustfs/src/site_replication/retry.rs b/rustfs/src/site_replication/retry.rs index fc17aaf01..0139d03e7 100644 --- a/rustfs/src/site_replication/retry.rs +++ b/rustfs/src/site_replication/retry.rs @@ -16,6 +16,10 @@ use super::*; pub(crate) const SITE_REPLICATION_RETRY_QUEUE_LIMIT: usize = 256; +pub(crate) const SITE_REPLICATION_RETRY_PROBE_CONCURRENCY: usize = 4; + +pub(crate) const SITE_REPLICATION_RETRY_PROBE_BATCH_TIMEOUT: Duration = Duration::from_secs(20); + /// Attempts before an entry reports as `failed` in retryStats. Visibility /// only: a `failed` entry stays drain-eligible, and the reachability probe /// short-circuits its backoff once the peer answers again — so an early @@ -1105,6 +1109,8 @@ pub(crate) async fn promote_reachable_deferred_retry_events( actionable: &mut Vec, deferred: Vec, ) { + use futures::StreamExt as _; + let due_peers: HashSet = actionable.iter().map(|event| event.peer_deployment_id.clone()).collect(); let mut deferred_by_peer: BTreeMap> = BTreeMap::new(); for event in deferred { @@ -1118,21 +1124,28 @@ pub(crate) async fn promote_reachable_deferred_retry_events( .or_default() .push(event); } - for (deployment_id, events) in deferred_by_peer { - let Some(peer) = runtime.state.peers.get(&deployment_id) else { - continue; - }; + let probes = deferred_by_peer.into_iter().filter_map(|(deployment_id, events)| { + let peer = runtime.state.peers.get(&deployment_id)?; if deployment_id == runtime.local_peer.deployment_id || same_identity_endpoint(&peer.endpoint, &runtime.local_peer.endpoint) { - continue; + return None; } - if probe_site_replication_peer_reachable(runtime, peer).await { + Some(async move { + let reachable = probe_site_replication_peer_reachable(runtime, peer).await; + (deployment_id, peer.endpoint.clone(), events, reachable) + }) + }); + let mut probes = futures::stream::iter(probes).buffer_unordered(SITE_REPLICATION_RETRY_PROBE_CONCURRENCY); + let deadline = tokio::time::Instant::now() + SITE_REPLICATION_RETRY_PROBE_BATCH_TIMEOUT; + while let Ok(Some((deployment_id, peer_endpoint, events, reachable))) = tokio::time::timeout_at(deadline, probes.next()).await + { + if reachable { info!( component = LOG_COMPONENT_ADMIN, subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, event = EVENT_ADMIN_SITE_REPLICATION_STATE, - peer = %peer.endpoint, + peer = %peer_endpoint, deployment_id = %deployment_id, promoted = events.len(), result = "retry_backoff_probe_promoted", @@ -1214,7 +1227,7 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> { // escalated markers are exactly the entries the drain skips. log_site_replication_retry_liabilities(&runtime.state); let now = OffsetDateTime::now_utc(); - let mut actionable = actionable_site_replication_retry_events(&runtime.state, now); + let actionable = actionable_site_replication_retry_events(&runtime.state, now); let deferred = deferred_site_replication_retry_events(&runtime.state, now); if actionable.is_empty() && deferred.is_empty() { return Ok(()); @@ -1231,21 +1244,35 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> { // guard) may have started since. Re-check on the fresh state. return Ok(()); } - // Probe before taking the repair lock: probes are read-only peer traffic - // and a dead peer's connect timeout must not hold the lock. - promote_reachable_deferred_retry_events(&runtime, &mut actionable, deferred).await; - if actionable.is_empty() { - return Ok(()); - } // Serialize against operator repair execution. This does NOT close the // dry-run -> execute window (dry-run takes no lock): a drain settling a // replayable bucket-op entry in that window changes the preflight token // and execute fails safe with "preflight is stale" — the operator // re-runs the dry-run. Lock order matches repair: lifecycle guard (held // by the reconcile tick) -> repair execution lock -> state object lock - // inside the send bookkeeping. An operator repair holding the lock makes - // this tick skip after the lock-acquire timeout. + // inside the send bookkeeping. The lock also elects one server to probe + // and replay the queue; after acquiring it, reload state so a settled + // event or deleted bucket cannot be replayed from this admission snapshot. with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move { + // Runtime and queue snapshots captured before this distributed lock + // are only admission hints. Another node may have settled the event, + // or a local bucket may have been deleted, while this node waited. + let Some(runtime) = runtime_site_replication_targets().await? else { + return Ok(()); + }; + if runtime.state.pending_endpoint_refresh.is_some() + || runtime.state.pending_remove.is_some() + || runtime.state.pending_rotation.is_some() + { + return Ok(()); + } + let now = OffsetDateTime::now_utc(); + let mut actionable = actionable_site_replication_retry_events(&runtime.state, now); + let deferred = deferred_site_replication_retry_events(&runtime.state, now); + promote_reachable_deferred_retry_events(&runtime, &mut actionable, deferred).await; + if actionable.is_empty() { + return Ok(()); + } drain_site_replication_retry_queue_locked(runtime, actionable).await }) .await