fix(site-replication): bound retry recovery rounds

This commit is contained in:
cxymds
2026-09-05 10:40:10 +08:00
parent b4b68060c9
commit cd810723b7
3 changed files with 167 additions and 58 deletions
+21 -23
View File
@@ -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::<Vec<_>>();
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::<Vec<_>>();
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)
}
}
+93 -35
View File
@@ -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<SiteReplicationRetryEvent>,
actionable: &[SiteReplicationRetryEvent],
deferred: Vec<SiteReplicationRetryEvent>,
) {
) -> S3Result<()> {
let due_peers: HashSet<String> = actionable.iter().map(|event| event.peer_deployment_id.clone()).collect();
let mut deferred_by_peer: BTreeMap<String, Vec<SiteReplicationRetryEvent>> = 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 {
+53
View File
@@ -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::<Vec<_>>();
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 {