From cd810723b7e051d7c1198c1fdfe866e2ae358069 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 5 Sep 2026 10:40:10 +0800 Subject: [PATCH] fix(site-replication): bound retry recovery rounds --- rustfs/src/site_replication/hooks.rs | 44 +++++---- rustfs/src/site_replication/retry.rs | 128 +++++++++++++++++++-------- rustfs/src/site_replication/tests.rs | 53 +++++++++++ 3 files changed, 167 insertions(+), 58 deletions(-) diff --git a/rustfs/src/site_replication/hooks.rs b/rustfs/src/site_replication/hooks.rs index 8ff6079c1..76ef318c8 100644 --- a/rustfs/src/site_replication/hooks.rs +++ b/rustfs/src/site_replication/hooks.rs @@ -393,20 +393,12 @@ pub(crate) async fn broadcast_site_replication_make_bucket( broadcast_site_replication_json_using_runtime(runtime, &configure_path, &serde_json::json!({})).await } -async fn broadcast_site_replication_destructive_bucket_op(runtime: &SiteReplicationRuntime, path: &str) -> S3Result<()> { - let 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::>(); - prequeue_site_replication_destructive_events(&peers, path).await?; - let mut first_error = None; - for peer in &peers { +async fn broadcast_site_replication_destructive_bucket_op( + runtime: &SiteReplicationRuntime, + peers: &[PeerInfo], + path: &str, +) -> S3Result<()> { + let sends = peers.iter().map(|peer| async move { let result = async { let transport = PeerTransport::for_runtime_peer(peer).await?; PeerAdminRequest::put(&transport.connection, path, &runtime.state.service_account_access_key) @@ -416,14 +408,22 @@ async fn broadcast_site_replication_destructive_bucket_op(runtime: &SiteReplicat } .await; match result { - Ok(_) => dequeue_site_replication_retry_event(peer, path).await, + Ok(_) => { + dequeue_site_replication_retry_event(peer, path).await; + None + } Err(err) => { enqueue_site_replication_retry_event(peer, path, &err).await; - first_error.get_or_insert(err); + Some(err) } } - } - first_error.map_or(Ok(()), Err) + }); + futures::future::join_all(sends) + .await + .into_iter() + .flatten() + .next() + .map_or(Ok(()), Err) } pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: bool) -> S3Result<()> { @@ -454,19 +454,17 @@ pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: boo }) .cloned() .collect::>(); - let retry_path = path.clone(); + prequeue_site_replication_destructive_events(&retry_peers, &path).await?; + let locked_peers = retry_peers; let result = with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move { - broadcast_site_replication_destructive_bucket_op(&runtime, &path).await + broadcast_site_replication_destructive_bucket_op(&runtime, &locked_peers, &path).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) } } diff --git a/rustfs/src/site_replication/retry.rs b/rustfs/src/site_replication/retry.rs index e96811554..2d269276d 100644 --- a/rustfs/src/site_replication/retry.rs +++ b/rustfs/src/site_replication/retry.rs @@ -16,7 +16,7 @@ use super::*; pub(crate) const SITE_REPLICATION_RETRY_QUEUE_LIMIT: usize = 256; -pub(crate) const SITE_REPLICATION_RETRY_DRAIN_BATCH_TIMEOUT: Duration = Duration::from_secs(20); +pub(crate) const SITE_REPLICATION_RETRY_DRAIN_BATCH_TIMEOUT: Duration = Duration::from_secs(25); /// Attempts before an entry reports as `failed` in retryStats. Visibility /// only: a `failed` entry stays drain-eligible, and the reachability probe @@ -63,6 +63,13 @@ pub(crate) fn retry_event_matches(event: &SiteReplicationRetryEvent, peer: &Peer (event.peer_deployment_id == peer.deployment_id || event.peer_endpoint == peer.endpoint) && event.path == path } +pub(crate) fn retry_event_is_destructive_bucket_op(event: &SiteReplicationRetryEvent) -> bool { + matches!( + retry_bucket_operation(&event.path).as_deref(), + Some("delete-bucket" | "force-delete-bucket" | "purge-deleted-bucket") + ) +} + pub(crate) const SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH: &str = "internal:retry-snapshot:iam"; pub(crate) const SITE_REPLICATION_RETRY_BUCKET_METADATA_SNAPSHOT_PATH: &str = "internal:retry-snapshot:bucket-metadata"; @@ -241,8 +248,15 @@ pub(crate) fn upsert_site_replication_retry_event( deletions_recorded: false, }); if queue.len() > SITE_REPLICATION_RETRY_QUEUE_LIMIT { - let overflow = queue.len() - SITE_REPLICATION_RETRY_QUEUE_LIMIT; - queue.drain(0..overflow); + while queue.len() > SITE_REPLICATION_RETRY_QUEUE_LIMIT { + let Some(index) = queue.iter().position(|event| !retry_event_is_destructive_bucket_op(event)) else { + // Destructive rows are a durable outbox, not a replay cache. + // Preserve every unacknowledged peer even if that temporarily + // exceeds the soft limit so none becomes a silent orphan. + break; + }; + queue.remove(index); + } } } @@ -1133,11 +1147,31 @@ pub(crate) async fn probe_site_replication_peer_reachable(runtime: &SiteReplicat /// every peer that answers. A probe failure advances nothing: retry counts /// only move on real delivery attempts, so the per-event backoff is intact /// when the peer is genuinely down. +pub(crate) fn mark_reachable_deferred_retry_events( + state: &mut SiteReplicationState, + recovered: &[SiteReplicationRetryEvent], +) -> usize { + let mut promoted = 0; + for recovered in recovered { + if let Some(current) = state.retry_queue.iter_mut().find(|current| { + current.id == recovered.id + && current.peer_deployment_id == recovered.peer_deployment_id + && current.path == recovered.path + && current.updated_at == recovered.updated_at + }) { + current.updated_at = None; + current.peer_unreachable = false; + promoted += 1; + } + } + promoted +} + pub(crate) async fn promote_reachable_deferred_retry_events( runtime: &SiteReplicationRuntime, - actionable: &mut Vec, + actionable: &[SiteReplicationRetryEvent], deferred: Vec, -) { +) -> S3Result<()> { 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 { @@ -1163,6 +1197,7 @@ pub(crate) async fn promote_reachable_deferred_retry_events( (deployment_id, peer.endpoint.clone(), events, reachable) }) }); + let mut recovered = Vec::new(); for (deployment_id, peer_endpoint, events, reachable) in futures::future::join_all(probes).await { if reachable { info!( @@ -1173,11 +1208,23 @@ pub(crate) async fn promote_reachable_deferred_retry_events( deployment_id = %deployment_id, promoted = events.len(), result = "retry_backoff_probe_promoted", - "peer reachable again; replaying its backed-off retry events this tick" + "peer reachable again; promoting backed-off retry events for the next tick" ); - actionable.extend(events); + recovered.extend(events); } } + if recovered.is_empty() { + return Ok(()); + } + update_site_replication_state_when_changed(move |state| { + let promoted = mark_reachable_deferred_retry_events(state, &recovered); + Ok(if promoted == 0 { + StateCommit::Unchanged(()) + } else { + StateCommit::Changed(()) + }) + }) + .await } /// Operator-visible per-tick alert for retry entries that no longer converge @@ -1291,10 +1338,10 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> { return Ok(()); } 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); let drain = async { - promote_reachable_deferred_retry_events(&runtime, &mut actionable, deferred).await; + promote_reachable_deferred_retry_events(&runtime, &actionable, deferred).await?; if actionable.is_empty() { return Ok(()); } @@ -1333,39 +1380,50 @@ pub(crate) async fn drain_site_replication_retry_queue_locked( .push(event); } - let mut settled = 0usize; - let mut failures = 0usize; - for (deployment_id, peer_events) in events_by_peer { - let Some(peer) = runtime.state.peers.get(&deployment_id) else { - continue; - }; + let runtime = Arc::new(runtime); + let plan = plan.map(Arc::new); + let peer_replays = events_by_peer.into_iter().filter_map(|(deployment_id, peer_events)| { + let peer = runtime.state.peers.get(&deployment_id)?.clone(); if deployment_id == runtime.local_peer.deployment_id || same_identity_endpoint(&peer.endpoint, &runtime.local_peer.endpoint) { - continue; + return None; } - let transport = match PeerTransport::for_runtime_peer(peer).await { - Ok(transport) => transport, - Err(err) => { - // Record the attempt so backoff advances for an unreachable - // peer instead of re-dialing it every tick. - for event in &peer_events { - enqueue_site_replication_retry_event(peer, &event.path, &err).await; + let runtime = Arc::clone(&runtime); + let plan = plan.as_ref().map(Arc::clone); + Some(async move { + let mut settled = 0usize; + let mut failures = 0usize; + let transport = match PeerTransport::for_runtime_peer(&peer).await { + Ok(transport) => transport, + Err(err) => { + // Record the attempt so backoff advances for an unreachable + // peer instead of re-dialing it every tick. + for event in &peer_events { + enqueue_site_replication_retry_event(&peer, &event.path, &err).await; + } + return (0, peer_events.len()); } - failures += peer_events.len(); - continue; - } - }; - for event in peer_events { - let Some(action) = classify_site_replication_retry_event(&event) else { - continue; }; - match drain_one_site_replication_retry_event(&runtime, peer, &transport, &event, action, plan.as_ref()).await { - Ok(true) => settled += 1, - Ok(false) => {} - Err(_) => failures += 1, + for event in peer_events { + let Some(action) = classify_site_replication_retry_event(&event) else { + continue; + }; + match drain_one_site_replication_retry_event(&runtime, &peer, &transport, &event, action, plan.as_deref()).await { + Ok(true) => settled += 1, + Ok(false) => {} + Err(_) => failures += 1, + } } - } + (settled, failures) + }) + }); + + let mut settled = 0usize; + let mut failures = 0usize; + for (peer_settled, peer_failures) in futures::future::join_all(peer_replays).await { + settled += peer_settled; + failures += peer_failures; } if settled > 0 || failures > 0 { diff --git a/rustfs/src/site_replication/tests.rs b/rustfs/src/site_replication/tests.rs index 8960634e1..78cb55d52 100644 --- a/rustfs/src/site_replication/tests.rs +++ b/rustfs/src/site_replication/tests.rs @@ -822,6 +822,59 @@ fn test_destructive_bucket_op_prequeues_every_remote_peer() { ); } +#[test] +fn test_destructive_retry_intents_survive_the_soft_queue_limit() { + let mut queue = (0..SITE_REPLICATION_RETRY_QUEUE_LIMIT) + .map(|index| { + drain_event( + &format!("peer-{index}"), + &format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=bucket-{index}&operation=delete-bucket"), + 1, + None, + ) + }) + .collect::>(); + let ordinary_peer = PeerInfo { + deployment_id: "ordinary".to_string(), + ..peer("ordinary", "https://ordinary.example.com") + }; + upsert_site_replication_retry_event(&mut queue, &ordinary_peer, SITE_REPLICATION_PEER_EDIT_PATH, "ordinary retry", None); + assert_eq!(queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT); + assert!(queue.iter().all(retry_event_is_destructive_bucket_op)); + + let extra_peer = PeerInfo { + deployment_id: "extra".to_string(), + ..peer("extra", "https://extra.example.com") + }; + let extra_path = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=extra&operation=force-delete-bucket"; + upsert_site_replication_retry_event(&mut queue, &extra_peer, extra_path, "pending destructive retry", None); + assert_eq!(queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT + 1); + assert!(queue.iter().any(|event| event.path == extra_path)); +} + +#[test] +fn test_reachable_probe_promotion_is_fenced_by_the_observed_event() { + let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp"); + let path = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning"; + let mut event = drain_event("remote", path, 3, Some(now)); + event.peer_unreachable = true; + let recovered = event.clone(); + let mut state = SiteReplicationState { + retry_queue: vec![event], + ..Default::default() + }; + + assert_eq!(mark_reachable_deferred_retry_events(&mut state, &[recovered.clone()]), 1); + assert_eq!(state.retry_queue[0].updated_at, None); + assert!(!state.retry_queue[0].peer_unreachable); + + state.retry_queue[0].updated_at = Some(now + time::Duration::seconds(1)); + state.retry_queue[0].peer_unreachable = true; + assert_eq!(mark_reachable_deferred_retry_events(&mut state, &[recovered]), 0); + assert_eq!(state.retry_queue[0].updated_at, Some(now + time::Duration::seconds(1))); + assert!(state.retry_queue[0].peer_unreachable); +} + #[test] fn test_retry_snapshot_fingerprint_detects_concurrent_iam_change() { let old = SRIAMItem {