diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 8ebf3c02a..fce86d24b 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -9573,6 +9573,10 @@ 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 7dada3492..781f5c3c0 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -75,7 +75,8 @@ use crate::auth::get_condition_values_with_client_info; use crate::error::ApiError; use crate::shared_types::RemoteAddr; use crate::site_replication::{ - site_replication_bucket_meta_hook, site_replication_delete_bucket_hook, site_replication_make_bucket_hook, + 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; use http::StatusCode; @@ -1397,7 +1398,8 @@ impl DefaultBucketUsecase { authorize_request(&mut req, Action::S3Action(S3Action::ForceDeleteBucketAction)).await?; } - store + let replication_delete_intent = prepare_site_replication_delete_bucket_hook(&input.bucket, force).await?; + let delete_result = store .delete_bucket( &input.bucket, &DeleteBucketOptions { @@ -1406,7 +1408,13 @@ impl DefaultBucketUsecase { }, ) .await - .map_err(ApiError::from)?; + .map_err(ApiError::from); + if let Err(err) = delete_result { + if let Some(intent) = replication_delete_intent.as_ref() { + cancel_site_replication_delete_bucket_hook(intent).await; + } + return Err(err.into()); + } // 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 @@ -1421,7 +1429,9 @@ impl DefaultBucketUsecase { // Re-evaluate lifecycle and replication after bucket removal. rustfs_scanner::record_scanner_maintenance_change(&input.bucket); - if let Err(err) = site_replication_delete_bucket_hook(&input.bucket, force).await { + if let Some(intent) = replication_delete_intent + && let Err(err) = finish_site_replication_delete_bucket_hook(intent).await + { warn!(bucket = %input.bucket, error = ?err, "site replication delete bucket hook failed"); } diff --git a/rustfs/src/site_replication/hooks.rs b/rustfs/src/site_replication/hooks.rs index fc8b1764f..1ce8d38ed 100644 --- a/rustfs/src/site_replication/hooks.rs +++ b/rustfs/src/site_replication/hooks.rs @@ -393,33 +393,30 @@ 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 sends = runtime.state.peers.values().filter_map(|peer| { - if peer.deployment_id == runtime.local_peer.deployment_id - || same_identity_endpoint(&peer.endpoint, &runtime.local_peer.endpoint) - { - return None; +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) + .with_client(&transport.client) + .send(&runtime.service_account_secret_key, &serde_json::json!({})) + .await } - Some(async move { - let result = async { - let transport = PeerTransport::for_runtime_peer(peer).await?; - PeerAdminRequest::put(&transport.connection, path, &runtime.state.service_account_access_key) - .with_client(&transport.client) - .send(&runtime.service_account_secret_key, &serde_json::json!({})) - .await + .await; + match result { + Ok(_) => { + dequeue_site_replication_retry_event(peer, path).await; + None } - .await; - match result { - Ok(_) => { - dequeue_site_replication_retry_event(peer, path).await; - None - } - Err(err) => { - enqueue_site_replication_retry_event(peer, path, &err).await; - Some(err) - } + Err(err) => { + enqueue_site_replication_retry_event(peer, path, &err).await; + Some(err) } - }) + } }); futures::future::join_all(sends) .await @@ -429,25 +426,39 @@ async fn broadcast_site_replication_destructive_bucket_op(runtime: &SiteReplicat .map_or(Ok(()), Err) } -pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: bool) -> S3Result<()> { +pub(crate) struct SiteReplicationDeleteBucketIntent { + runtime: SiteReplicationRuntime, + peers: Vec, + reserved_peers: Vec, + path: String, +} + +pub(crate) async fn prepare_site_replication_delete_bucket_hook( + bucket: &str, + force_delete: bool, +) -> S3Result> { let operation = if force_delete { "force-delete-bucket" } else { "delete-bucket" }; + // Keep concurrent delete reservations independent so rolling back one + // failed local delete cannot clear another request's durable intent. + let retry_intent = Uuid::new_v4().to_string(); let path = format!( "/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", form_urlencoded::Serializer::new(String::new()) .append_pair("bucket", bucket) .append_pair("operation", operation) + .append_pair("retryIntent", &retry_intent) .finish() ); let Some(runtime) = runtime_site_replication_targets().await? else { - return Ok(()); + return Ok(None); }; let store = current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?; - let retry_peers = runtime + let peers = runtime .state .peers .values() @@ -457,24 +468,41 @@ pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: boo }) .cloned() .collect::>(); - let retry_path = path.clone(); + 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 { - broadcast_site_replication_destructive_bucket_op(&runtime, &path).await + prequeue_site_replication_destructive_events(&reserve_peers, &reserve_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) - } + let reserved_peers = match result { + Ok(result) => result?, + Err(err) => return Err(ApiError::from(err).into()), + }; + Ok(Some(SiteReplicationDeleteBucketIntent { + runtime, + peers, + reserved_peers, + path, + })) +} + +pub(crate) async fn cancel_site_replication_delete_bucket_hook(intent: &SiteReplicationDeleteBucketIntent) { + for peer in &intent.reserved_peers { + dequeue_site_replication_retry_event(peer, &intent.path).await; } } +pub(crate) async fn finish_site_replication_delete_bucket_hook(intent: SiteReplicationDeleteBucketIntent) -> S3Result<()> { + let store = + current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?; + with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move { + broadcast_site_replication_destructive_bucket_op(&intent.runtime, &intent.peers, &intent.path).await + }) + .await + .map_err(ApiError::from)? +} + pub async fn site_replication_bucket_meta_hook(mut item: SRBucketMeta) -> S3Result<()> { let Some(runtime) = runtime_site_replication_targets().await? else { return Ok(()); diff --git a/rustfs/src/site_replication/retry.rs b/rustfs/src/site_replication/retry.rs index 990230b14..8fd52b4e0 100644 --- a/rustfs/src/site_replication/retry.rs +++ b/rustfs/src/site_replication/retry.rs @@ -16,6 +16,11 @@ 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 @@ -42,7 +47,7 @@ pub(crate) struct SiteReplicationRetryEvent { #[serde(default, skip_serializing_if = "Option::is_none")] pub(crate) edit_generation: Option, /// The latest delivery failure happened before an authenticated peer - /// response was received (connect, timeout, DNS, or TLS). Such failures + /// response was received (connect, DNS, or TLS). Such failures /// may bypass the expensive replay backoff only after a cheap devnull /// reachability probe proves the peer is back. #[serde(default, skip_serializing_if = "std::ops::Not::not")] @@ -246,10 +251,7 @@ pub(crate) fn upsert_site_replication_retry_event( pub(crate) fn retry_error_indicates_peer_unreachable(error: &str) -> bool { let error = error.to_ascii_lowercase(); - error.contains("failed (connect)") - || error.contains("failed (timeout)") - || error.contains("failed (dns resolution)") - || error.contains("failed (tls handshake)") + error.contains("failed (connect)") || error.contains("failed (dns resolution)") || error.contains("failed (tls handshake)") } pub(crate) fn retry_stats_for_state(state: &SiteReplicationState) -> Option { @@ -274,6 +276,46 @@ pub(crate) async fn enqueue_site_replication_retry_event(peer: &PeerInfo, path: enqueue_site_replication_retry_event_for_generation(peer, path, error, None).await } +pub(crate) fn prequeue_site_replication_destructive_events_in_state( + state: &mut SiteReplicationState, + peers: &[PeerInfo], + path: &str, +) -> S3Result> { + let missing = peers + .iter() + .filter(|peer| { + state.peers.contains_key(&peer.deployment_id) + && !state.retry_queue.iter().any(|event| retry_event_matches(event, peer, path)) + }) + .cloned() + .collect::>(); + if state.retry_queue.len().saturating_add(missing.len()) > SITE_REPLICATION_RETRY_QUEUE_LIMIT { + return Err(S3Error::with_message( + S3ErrorCode::InternalError, + "site replication retry queue is full; destructive operation was not reserved".to_string(), + )); + } + for peer in &missing { + upsert_site_replication_retry_event( + &mut state.retry_queue, + peer, + path, + "destructive site replication bucket operation pending delivery", + None, + ); + } + Ok(missing) +} + +pub(crate) async fn prequeue_site_replication_destructive_events(peers: &[PeerInfo], path: &str) -> S3Result> { + if peers.is_empty() { + return Ok(Vec::new()); + } + let peers = peers.to_vec(); + let path = path.to_string(); + update_site_replication_state(move |state| prequeue_site_replication_destructive_events_in_state(state, &peers, &path)).await +} + pub(crate) async fn enqueue_site_replication_retry_event_for_generation( peer: &PeerInfo, path: &str, @@ -918,6 +960,12 @@ pub(crate) fn bucket_op_retry_replay_tasks<'a>( format!("site replication retry plan has no configure operation for bucket {bucket:?}"), )); } + tasks.extend( + plan.bucket_items + .iter() + .filter(|item| item.bucket == bucket) + .map(SiteReplicationRepairTask::BucketMetadata), + ); tasks.extend(configure_tasks); Ok(tasks) } @@ -1051,10 +1099,10 @@ pub(crate) fn actionable_site_replication_retry_events( /// backoff exists to spare a *dead* peer the expensive replay (plan build, /// snapshot resend) — it must not delay convergence to a peer that has /// already RECOVERED, or a failure window ends in up to a day of silent -/// divergence (backlog#2071). Transport failures may be probed before the -/// normal replay backoff elapses; application failures still wait at least -/// one base interval so a reachable peer that keeps rejecting a replay is not -/// hammered faster than before. +/// divergence (backlog#2071). Peer connection failures may be probed before +/// the normal replay backoff elapses; request timeouts and application +/// failures still wait at least one base interval so a reachable peer that +/// keeps rejecting a replay is not hammered faster than before. pub(crate) fn deferred_site_replication_retry_events( state: &SiteReplicationState, now: OffsetDateTime, @@ -1293,11 +1341,17 @@ 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); - promote_reachable_deferred_retry_events(&runtime, &actionable, deferred).await?; - if actionable.is_empty() { - return Ok(()); + 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(()), } - drain_site_replication_retry_queue_locked(runtime, actionable).await }) .await .map_err(ApiError::from)? diff --git a/rustfs/src/site_replication/tests.rs b/rustfs/src/site_replication/tests.rs index 44b1a17e1..65e7e8ee4 100644 --- a/rustfs/src/site_replication/tests.rs +++ b/rustfs/src/site_replication/tests.rs @@ -764,6 +764,18 @@ fn test_bucket_make_retry_replays_matching_configure_before_settlement() { "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=videos&operation=configure-replication".to_string(), configure_photos.clone(), ], + bucket_items: vec![ + SRBucketMeta { + bucket: "videos".to_string(), + r#type: "tags".to_string(), + ..Default::default() + }, + SRBucketMeta { + bucket: "photos".to_string(), + r#type: "policy".to_string(), + ..Default::default() + }, + ], ..Default::default() }; @@ -771,10 +783,15 @@ fn test_bucket_make_retry_replays_matching_configure_before_settlement() { .expect("make retry plan should include its configure follow-up"); assert_eq!( tasks.iter().map(SiteReplicationRepairTask::path).collect::>(), - vec![make_photos.as_str(), configure_photos.as_str()] + vec![ + make_photos.as_str(), + "/rustfs/admin/v3/site-replication/peer/bucket-meta", + configure_photos.as_str() + ] ); assert!(matches!(tasks[0], SiteReplicationRepairTask::BucketMake(_))); - assert!(matches!(tasks[1], SiteReplicationRepairTask::Replication(_))); + assert!(matches!(&tasks[1], SiteReplicationRepairTask::BucketMetadata(item) if item.bucket == "photos")); + assert!(matches!(tasks[2], SiteReplicationRepairTask::Replication(_))); } #[test] @@ -816,6 +833,35 @@ fn test_reachable_probe_promotion_is_fenced_by_the_observed_event() { assert!(state.retry_queue[0].peer_unreachable); } +#[test] +fn test_destructive_bucket_op_prequeue_does_not_evict_existing_liability() { + let peer_b = PeerInfo { + deployment_id: "peer-b".to_string(), + ..peer("peer-b", "https://peer-b.example.com") + }; + let mut state = SiteReplicationState::default(); + state.peers.insert(peer_b.deployment_id.clone(), peer_b.clone()); + for index in 0..SITE_REPLICATION_RETRY_QUEUE_LIMIT { + let existing_peer = peer(&format!("existing-{index}"), &format!("https://existing-{index}.example.com")); + upsert_site_replication_retry_event( + &mut state.retry_queue, + &existing_peer, + &format!("/existing/{index}"), + "failed (connect)", + None, + ); + } + let existing_ids = state.retry_queue.iter().map(|event| event.id.clone()).collect::>(); + let path = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=delete-bucket"; + + let err = prequeue_site_replication_destructive_events_in_state(&mut state, &[peer_b], path) + .expect_err("a destructive operation must not evict existing retry liability"); + + assert_eq!(err.code(), &S3ErrorCode::InternalError); + assert_eq!(state.retry_queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT); + assert_eq!(state.retry_queue.iter().map(|event| event.id.clone()).collect::>(), existing_ids); +} + #[test] fn test_retry_snapshot_fingerprint_detects_concurrent_iam_change() { let old = SRIAMItem { @@ -897,7 +943,7 @@ fn test_site_replication_retry_backoff_schedule() { } #[test] -fn test_retry_error_marks_peer_unreachable_only_for_transport_failures() { +fn test_retry_error_marks_peer_unreachable_only_for_connection_failures() { let mut queue = Vec::new(); let peer = peer("remote", "https://remote.example.com"); let bucket_make = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning"; @@ -911,6 +957,18 @@ fn test_retry_error_marks_peer_unreachable_only_for_transport_failures() { ); assert!(queue[0].peer_unreachable); + upsert_site_replication_retry_event( + &mut queue, + &peer, + bucket_make, + "peer request to https://remote.example.com failed (timeout): request exceeded 10 seconds", + None, + ); + assert!( + !queue[0].peer_unreachable, + "a whole-request timeout does not prove the peer is unreachable" + ); + upsert_site_replication_retry_event( &mut queue, &peer,