diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index fce86d24b..50b9e99b9 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -1824,7 +1824,7 @@ async fn site_replication_reconcile_prerequisites_ready() -> bool { fn reconcile_site_replication_retry_drain() -> std::pin::Pin + Send>> { Box::pin(async { - let Some(_lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else { + let Some(lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else { return; }; if !site_replication_reconcile_prerequisites_ready().await { @@ -1839,6 +1839,10 @@ fn reconcile_site_replication_retry_drain() -> std::pin::Pin return, } + // Retry sends re-check membership under the bucket-op read lock for + // each bounded peer request. Do not hold the lifecycle guard across + // an arbitrarily large snapshot replay. + drop(lifecycle); drain_site_replication_retry_queue().await; }) } @@ -9573,10 +9577,6 @@ mod tests { 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" ); - assert!( - crate::site_replication::SITE_REPLICATION_RETRY_DRAIN_BATCH_TIMEOUT < SITE_REPLICATION_LIFECYCLE_LOCK_TIMEOUT, - "the retry drainer must release the lifecycle guard before operator waiters time out" - ); } #[tokio::test(start_paused = true)] diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index 781f5c3c0..1201e4a32 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -75,7 +75,7 @@ use crate::auth::get_condition_values_with_client_info; use crate::error::ApiError; use crate::shared_types::RemoteAddr; use crate::site_replication::{ - cancel_site_replication_delete_bucket_hook, finish_site_replication_delete_bucket_hook, + SITE_REPLICATION_BUCKET_OP_LOCK, cancel_site_replication_delete_bucket_hook, finish_site_replication_delete_bucket_hook, prepare_site_replication_delete_bucket_hook, site_replication_bucket_meta_hook, site_replication_make_bucket_hook, }; use crate::storage::storage_api::lock_bucket_targets_metadata; @@ -1398,6 +1398,7 @@ impl DefaultBucketUsecase { authorize_request(&mut req, Action::S3Action(S3Action::ForceDeleteBucketAction)).await?; } + let bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.read().await; let replication_delete_intent = prepare_site_replication_delete_bucket_hook(&input.bucket, force).await?; let delete_result = store .delete_bucket( @@ -1415,6 +1416,7 @@ impl DefaultBucketUsecase { } return Err(err.into()); } + drop(bucket_op_guard); // Drop every cached object body for the now-deleted bucket so dead // bytes do not sit resident until TTL. Covers both the normal and the diff --git a/rustfs/src/site_replication/hooks.rs b/rustfs/src/site_replication/hooks.rs index 1ce8d38ed..dfcc30db7 100644 --- a/rustfs/src/site_replication/hooks.rs +++ b/rustfs/src/site_replication/hooks.rs @@ -409,7 +409,7 @@ async fn broadcast_site_replication_destructive_bucket_op( .await; match result { Ok(_) => { - dequeue_site_replication_retry_event(peer, path).await; + dequeue_site_replication_destructive_retry_events(peer, path).await; None } Err(err) => { @@ -456,8 +456,6 @@ pub(crate) async fn prepare_site_replication_delete_bucket_hook( let Some(runtime) = runtime_site_replication_targets().await? else { return Ok(None); }; - let store = - current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?; let peers = runtime .state .peers @@ -468,17 +466,7 @@ pub(crate) async fn prepare_site_replication_delete_bucket_hook( }) .cloned() .collect::>(); - let reserve_peers = peers.clone(); - let reserve_path = path.clone(); - let result = - with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move { - prequeue_site_replication_destructive_events(&reserve_peers, &reserve_path).await - }) - .await; - let reserved_peers = match result { - Ok(result) => result?, - Err(err) => return Err(ApiError::from(err).into()), - }; + let reserved_peers = prequeue_site_replication_destructive_events(&peers, &path).await?; Ok(Some(SiteReplicationDeleteBucketIntent { runtime, peers, diff --git a/rustfs/src/site_replication/retry.rs b/rustfs/src/site_replication/retry.rs index 8fd52b4e0..e275717d9 100644 --- a/rustfs/src/site_replication/retry.rs +++ b/rustfs/src/site_replication/retry.rs @@ -16,11 +16,6 @@ use super::*; pub(crate) const SITE_REPLICATION_RETRY_QUEUE_LIMIT: usize = 256; -/// Keep the reconcile tick's lifecycle guard below the 30-second operator -/// wait bound. Successfully settled events are persisted as the batch runs, -/// so a later tick continues with the remaining work. -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 /// short-circuits its backoff once the peer answers again — so an early @@ -66,6 +61,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"; @@ -244,8 +246,12 @@ 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 { + break; + }; + queue.remove(index); + } } } @@ -307,6 +313,52 @@ pub(crate) fn prequeue_site_replication_destructive_events_in_state( Ok(missing) } +pub(crate) fn settle_site_replication_destructive_retry_events( + queue: &mut Vec, + peer: &PeerInfo, + path: &str, +) -> usize { + let Some(bucket) = retry_bucket_name(path) else { + return 0; + }; + let before = queue.len(); + queue.retain(|event| { + !(retry_event_matches_peer(event, peer) + && retry_event_is_destructive_bucket_op(event) + && retry_bucket_name(&event.path).as_deref() == Some(bucket.as_str())) + }); + before.saturating_sub(queue.len()) +} + +fn retry_event_matches_peer(event: &SiteReplicationRetryEvent, peer: &PeerInfo) -> bool { + event.peer_deployment_id == peer.deployment_id || event.peer_endpoint == peer.endpoint +} + +pub(crate) async fn dequeue_site_replication_destructive_retry_events(peer: &PeerInfo, path: &str) { + let peer_owned = peer.clone(); + let path_owned = path.to_string(); + let result = update_site_replication_state_when_changed(move |state| { + let removed = settle_site_replication_destructive_retry_events(&mut state.retry_queue, &peer_owned, &path_owned); + Ok(if removed == 0 { + StateCommit::Unchanged(()) + } else { + StateCommit::Changed(()) + }) + }) + .await; + if let Err(err) = result { + warn!( + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + peer = %peer.endpoint, + path, + error = ?err, + "failed to settle destructive site replication retry events" + ); + } +} + pub(crate) async fn prequeue_site_replication_destructive_events(peers: &[PeerInfo], path: &str) -> S3Result> { if peers.is_empty() { return Ok(Vec::new()); @@ -783,27 +835,69 @@ impl RetrySnapshot { } } - pub(crate) async fn send(&self, transport: &PeerTransport, access_key: &str, secret_key: &str) -> S3Result<()> { + pub(crate) async fn send( + &self, + peer: &PeerInfo, + transport: &PeerTransport, + access_key: &str, + secret_key: &str, + ) -> S3Result { match self { Self::Iam(items) => { for item in items { - SiteReplicationRepairTask::Iam(item) - .send(transport, access_key, secret_key) - .await?; + if !send_retry_task_if_peer_current( + peer, + &SiteReplicationRepairTask::Iam(item), + transport, + access_key, + secret_key, + ) + .await? + { + return Ok(false); + } } } Self::BucketMetadata(items) => { for item in items { - SiteReplicationRepairTask::BucketMetadata(item) - .send(transport, access_key, secret_key) - .await?; + if !send_retry_task_if_peer_current( + peer, + &SiteReplicationRepairTask::BucketMetadata(item), + transport, + access_key, + secret_key, + ) + .await? + { + return Ok(false); + } } } } - Ok(()) + Ok(true) } } +pub(crate) async fn send_retry_task_if_peer_current( + peer: &PeerInfo, + task: &SiteReplicationRepairTask<'_>, + transport: &PeerTransport, + access_key: &str, + secret_key: &str, +) -> S3Result { + let _bucket_op = SITE_REPLICATION_BUCKET_OP_LOCK.read().await; + let state = load_site_replication_state().await?; + let current = state + .peers + .get(&peer.deployment_id) + .is_some_and(|current| same_identity_endpoint(¤t.endpoint, &peer.endpoint)); + if !current { + return Ok(false); + } + task.send(transport, access_key, secret_key).await?; + Ok(true) +} + #[derive(Hash, PartialEq, Eq)] pub(crate) enum IamSnapshotKey { Policy(String), @@ -1316,13 +1410,14 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> { // guard) may have started since. Re-check on the fresh state. return Ok(()); } - // Serialize against operator repair execution. This does NOT close the + // Serialize against operator repair execution. Peer membership is + // re-checked under the bucket-op read lock immediately before each + // network request, so the caller need not hold the lifecycle guard while + // a large snapshot is replayed. 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. The lock also elects one server to probe + // re-runs the dry-run. 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 { @@ -1341,17 +1436,11 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> { let now = OffsetDateTime::now_utc(); 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, &actionable, deferred).await?; - if actionable.is_empty() { - return Ok(()); - } - drain_site_replication_retry_queue_locked(runtime, actionable).await - }; - match tokio::time::timeout(SITE_REPLICATION_RETRY_DRAIN_BATCH_TIMEOUT, drain).await { - Ok(result) => result, - Err(_) => Ok(()), + promote_reachable_deferred_retry_events(&runtime, &actionable, deferred).await?; + if actionable.is_empty() { + return Ok(()); } + drain_site_replication_retry_queue_locked(runtime, actionable).await }) .await .map_err(ApiError::from)? @@ -1472,12 +1561,21 @@ pub(crate) async fn drain_one_site_replication_retry_event( drop_corrupt_iam_deletion_replay(peer, &record.id).await; continue; }; - if let Err(err) = SiteReplicationRepairTask::Iam(&item) - .send(transport, access_key, secret_key) - .await + match send_retry_task_if_peer_current( + peer, + &SiteReplicationRepairTask::Iam(&item), + transport, + access_key, + secret_key, + ) + .await { - enqueue_site_replication_retry_event(peer, &event.path, &err).await; - return Err(err); + Ok(true) => {} + Ok(false) => return Ok(false), + Err(err) => { + enqueue_site_replication_retry_event(peer, &event.path, &err).await; + return Err(err); + } } replayed_record_ids.push(record.id.clone()); } @@ -1486,9 +1584,13 @@ pub(crate) async fn drain_one_site_replication_retry_event( let mut replay = current_snapshot.clone(); for _ in 0..SITE_REPLICATION_RETRY_SNAPSHOT_STABILITY_ATTEMPTS { let current_fingerprint = current_snapshot.fingerprint()?; - if let Err(err) = replay.send(transport, access_key, secret_key).await { - enqueue_site_replication_retry_event(peer, &event.path, &err).await; - return Err(err); + match replay.send(peer, transport, access_key, secret_key).await { + Ok(true) => {} + Ok(false) => return Ok(false), + Err(err) => { + enqueue_site_replication_retry_event(peer, &event.path, &err).await; + return Err(err); + } } let fresh_info = build_sr_info(&runtime.state, &runtime.local_peer).await?; let fresh_plan = site_replication_bootstrap_plan(&fresh_info)?; @@ -1533,9 +1635,13 @@ pub(crate) async fn drain_one_site_replication_retry_event( return Ok(true); } for task in &tasks { - if let Err(err) = task.send(transport, access_key, secret_key).await { - enqueue_site_replication_retry_event(peer, &event.path, &err).await; - return Err(err); + match send_retry_task_if_peer_current(peer, task, transport, access_key, secret_key).await { + Ok(true) => {} + Ok(false) => return Ok(false), + Err(err) => { + enqueue_site_replication_retry_event(peer, &event.path, &err).await; + return Err(err); + } } } dequeue_site_replication_retry_event(peer, &event.path).await; @@ -1564,6 +1670,15 @@ pub(crate) async fn drain_one_site_replication_retry_event( let edit_path = peer_edit_path_with_fence(local_deployment_id, generation); let delivery_fence = local_deployment_id.is_some().then_some(generation); for body in &bodies { + let _bucket_op = SITE_REPLICATION_BUCKET_OP_LOCK.read().await; + let current_state = load_site_replication_state().await?; + let current = current_state + .peers + .get(&peer.deployment_id) + .is_some_and(|current| same_identity_endpoint(¤t.endpoint, &peer.endpoint)); + if !current { + return Ok(false); + } if let Err(err) = PeerAdminRequest::put(&transport.connection, &edit_path, access_key) .with_client(&transport.client) .send(secret_key, body) diff --git a/rustfs/src/site_replication/tests.rs b/rustfs/src/site_replication/tests.rs index 65e7e8ee4..8936b3df2 100644 --- a/rustfs/src/site_replication/tests.rs +++ b/rustfs/src/site_replication/tests.rs @@ -862,6 +862,54 @@ fn test_destructive_bucket_op_prequeue_does_not_evict_existing_liability() { assert_eq!(state.retry_queue.iter().map(|event| event.id.clone()).collect::>(), existing_ids); } +#[test] +fn test_destructive_success_settles_equivalent_bucket_intents() { + let peer = peer("remote", "https://remote.example.com"); + let failed = drain_event( + "remote", + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=delete-bucket&retryIntent=a", + 1, + None, + ); + let current = drain_event( + "remote", + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=force-delete-bucket&retryIntent=b", + 1, + None, + ); + let unrelated = drain_event( + "remote", + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=videos&operation=delete-bucket&retryIntent=c", + 1, + None, + ); + let mut queue = vec![failed, current.clone(), unrelated.clone()]; + + assert_eq!(settle_site_replication_destructive_retry_events(&mut queue, &peer, ¤t.path), 2); + assert_eq!(queue.len(), 1); + assert_eq!(queue[0].id, unrelated.id); +} + +#[test] +fn test_ordinary_retry_does_not_evict_destructive_liability() { + 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("ordinary", "https://ordinary.example.com"); + + upsert_site_replication_retry_event(&mut queue, &ordinary, SITE_REPLICATION_PEER_EDIT_PATH, "failed (connect)", None); + + assert_eq!(queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT); + assert!(queue.iter().all(retry_event_is_destructive_bucket_op)); +} + #[test] fn test_retry_snapshot_fingerprint_detects_concurrent_iam_change() { let old = SRIAMItem {