From bbd7b9ef17b5b9bce2dd2134fcfc18510bb8c5f8 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 5 Sep 2026 15:55:54 +0800 Subject: [PATCH] fix(site-replication): bound and order outage recovery (#7148) * fix(site-replication): wake retry drain after peer recovery * fix(site-replication): replay configure after bucket make * fix(site-replication): serialize retry replay state * fix(site-replication): persist destructive retry intents * fix(site-replication): bound retry recovery rounds * fix(site-replication): keep recovery replay live * fix(site-replication): preserve retry ordering * fix(site-replication): bound retry coordination * fix(site-replication): serialize topology replay * fix(site-replication): fence distributed retry state * fix(site-replication): bound outage retry drain * fix(site-replication): drop unsafe delete retry intents * fix(site-replication): order bucket mutation replay * fix(site-replication): harden outage retry replay * fix(site-replication): fence destructive peer delivery * fix(site-replication): avoid peer edit retry deadlock * fix(site-replication): fence retry error classification * fix(site-replication): classify connect timeouts * fix(site-replication): close recovery review races * test(site-replication): cover timeout endpoint text * fix(site-replication): close destructive recovery gaps * fix(site-replication): fence recovery revisions * fix(site-replication): replay bucket metadata on recovery * fix(site-replication): preserve s3gate boundary --------- Co-authored-by: overtrue --- .../src/replication_extension_test.rs | 93 ++ rustfs/src/admin/handlers/site_replication.rs | 662 +++++++++----- rustfs/src/app/bucket_usecase.rs | 100 ++- rustfs/src/site_replication/hooks.rs | 462 +++++++++- rustfs/src/site_replication/mod.rs | 13 +- rustfs/src/site_replication/repair.rs | 4 +- rustfs/src/site_replication/retry.rs | 849 +++++++++++++++--- rustfs/src/site_replication/state_lock.rs | 37 +- rustfs/src/site_replication/tests.rs | 609 ++++++++++++- rustfs/src/site_replication/transport.rs | 30 +- rustfs/src/site_replication_reconcile.rs | 48 +- rustfs/src/storage_api.rs | 4 +- 12 files changed, 2375 insertions(+), 536 deletions(-) diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index f5d308c7e..198fa5166 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -6743,6 +6743,99 @@ async fn test_site_replication_replicates_object_with_bucket_versioning_real_dua Ok(()) } +#[tokio::test] +async fn test_site_replication_replays_bucket_created_during_peer_outage_real_dual_node() -> TestResult { + init_logging(); + + // Keep compilation outside the scenario timeout. Recovery itself waits + // for the production 30-second lightweight retry tick. + let _rustfs_binary = rustfs_binary_path(); + + match timeout(Duration::from_secs(150), async { + let mut site_env = replication_fast_env(); + site_env.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV); + + let mut site_a_env = RustFSTestEnvironment::new().await?; + site_a_env.start_rustfs_server_with_env(vec![], &site_env).await?; + + let mut site_b_env = RustFSTestEnvironment::new().await?; + site_b_env.start_rustfs_server_without_cleanup_with_env(&site_env).await?; + + let site_a_client = site_a_env.create_s3_client(); + let site_b_client = site_b_env.create_s3_client(); + let bucket = "site-repl-peer-outage"; + let key = "after-recovery.txt"; + let payload = b"site replication recovered the missed bucket".to_vec(); + + let add_status = site_replication_add( + &site_a_env, + &[ + PeerSite { + name: "outage-site-a".to_string(), + endpoint: site_a_env.url.clone(), + access_key: site_a_env.access_key.clone(), + secret_key: site_a_env.secret_key.clone(), + ..Default::default() + }, + PeerSite { + name: "outage-site-b".to_string(), + endpoint: site_b_env.url.clone(), + access_key: site_b_env.access_key.clone(), + secret_key: site_b_env.secret_key.clone(), + ..Default::default() + }, + ], + ) + .await?; + assert!(add_status.success, "unexpected site add result: {add_status:?}"); + wait_for_site_replication_enabled(&site_a_env, 2).await?; + wait_for_site_replication_enabled(&site_b_env, 2).await?; + + site_b_env.stop_server(); + site_a_client.create_bucket().bucket(bucket).send().await?; + site_a_client.head_bucket().bucket(bucket).send().await?; + + let queued = site_replication_info(&site_a_env) + .await? + .retry_stats + .ok_or("peer outage did not persist a site replication retry event")?; + assert!(queued.pending + queued.failed > 0, "peer outage retry queue was unexpectedly empty"); + + site_b_env.restart_server_preserving_data(vec![], &site_env).await?; + let recovery_deadline = tokio::time::Instant::now() + Duration::from_secs(75); + loop { + let bucket_recovered = site_b_client.head_bucket().bucket(bucket).send().await.is_ok(); + let queue_empty = site_replication_info(&site_a_env).await?.retry_stats.is_none(); + if bucket_recovered && queue_empty { + break; + } + if tokio::time::Instant::now() >= recovery_deadline { + return Err(format!( + "site replication retry did not settle after peer recovery; bucket_recovered={bucket_recovered}, queue_empty={queue_empty}" + ) + .into()); + } + sleep(Duration::from_millis(250)).await; + } + + site_a_client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from(payload.clone())) + .send() + .await?; + assert_eq!(wait_for_object_on_target(&site_b_client, bucket, key).await?, payload); + + Ok(()) + }) + .await + { + Ok(result) => result, + Err(_) => Err("site replication peer-outage recovery timed out after 150 seconds".into()), + } +} + /// Re-applying a site's own replication config must not disable the peer's reverse direction. /// /// `PutBucketReplication` broadcasts the config to every peer — the console's replication diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index d75e3b850..51c7dda99 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -190,7 +190,8 @@ fn site_replicator_service_account_policy() -> S3Result { .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("parse site replicator policy failed: {e}"))) } -// Lock order: lifecycle -> bucket operation -> repair admission -> state -> per-bucket metadata. +// Lock order: lifecycle -> bucket-mutation admission -> per-bucket mutation +// -> bucket operation -> repair admission -> state -> per-bucket metadata. // "state" is the distributed state-object lock in // crate::site_replication::state_lock, entered through // update_site_replication_state (P1-15). There is no process-local state @@ -434,6 +435,7 @@ pub fn register_site_replication_route(r: &mut S3Router) -> std: // into this module: startup sits below this layer and must not depend upwards. The admin // router is built before startup reconciles, so the hook is always installed in time. crate::site_replication_reconcile::register_site_replication_reconciler(reconcile_site_replication_wiring); + crate::site_replication_reconcile::register_site_replication_retry_drainer(reconcile_site_replication_retry_drain); for (method, path, operation) in [ (Method::PUT, "/v3/site-replication/add", AdminOperation(&SiteReplicationAddHandler {})), @@ -1803,28 +1805,61 @@ async fn reconcile_site_replication_buckets() -> S3Result<()> { /// (`SiteReplicationEditHandler`), so a tick landing between them would rewrite the targets /// from the stale endpoint. The pending marker in the persisted state closes that window. /// Skipping costs nothing — the timer comes back. +async fn site_replication_reconcile_prerequisites_ready() -> bool { + if current_iam_handle().is_none() || current_object_store_handle().is_none() { + return false; + } + if let Err(err) = migrate_collapsed_retry_queue_paths().await { + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + result = "retry_queue_migration_failed", + error = ?err, + "admin site replication state" + ); + return false; + } + true +} + +fn reconcile_site_replication_retry_drain() -> std::pin::Pin + Send>> { + Box::pin(async { + let Some(lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else { + return; + }; + if !site_replication_reconcile_prerequisites_ready().await { + return; + } + match load_site_replication_state().await { + Ok(state) => { + if state.pending_endpoint_refresh.is_some() || state.pending_rotation.is_some() || state.pending_remove.is_some() + { + return; + } + } + Err(_) => return, + } + // Admission above observes a lifecycle-stable state. The lightweight + // drain itself handles only idempotent bucket setup, reloads state + // under the distributed repair lock, and shares that lock with bucket + // deletion. Do not hold this process-local guard across peer I/O: an + // outage recovery must not make admin add/edit/remove time out. + drop(lifecycle); + drain_site_replication_retry_queue_lightweight().await; + }) +} + fn reconcile_site_replication_wiring() -> std::pin::Pin + Send>> { Box::pin(async { // The scheduler starts before IAM and the object store are guaranteed ready (IAM // bootstrap may still be recovering), so an early tick returns quietly instead of // logging a failure for every reconciler. - if current_iam_handle().is_none() || current_object_store_handle().is_none() { - return; - } - - let Some(_lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else { + let Some(lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else { return; }; - if let Err(err) = migrate_collapsed_retry_queue_paths().await { - warn!( - event = EVENT_ADMIN_SITE_REPLICATION_STATE, - component = LOG_COMPONENT_ADMIN, - subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, - result = "retry_queue_migration_failed", - error = ?err, - "admin site replication state" - ); + if !site_replication_reconcile_prerequisites_ready().await { return; } @@ -1878,8 +1913,9 @@ fn reconcile_site_replication_wiring() -> std::pin::Pin bool { let (origin, generation) = fence; if origin != local_deployment_id && state.peers.contains_key(origin) { @@ -4791,105 +4827,135 @@ async fn backfill_existing_buckets_after_add( let resync_id = Uuid::new_v4().to_string(); for bucket in &buckets { - let name = &bucket.name; + let operation_name = bucket.name.clone(); + let lock_bucket = operation_name.clone(); + let operation_state = state.clone(); + let operation_local_peer = local_peer.clone(); + let operation_resync_id = resync_id.clone(); + let operation_bootstrap_token = bootstrap_token.map(str::to_owned); + let bucket_errors = with_site_replication_bucket_mutation_lock(store.clone(), &lock_bucket, move || async move { + let mut errors = SiteReplicationErrorSummary::default(); + let name = &operation_name; - if let Err(err) = ensure_site_replication_bucket_versioning(name).await { - warn!( - event = EVENT_ADMIN_SITE_REPLICATION_STATE, - component = LOG_COMPONENT_ADMIN, - subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, - bucket = %name, - result = "backfill_versioning_setup_failed", - error = ?err, - "admin site replication state" - ); - errors.push(format!("{name}: versioning setup failed: {err}")); - continue; - } - match ensure_site_replication_bucket_setup(name).await { - Ok(true) => {} - Ok(false) => { - // Runtime targets unavailable: the setup silently no-ops, which would make the - // downstream make-bucket broadcast and resync fail. Record it and skip so the - // operator sees this bucket was not propagated instead of an unqualified success. + if let Err(err) = ensure_site_replication_bucket_versioning(name).await { warn!( event = EVENT_ADMIN_SITE_REPLICATION_STATE, component = LOG_COMPONENT_ADMIN, subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, bucket = %name, - result = "backfill_bucket_setup_skipped", - "admin site replication state" - ); - errors.push(format!("{name}: replication setup skipped (site replication runtime unavailable)")); - continue; - } - Err(err) => { - warn!( - event = EVENT_ADMIN_SITE_REPLICATION_STATE, - component = LOG_COMPONENT_ADMIN, - subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, - bucket = %name, - result = "backfill_bucket_setup_failed", + result = "backfill_versioning_setup_failed", error = ?err, "admin site replication state" ); - errors.push(format!("{name}: bucket setup failed: {err}")); + errors.push(format!("{name}: versioning setup failed: {err}")); + return errors; } - } - // Broadcast the bucket to peers so they create it too (idempotent on the peer side). - // Read the real lock_enabled flag so peers recreate the bucket with the same object-lock - // setting — object lock cannot be added after bucket creation. - let lock_enabled = match metadata_sys::get(name).await { - Ok(bm) => bm.lock_enabled, - Err(err) => { - warn!( - event = EVENT_ADMIN_SITE_REPLICATION_STATE, - component = LOG_COMPONENT_ADMIN, - subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, - bucket = %name, - result = "backfill_bucket_metadata_read_failed", - fallback = "lock_enabled=false", - error = ?err, - "admin site replication state" - ); - false + match ensure_site_replication_bucket_setup(name).await { + Ok(true) => {} + Ok(false) => { + // Runtime targets unavailable: the setup silently no-ops, which would make the + // downstream make-bucket broadcast and resync fail. Record it and skip so the + // operator sees this bucket was not propagated instead of an unqualified success. + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + bucket = %name, + result = "backfill_bucket_setup_skipped", + "admin site replication state" + ); + errors.push(format!("{name}: replication setup skipped (site replication runtime unavailable)")); + return errors; + } + Err(err) => { + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + bucket = %name, + result = "backfill_bucket_setup_failed", + error = ?err, + "admin site replication state" + ); + errors.push(format!("{name}: bucket setup failed: {err}")); + } } - }; - if let Err(err) = broadcast_site_replication_make_bucket(name, lock_enabled, None, bootstrap_token).await { - warn!( - event = EVENT_ADMIN_SITE_REPLICATION_STATE, - component = LOG_COMPONENT_ADMIN, - subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, - bucket = %name, - result = "backfill_make_bucket_broadcast_failed", - error = ?err, - "admin site replication state" - ); - errors.push(format!("{name}: make-bucket broadcast failed: {err}")); - } - // Kick a resync toward every remote peer so existing objects travel across. - for peer in state.peers.values() { - if peer.deployment_id == local_peer.deployment_id || same_identity_endpoint(&peer.endpoint, &local_peer.endpoint) { - continue; - } - let manifest = site_bucket_resync_manifest_entry(name, peer, OffsetDateTime::now_utc()).await; - let result = if manifest.target_arn.is_empty() { - manifest - } else { - start_site_bucket_resync(name, &manifest.target_arn, &resync_id).await + // Broadcast the bucket to peers so they create it too (idempotent on the peer side). + // Read the real lock_enabled flag so peers recreate the bucket with the same object-lock + // setting — object lock cannot be added after bucket creation. + let lock_enabled = match metadata_sys::get(name).await { + Ok(bm) => bm.lock_enabled, + Err(err) => { + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + bucket = %name, + result = "backfill_bucket_metadata_read_failed", + fallback = "lock_enabled=false", + error = ?err, + "admin site replication state" + ); + false + } }; - if result.status == "failed" { + if let Err(err) = + broadcast_site_replication_make_bucket(name, lock_enabled, None, operation_bootstrap_token.as_deref()).await + { warn!( event = EVENT_ADMIN_SITE_REPLICATION_STATE, component = LOG_COMPONENT_ADMIN, subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, bucket = %name, - peer = %peer.endpoint, - result = "backfill_resync_kick_failed", - detail = %result.err_detail, + result = "backfill_make_bucket_broadcast_failed", + error = ?err, "admin site replication state" ); - errors.push(format!("{name} -> {}: resync kick failed: {}", peer.endpoint, result.err_detail)); + errors.push(format!("{name}: make-bucket broadcast failed: {err}")); + } + // Kick a resync toward every remote peer so existing objects travel across. + for peer in operation_state.peers.values() { + if peer.deployment_id == operation_local_peer.deployment_id + || same_identity_endpoint(&peer.endpoint, &operation_local_peer.endpoint) + { + continue; + } + let manifest = site_bucket_resync_manifest_entry(name, peer, OffsetDateTime::now_utc()).await; + let result = if manifest.target_arn.is_empty() { + manifest + } else { + start_site_bucket_resync(name, &manifest.target_arn, &operation_resync_id).await + }; + if result.status == "failed" { + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + bucket = %name, + peer = %peer.endpoint, + result = "backfill_resync_kick_failed", + detail = %result.err_detail, + "admin site replication state" + ); + errors.push(format!("{name} -> {}: resync kick failed: {}", peer.endpoint, result.err_detail)); + } + } + errors + }) + .await; + match bucket_errors { + Ok(bucket_errors) => errors.extend(bucket_errors), + Err(err) => { + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + bucket = %lock_bucket, + result = "backfill_bucket_mutation_lock_failed", + error = ?err, + "admin site replication state" + ); + errors.push(format!("{lock_bucket}: bucket mutation lock failed: {err}")); } } } @@ -6072,146 +6138,204 @@ fn parse_peer_join_response(body: &[u8], fallback_peer: PeerInfo) -> Result, present: &HashSet) -> S3Result<()> { + let mut missing = expected.difference(present).cloned().collect::>(); + if !missing.is_empty() { + missing.sort_unstable(); + return Err(S3Error::with_message( + S3ErrorCode::InvalidRequest, + format!( + "bucket `{}` disappeared while site replication was being added; peers may already be joined — re-run replicate add", + missing[0] + ), + )); + } + + let mut unexpected = present.difference(expected).cloned().collect::>(); + if !unexpected.is_empty() { + unexpected.sort_unstable(); + return Err(S3Error::with_message( + S3ErrorCode::InvalidRequest, + format!( + "bucket `{}` appeared while site replication was being added; peers may already be joined — re-run replicate add", + unexpected[0] + ), + )); + } + + Ok(()) +} + #[async_trait::async_trait] impl Operation for SiteReplicationAddHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationAddAction).await?; reject_site_replicator_on_public_admin(&cred)?; let replicate_ilm_expiry = sr_add_replicate_ilm_expiry(&req.uri); + let local_endpoint = site_replication_local_endpoint(&req.uri, &req.headers); let lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await?; - // Everything up to the commit below is preflight: peer probes, IAM - // work and the join fan-out all talk to the network, so none of it may - // run inside the state transaction. The snapshot read here is what the - // `updated_at` CAS in the commit validates. - let current_state = load_site_replication_state().await?; - if pending_endpoint_refresh(¤t_state).is_some() { - return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); - } - let local_peer = current_local_peer(&req, ¤t_state); let mut sites: Vec = read_site_replication_json(req, &cred.secret_key, true).await?; - // The web console's "Set Up Site Replication" omits the local deployment from the payload; - // inject it so the add preflight (which requires the local deployment) succeeds. No-op for `mc`. - ensure_local_site_present(&mut sites, &local_peer); - validate_add_sites(&sites, &local_peer)?; - let preflight_infos = add_preflight_infos(&sites, ¤t_state, &local_peer).await?; - validate_add_preflight_topology(&preflight_infos, &local_peer)?; - let expected_updated_at = current_state.updated_at; - require_add_peer_tls_capability(&sites, &local_peer).await?; - // Early exit on a state that moved under the preflight probes, BEFORE - // the IAM write and the join fan-out change anything remote. Advisory - // only — the binding check is the CAS inside the commit — but it fences - // the common race off the side-effect path and refreshes the merge - // base so the CAS window is only the join round trips. - let latest_state = load_site_replication_state().await?; - ensure_edit_precondition(&latest_state, expected_updated_at, None, "add preflight")?; - let current_state = latest_state; - let (service_account_access_key, service_account_secret_key) = - ensure_site_replicator_service_account(&cred.access_key, false).await?; - let bootstrap_buckets = preflight_infos - .iter() - .filter(|info| !same_identity_endpoint(&info.endpoint, &local_peer.endpoint)) - .flat_map(|info| info.buckets.keys().cloned()) - .collect(); - let add_in_progress_guard = SiteReplicationAddInProgressGuard::start(lifecycle_guard, bootstrap_buckets)?; - let mut state = merge_add_sites( - current_state, - local_peer.clone(), - sites.clone(), - service_account_access_key.clone(), - cred.access_key.clone(), - replicate_ilm_expiry, - ); - state.sync_state_initialized = true; - let join_req = SRPeerJoinEnvelope { - request: SRPeerJoinReq { - svc_acct_access_key: service_account_access_key, - svc_acct_secret_key: service_account_secret_key.clone(), - svc_acct_parent: String::new(), - peers: state.peers.clone(), - updated_at: state.updated_at, - }, - defer_sync_state_enable: true, - }; - let peer_join_path = - with_site_replication_bootstrap_token(SITE_REPLICATION_PEER_JOIN_PATH, &add_in_progress_guard.token.to_string()); + let admin_access_key = cred.access_key.clone(); + let admission_store = current_object_store_handle() + .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?; + let list_store = admission_store.clone(); + let (state, edit_generation, local_peer, service_account_secret_key, mut initial_sync_errors, _add_guard) = + with_site_replication_bucket_mutation_admission_lock(admission_store, move || async move { + // The writer starts before the local bucket snapshot and stays + // held through every peer join and the topology commit. A + // delete followed by a same-name create therefore cannot hide + // behind an unchanged final name set. Peer bootstrap callbacks + // use their internal path and do not acquire this public- + // mutation admission lock. + let current_state = load_site_replication_state().await?; + if pending_endpoint_refresh(¤t_state).is_some() { + return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); + } + let local_peer = local_peer_at_endpoint(local_endpoint, ¤t_state); + // The web console's "Set Up Site Replication" omits the local deployment from the payload; + // inject it so the add preflight (which requires the local deployment) succeeds. No-op for `mc`. + ensure_local_site_present(&mut sites, &local_peer); + validate_add_sites(&sites, &local_peer)?; + let preflight_infos = add_preflight_infos(&sites, ¤t_state, &local_peer).await?; + validate_add_preflight_topology(&preflight_infos, &local_peer)?; + let expected_updated_at = current_state.updated_at; + require_add_peer_tls_capability(&sites, &local_peer).await?; + // Early exit on a state that moved under the preflight probes, BEFORE + // the IAM write and the join fan-out change anything remote. Advisory + // only — the binding check is the CAS inside the commit — but it fences + // the common race off the side-effect path and refreshes the merge + // base so the CAS window is only the join round trips. + let latest_state = load_site_replication_state().await?; + ensure_edit_precondition(&latest_state, expected_updated_at, None, "add preflight")?; + let current_state = latest_state; + let (service_account_access_key, service_account_secret_key) = + ensure_site_replicator_service_account(&admin_access_key, false).await?; + let expected_buckets: HashSet = + preflight_infos.iter().flat_map(|info| info.buckets.keys().cloned()).collect(); + let bootstrap_buckets: HashSet = preflight_infos + .iter() + .filter(|info| !same_identity_endpoint(&info.endpoint, &local_peer.endpoint)) + .flat_map(|info| info.buckets.keys().cloned()) + .collect(); + let add_in_progress_guard = + SiteReplicationAddInProgressGuard::start(lifecycle_guard, bootstrap_buckets.clone())?; + let mut state = merge_add_sites( + current_state, + local_peer.clone(), + sites.clone(), + service_account_access_key.clone(), + admin_access_key, + replicate_ilm_expiry, + ); + state.sync_state_initialized = true; + let join_req = SRPeerJoinEnvelope { + request: SRPeerJoinReq { + svc_acct_access_key: service_account_access_key, + svc_acct_secret_key: service_account_secret_key.clone(), + svc_acct_parent: String::new(), + peers: state.peers.clone(), + updated_at: state.updated_at, + }, + defer_sync_state_enable: true, + }; + let peer_join_path = with_site_replication_bootstrap_token( + SITE_REPLICATION_PEER_JOIN_PATH, + &add_in_progress_guard.token.to_string(), + ); - let mut joined_endpoints = HashSet::new(); - let mut initial_sync_errors = SiteReplicationErrorSummary::default(); - for (site, preflight) in sites.iter().zip(preflight_infos.iter()) { - if same_identity_endpoint(&site.endpoint, &local_peer.endpoint) - || !joined_endpoints.insert(site_identity_key(&site.endpoint)) - { - continue; - } + let mut joined_endpoints = HashSet::new(); + let mut initial_sync_errors = SiteReplicationErrorSummary::default(); + for (site, preflight) in sites.iter().zip(preflight_infos.iter()) { + if same_identity_endpoint(&site.endpoint, &local_peer.endpoint) + || !joined_endpoints.insert(site_identity_key(&site.endpoint)) + { + continue; + } - let mut peer_join_req = join_req.clone(); - peer_join_req.request.svc_acct_parent = site.access_key.clone(); - let connection = PeerConnection::try_from(site)?; - let body = PeerAdminRequest::put(&connection, &peer_join_path, &site.access_key) - .send(&site.secret_key, &peer_join_req) + let mut peer_join_req = join_req.clone(); + peer_join_req.request.svc_acct_parent = site.access_key.clone(); + let connection = PeerConnection::try_from(site)?; + let body = PeerAdminRequest::put(&connection, &peer_join_path, &site.access_key) + .send(&site.secret_key, &peer_join_req) + .await?; + + let mut fallback_peer = existing_peer_for_endpoint(&state, &site.endpoint) + .unwrap_or_else(|| normalize_peer_site(site.clone(), replicate_ilm_expiry)); + fallback_peer.deployment_id = preflight.deployment_id.clone(); + let join_response = parse_peer_join_response(&body, fallback_peer).map_err(|e| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("parse peer join response from {} failed: {e}", site.endpoint), + ) + })?; + if !join_response.initial_sync_error_message.is_empty() { + initial_sync_errors.push(format!("{}: {}", site.endpoint, join_response.initial_sync_error_message)); + } + // An explicit no-op join. The peer answered 200 but wrote nothing — + // its persisted state is already newer than the snapshot it was + // sent — so the add is only PARTIALLY configured and saying + // "configured successfully" would be a lie (rustfs/rustfs#5963). + // `None` (a MinIO peer, or one older than the field) is not a + // no-op signal and is deliberately not reported. + if join_response.applied == Some(false) { + initial_sync_errors.push(format!( + "{}: peer did not apply the join (its site replication state is newer than the snapshot it was sent); \ + the site is not configured against this peer", + site.endpoint + )); + } + state = reconcile_peer_with_actual_identity(state, join_response.peer); + let reconciled_peer = existing_peer_for_endpoint(&state, &site.endpoint).ok_or_else(|| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("peer join response from {} did not identify the requested site", site.endpoint), + ) + })?; + validate_proposed_peer(&reconciled_peer).map_err(|err| { + S3Error::with_message( + S3ErrorCode::InvalidRequest, + format!("invalid peer join response from {}: {err}", site.endpoint), + ) + })?; + } + + mark_unknown_peer_sync_enabled(&mut state.peers); + + // Commit. The state transaction's CAS still fences topology + // writers that do not use bucket admission. By this point + // remote sites may already have accepted their joins, so a + // mismatch asks the operator to re-run add and reconverge. + let next_state = state; + let present = list_store + .list_bucket(&BucketOptions::default()) + .await + .map_err(ApiError::from)? + .into_iter() + .map(|bucket| bucket.name) + .collect::>(); + ensure_add_bucket_set_matches_preflight(&expected_buckets, &present)?; + let (state, edit_generation) = update_site_replication_state(move |state| { + if state.updated_at != expected_updated_at || pending_endpoint_refresh(state).is_some() { + return Err(s3_error!( + InvalidRequest, + "site replication state changed during peer join; the peers may already be joined — re-run replicate add" + )); + } + adopt_add_commit_state(state, next_state); + let edit_generation = next_peer_edit_generation(state); + Ok((state.clone(), edit_generation)) + }) .await?; - - let mut fallback_peer = existing_peer_for_endpoint(&state, &site.endpoint) - .unwrap_or_else(|| normalize_peer_site(site.clone(), replicate_ilm_expiry)); - fallback_peer.deployment_id = preflight.deployment_id.clone(); - let join_response = parse_peer_join_response(&body, fallback_peer).map_err(|e| { - S3Error::with_message( - S3ErrorCode::InternalError, - format!("parse peer join response from {} failed: {e}", site.endpoint), - ) - })?; - if !join_response.initial_sync_error_message.is_empty() { - initial_sync_errors.push(format!("{}: {}", site.endpoint, join_response.initial_sync_error_message)); - } - // An explicit no-op join. The peer answered 200 but wrote nothing — - // its persisted state is already newer than the snapshot it was - // sent — so the add is only PARTIALLY configured and saying - // "configured successfully" would be a lie (rustfs/rustfs#5963). - // `None` (a MinIO peer, or one older than the field) is not a - // no-op signal and is deliberately not reported. - if join_response.applied == Some(false) { - initial_sync_errors.push(format!( - "{}: peer did not apply the join (its site replication state is newer than the snapshot it was sent); \ - the site is not configured against this peer", - site.endpoint - )); - } - state = reconcile_peer_with_actual_identity(state, join_response.peer); - let reconciled_peer = existing_peer_for_endpoint(&state, &site.endpoint).ok_or_else(|| { - S3Error::with_message( - S3ErrorCode::InternalError, - format!("peer join response from {} did not identify the requested site", site.endpoint), - ) - })?; - validate_proposed_peer(&reconciled_peer).map_err(|err| { - S3Error::with_message( - S3ErrorCode::InvalidRequest, - format!("invalid peer join response from {}: {err}", site.endpoint), - ) - })?; - } - - mark_unknown_peer_sync_enabled(&mut state.peers); - - // Commit. The CAS runs inside the transaction, against the state the - // transaction itself loaded — the peer round trips above took however - // long they took, and only this check can tell whether the topology - // this add was planned against is still the current one. The error - // says so: by this point the remote sites already accepted their - // joins, and re-running the add is what reconverges the local side. - let next_state = state; - let (state, edit_generation) = update_site_replication_state(move |state| { - if state.updated_at != expected_updated_at || pending_endpoint_refresh(state).is_some() { - return Err(s3_error!( - InvalidRequest, - "site replication state changed during peer join; the peers may already be joined — re-run replicate add" - )); - } - adopt_add_commit_state(state, next_state); - let edit_generation = next_peer_edit_generation(state); - Ok((state.clone(), edit_generation)) - }) - .await?; + Ok(( + state, + edit_generation, + local_peer, + service_account_secret_key, + initial_sync_errors, + add_in_progress_guard, + )) + }) + .await?; // The finalize fan-out delivers peer-edit payloads, so it carries the // generation allocated in the commit above: the receiving site orders @@ -7185,8 +7309,14 @@ impl Operation for SRPeerEditHandler { // The fence is self-reported — the shared service account means // the sender cannot be identified — so it is honoured only after // the admissibility check, against the same state it will gate. - let commit_fence = - commit_fence.filter(|fence| peer_edit_fence_is_admissible(state, &local_peer.deployment_id, fence)); + let commit_fence = match commit_fence { + Some(fence) if peer_edit_fence_is_admissible(state, &local_peer.deployment_id, &fence) => Some(fence), + // A fenced edit can only come from a current remote peer. If + // that origin left while the retry was in flight, applying + // its body here would resurrect the removed topology. + Some(_) => return Ok(StateCommit::Unchanged(PeerEditOutcome::Acked)), + None => None, + }; // Ordering fence: the sending site allocates the generation under // its state-object lock, so a delivery that lost the race carries // a generation this site has already passed. Applying it would @@ -8886,6 +9016,41 @@ mod tests { ); } + #[test] + fn add_admission_starts_before_preflight_and_rejects_bucket_set_changes() { + let expected = HashSet::from(["remote-owned".to_string(), "shared".to_string()]); + let present = HashSet::from(["shared".to_string()]); + + let err = ensure_add_bucket_set_matches_preflight(&expected, &present) + .expect_err("a missing bootstrap bucket must reject the topology commit"); + assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); + + let present = HashSet::from([ + "remote-owned".to_string(), + "shared".to_string(), + "created-during-add".to_string(), + ]); + let err = ensure_add_bucket_set_matches_preflight(&expected, &present) + .expect_err("a bucket created during add must reject the topology commit"); + assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); + + let src = include_str!("site_replication.rs"); + let add = src + .split("impl Operation for SiteReplicationAddHandler") + .nth(1) + .and_then(|rest| rest.split("pub struct SiteReplicationRemoveHandler").next()) + .expect("add handler block"); + let admission = add + .find("with_site_replication_bucket_mutation_admission_lock") + .expect("distributed mutation admission"); + let preflight = add.find("add_preflight_infos").expect("bucket preflight"); + let validation = add + .find("ensure_add_bucket_set_matches_preflight") + .expect("bucket-set validation"); + let commit = add.find("adopt_add_commit_state").expect("topology commit"); + assert!(admission < preflight && preflight < validation && validation < commit); + } + #[test] fn test_tls_capability_gates_run_before_add_or_edit_state_side_effects() { let src = include_str!("site_replication.rs"); @@ -9182,13 +9347,19 @@ mod tests { ); // Fence hardening: origin and generation are self-reported by a // caller the shared service account cannot identify, so the handler - // must pass the fence through the admissibility check — against the - // same state the fence gates, i.e. inside the transaction — before - // reading or raising any high-water mark. + // must admit the fence against the same state it gates. An origin + // removed while a retry was in flight is acknowledged without + // applying the stale body; otherwise it could recreate topology. assert!( - handler_block.contains(".filter(|fence| peer_edit_fence_is_admissible(state, &local_peer.deployment_id, fence))"), + handler_block.contains( + "Some(fence) if peer_edit_fence_is_admissible(state, &local_peer.deployment_id, &fence) => Some(fence)" + ), "SRPeerEditHandler must admit a fence only through peer_edit_fence_is_admissible inside the state transaction" ); + assert!( + handler_block.contains("Some(_) => return Ok(StateCommit::Unchanged(PeerEditOutcome::Acked))"), + "SRPeerEditHandler must not apply a fenced edit after its origin leaves the current topology" + ); // P1-15 PR2: both halves of the fence and the edit they fence share // ONE transaction. Checking the fence against a state read outside the // lock would let the check pass on one snapshot and the write land on @@ -10352,8 +10523,9 @@ mod tests { /// A fence is self-reported: every site authenticates peer traffic with /// the same site-replicator credential, so a compromised peer can stamp /// ANY origin with ANY generation. An origin the receiver does not - /// replicate with — or the receiver itself — is ignored and plants no - /// mark; a mark a compromised peer plants for a CURRENT origin cannot + /// replicate with — or the receiver itself — is inadmissible and plants + /// no mark; the handler acknowledges such a request without applying its + /// body. A mark a compromised peer plants for a CURRENT origin cannot /// silence that origin, because the staleness window refuses to fence on /// a mark implausibly far above the genuine deliveries. #[test] @@ -12791,6 +12963,7 @@ mod tests { last_error: "site replication is not enabled".to_string(), updated_at: Some(OffsetDateTime::now_utc()), edit_generation: None, + peer_unreachable: false, deletions_recorded: false, }], ..Default::default() @@ -12989,6 +13162,7 @@ mod tests { last_error: "peer offline".to_string(), updated_at: Some(OffsetDateTime::now_utc()), edit_generation: None, + peer_unreachable: false, deletions_recorded: false, }], ..Default::default() diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index 7dada3492..e1951cd69 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, commit_site_replication_delete_bucket, prepare_site_replication_delete_bucket, + site_replication_bucket_meta_hook, site_replication_make_bucket_hook, with_site_replication_bucket_mutation_lock, }; use crate::storage::storage_api::lock_bucket_targets_metadata; use http::StatusCode; @@ -1331,23 +1332,34 @@ impl DefaultBucketUsecase { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; - let make_result = store - .make_bucket( - &bucket, - &MakeBucketOptions { - force_create: false, - lock_enabled, - ..Default::default() - }, - ) - .await; + // Keep the local namespace mutation and its peer hook ordered across + // every node in this site. Otherwise a delete waiting for repair + // coordination can arrive after this create on remote sites. + let operation_bucket = bucket.clone(); + let operation_store = store.clone(); + let make_result = with_site_replication_bucket_mutation_lock(store, &bucket, move || async move { + let make_result = operation_store + .make_bucket( + &operation_bucket, + &MakeBucketOptions { + force_create: false, + lock_enabled, + ..Default::default() + }, + ) + .await; + if make_result.is_ok() { + crate::storage::invalidate_bucket_validation_cache(&operation_bucket); + if let Err(err) = site_replication_make_bucket_hook(&operation_bucket, lock_enabled).await { + warn!(bucket = %operation_bucket, error = ?err, "site replication make bucket hook failed"); + } + } + make_result + }) + .await?; match make_result { - Ok(()) => { - // Invalidate the bucket validation cache so subsequent GETs - // see the newly created bucket immediately. - crate::storage::invalidate_bucket_validation_cache(&bucket); - } + Ok(()) => {} Err(StorageError::BucketExists(_)) => { // Per S3 spec: bucket namespace is global. Owner recreating returns 200 OK; // non-owner gets 409 BucketAlreadyExists. @@ -1358,10 +1370,6 @@ impl DefaultBucketUsecase { Err(e) => return Err(ApiError::from(e).into()), } - if let Err(err) = site_replication_make_bucket_hook(&bucket, lock_enabled).await { - warn!(bucket = %bucket, error = ?err, "site replication make bucket hook failed"); - } - let output = CreateBucketOutput::default(); counter!("rustfs_create_bucket_total").increment(1); let result = Ok(S3Response::new(output)); @@ -1397,16 +1405,41 @@ impl DefaultBucketUsecase { authorize_request(&mut req, Action::S3Action(S3Action::ForceDeleteBucketAction)).await?; } - store - .delete_bucket( - &input.bucket, - &DeleteBucketOptions { - force, - ..Default::default() - }, - ) - .await - .map_err(ApiError::from)?; + // Keep the local namespace mutation and its peer hook ordered across + // every node in this site so an older delete cannot overtake a new + // same-name make while it waits for repair coordination. + let operation_bucket = input.bucket.clone(); + let operation_store = store.clone(); + with_site_replication_bucket_mutation_lock(store, &input.bucket, move || async move { + let intent = prepare_site_replication_delete_bucket(&operation_bucket, force).await?; + let delete_result = operation_store + .delete_bucket( + &operation_bucket, + &DeleteBucketOptions { + force, + ..Default::default() + }, + ) + .await; + match delete_result { + Ok(()) => { + crate::storage::invalidate_bucket_validation_cache(&operation_bucket); + if let Some(intent) = intent + && let Err(err) = commit_site_replication_delete_bucket(&intent).await + { + warn!(bucket = %operation_bucket, error = ?err, "site replication delete bucket hook failed"); + } + Ok::<(), S3Error>(()) + } + Err(err) => { + if let Some(intent) = intent { + cancel_site_replication_delete_bucket(intent).await; + } + Err(S3Error::from(ApiError::from(err))) + } + } + }) + .await??; // 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 @@ -1415,16 +1448,9 @@ impl DefaultBucketUsecase { let cache_adapter = current_object_data_cache_for_context(self.context.as_deref()); let _ = invalidate_object_data_cache_bucket_after_delete(&cache_adapter, &input.bucket).await; - // Invalidate bucket validation cache - crate::storage::invalidate_bucket_validation_cache(&input.bucket); - // 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 { - warn!(bucket = %input.bucket, error = ?err, "site replication delete bucket hook failed"); - } - // Notify peers to drop their cached metadata for the now-deleted bucket. let request_context = req.extensions.get::().cloned(); notify_bucket_metadata_delete(input.bucket.clone(), request_context); diff --git a/rustfs/src/site_replication/hooks.rs b/rustfs/src/site_replication/hooks.rs index cd62704ee..95a830500 100644 --- a/rustfs/src/site_replication/hooks.rs +++ b/rustfs/src/site_replication/hooks.rs @@ -22,6 +22,57 @@ pub(crate) const SITE_REPLICATION_BUCKET_OP_CONFIGURE_REPLICATION: &str = "confi pub(crate) static SITE_REPLICATION_BUCKET_OP_LOCK: LazyLock> = LazyLock::new(|| RwLock::new(())); +const SITE_REPLICATION_BUCKET_MUTATION_LOCK_PREFIX: &str = "config/site-replication/bucket-mutation"; +pub(crate) const SITE_REPLICATION_BUCKET_MUTATION_ADMISSION_LOCK_PATH: &str = + "config/site-replication/bucket-mutation-admission.lock"; + +pub(crate) fn site_replication_bucket_mutation_lock_path(bucket: &str) -> String { + format!("{SITE_REPLICATION_BUCKET_MUTATION_LOCK_PREFIX}/{bucket}.lock") +} + +pub(crate) async fn with_site_replication_bucket_mutation_lock( + store: Arc, + bucket: &str, + operation: F, +) -> S3Result +where + F: FnOnce() -> Fut + Send + 'static, + Fut: std::future::Future + Send + 'static, + T: Send + 'static, +{ + let mutation_store = store.clone(); + let mutation_path = site_replication_bucket_mutation_lock_path(bucket); + with_config_object_read_lock( + store, + SITE_REPLICATION_BUCKET_MUTATION_ADMISSION_LOCK_PATH.to_string(), + move || async move { + with_config_object_write_lock(mutation_store, mutation_path, operation) + .await + .map_err(|err| S3Error::from(ApiError::from(err))) + }, + ) + .await + .map_err(|err| S3Error::from(ApiError::from(err)))? +} + +/// Exclude every local bucket namespace mutation from an add's local preflight +/// snapshot until its topology commit. Peer bootstrap callbacks do not enter +/// this public-mutation admission path, so they can finish while the writer is +/// held; post-commit fan-out and backfill must run after it is released. +pub(crate) async fn with_site_replication_bucket_mutation_admission_lock( + store: Arc, + operation: F, +) -> S3Result +where + F: FnOnce() -> Fut + Send + 'static, + Fut: std::future::Future> + Send + 'static, + T: Send + 'static, +{ + with_config_object_write_lock(store, SITE_REPLICATION_BUCKET_MUTATION_ADMISSION_LOCK_PATH.to_string(), operation) + .await + .map_err(|err| S3Error::from(ApiError::from(err)))? +} + #[derive(Debug, Default)] pub(crate) struct SiteReplicationBootstrapPlan { pub(crate) iam_items: Vec, @@ -329,6 +380,91 @@ pub(crate) fn site_replication_bootstrap_plan(info: &SRInfo) -> S3Result S3Result { + let mut plan = SiteReplicationBootstrapPlan { + bucket_make_ops: vec![bootstrap_bucket_make_op_path(bucket)], + bucket_configure_ops: vec![bootstrap_bucket_op_path( + &bucket.bucket, + SITE_REPLICATION_BUCKET_OP_CONFIGURE_REPLICATION, + )], + ..Default::default() + }; + append_bootstrap_bucket_items(&mut plan, bucket, replicate_ilm_expiry)?; + Ok(plan) +} + +pub(crate) fn site_replication_bucket_retry_plan_from_info( + bucket: &SRBucketInfo, + replicate_ilm_expiry: bool, +) -> S3Result { + let mut plan = site_replication_bucket_retry_plan_for(bucket, replicate_ilm_expiry)?; + // Omit only metadata the make/configure operations can reproduce exactly. + // Non-default versioning fields and operator-authored replication rules + // remain in the plan; their extra request cost intentionally defers the + // event to the complete drain when the lightweight budget is too small. + plan.bucket_items.retain(|item| !retry_bucket_metadata_is_redundant(item)); + Ok(plan) +} + +fn retry_bucket_metadata_is_redundant(item: &SRBucketMeta) -> bool { + match item.r#type.as_str() { + "version-config" => item.versioning.as_deref().is_some_and(|raw| { + deserialize::(&decode_bucket_meta_wire_value(raw)).is_ok_and(|config| { + config + == VersioningConfiguration { + status: Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)), + ..Default::default() + } + }) + }), + "replication-config" => item.replication_config.as_deref().is_some_and(|raw| { + deserialize::(&decode_bucket_meta_wire_value(raw)) + .is_ok_and(|config| config.role.trim().is_empty() && config.rules.iter().all(is_derived_site_replication_rule)) + }), + // `Some("")` is the in-memory sentinel used when the bucket is lock + // enabled but has no object-lock configuration body. The make query + // carries lockEnabled=true; sending an empty metadata body is neither + // useful nor parseable. + "object-lock-config" => item.object_lock_config.as_deref() == Some(""), + _ => false, + } +} + +pub(crate) async fn site_replication_bucket_retry_plan( + bucket: &str, + replicate_ilm_expiry: bool, +) -> S3Result { + let Some(store) = current_object_store_handle() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + let bucket_info = match store.get_bucket_info(bucket, &BucketOptions::default()).await { + Ok(bucket_info) => bucket_info, + Err(err) if is_err_bucket_not_found(&err) => return Ok(SiteReplicationBootstrapPlan::default()), + Err(err) => return Err(ApiError::from(err).into()), + }; + let lock_enabled = bucket_info.object_locking; + let metadata = metadata_sys::get(bucket).await.map_err(ApiError::from)?; + let mut bucket_info = SRBucketInfo { + bucket: bucket.to_string(), + created_at: bucket_info.created, + location: current_region().map(|region| region.to_string()).unwrap_or_default(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + ..Default::default() + }; + populate_sr_bucket_info_from_metadata(&mut bucket_info, &metadata).await; + if lock_enabled && bucket_info.object_lock_config.is_none() { + bucket_info.object_lock_config = Some(String::new()); + } + site_replication_bucket_retry_plan_from_info(&bucket_info, replicate_ilm_expiry) +} + pub async fn site_replication_make_bucket_hook(bucket: &str, lock_enabled: bool) -> S3Result<()> { let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.read().await; let runtime = { @@ -393,20 +529,273 @@ pub(crate) async fn broadcast_site_replication_make_bucket( broadcast_site_replication_json_using_runtime(runtime, &configure_path, &serde_json::json!({})).await } -pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: bool) -> S3Result<()> { +const SITE_REPLICATION_DELETE_INTENT_PENDING: &str = + "bucket deletion reserved; local completion and peer delivery are not yet known"; + +#[derive(Clone)] +struct SiteReplicationDeleteBucketReservation { + peer: PeerInfo, + previous: Option, + observed: SiteReplicationRetryEvent, +} + +pub(crate) struct SiteReplicationDeleteBucketIntent { + path: String, + reservations: Vec, + displaced: Vec, +} + +fn site_replication_delete_bucket_path(bucket: &str, force_delete: bool) -> String { let operation = if force_delete { "force-delete-bucket" } else { "delete-bucket" }; - let path = format!( + format!( "/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", form_urlencoded::Serializer::new(String::new()) .append_pair("bucket", bucket) .append_pair("operation", operation) .finish() - ); - broadcast_site_replication_json(&path, &serde_json::json!({})).await + ) +} + +/// Reserve every destructive peer delivery before the local namespace is +/// changed. The state transaction either persists the complete set or writes +/// nothing, so a full/unreadable queue fails the S3 delete closed. +pub(crate) async fn prepare_site_replication_delete_bucket( + bucket: &str, + force_delete: bool, +) -> S3Result> { + let path = site_replication_delete_bucket_path(bucket, force_delete); + let reservation_path = path.clone(); + update_site_replication_state_when_changed(move |state| { + if !state.enabled() { + return Ok(StateCommit::Unchanged(None)); + } + let local_peer = current_local_runtime_peer(state); + let peers = state + .peers + .values() + .filter(|peer| { + peer.deployment_id != local_peer.deployment_id && !same_identity_endpoint(&peer.endpoint, &local_peer.endpoint) + }) + .cloned() + .collect::>(); + if peers.is_empty() { + return Ok(StateCommit::Unchanged(None)); + } + + let mut reservations = Vec::with_capacity(peers.len()); + let mut displaced = Vec::new(); + for peer in peers { + let previous = state + .retry_queue + .iter() + .find(|event| retry_event_matches(event, &peer, &reservation_path)) + .cloned(); + displaced.extend(upsert_site_replication_retry_event( + &mut state.retry_queue, + &peer, + &reservation_path, + SITE_REPLICATION_DELETE_INTENT_PENDING, + None, + )?); + let observed = state + .retry_queue + .iter() + .find(|event| retry_event_matches(event, &peer, &reservation_path)) + .cloned() + .ok_or_else(|| { + S3Error::with_message( + S3ErrorCode::InternalError, + "site replication delete reservation disappeared before commit".to_string(), + ) + })?; + reservations.push(SiteReplicationDeleteBucketReservation { + peer, + previous, + observed, + }); + } + Ok(StateCommit::Changed(Some(SiteReplicationDeleteBucketIntent { + path: reservation_path, + reservations, + displaced, + }))) + }) + .await +} + +/// Roll back a reservation when the local storage delete definitively failed. +/// A concurrently revised reservation is preserved; it belongs to a newer +/// observation and this operation has no authority to settle it. +pub(crate) async fn cancel_site_replication_delete_bucket(intent: SiteReplicationDeleteBucketIntent) { + let path = intent.path.clone(); + let result = update_site_replication_state_when_changed(move |state| { + let mut changed = false; + for reservation in intent.reservations { + let Some(index) = state.retry_queue.iter().position(|event| { + retry_event_matches(event, &reservation.peer, &reservation.observed.path) + && event.id == reservation.observed.id + && event.updated_at == reservation.observed.updated_at + }) else { + continue; + }; + if let Some(previous) = reservation.previous { + state.retry_queue[index] = previous; + } else { + state.retry_queue.remove(index); + } + changed = true; + } + + let mut restored_all = true; + for displaced in intent.displaced { + let duplicate = state.retry_queue.iter().any(|event| { + event.id == displaced.id + || (event.peer_deployment_id == displaced.peer_deployment_id && event.path == displaced.path) + }); + if duplicate { + continue; + } + if state.retry_queue.len() >= SITE_REPLICATION_RETRY_QUEUE_LIMIT { + restored_all = false; + continue; + } + state.retry_queue.push(displaced); + changed = true; + } + Ok(if changed { + StateCommit::Changed(restored_all) + } else { + StateCommit::Unchanged(restored_all) + }) + }) + .await; + + match result { + Ok(true) => {} + Ok(false) => warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + path, + result = "delete_intent_cancel_incomplete", + "admin site replication state" + ), + Err(err) => warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + path, + result = "delete_intent_cancel_failed", + error = ?err, + "admin site replication state" + ), + } +} + +async fn broadcast_site_replication_delete_bucket(intent: &SiteReplicationDeleteBucketIntent) -> S3Result<()> { + let sends = intent.reservations.iter().cloned().map(|reservation| { + let request_path = intent.path.clone(); + async move { + let fallback_peer = reservation.peer.clone(); + let observed = reservation.observed.clone(); + let delivery_path = request_path.clone(); + let delivery = with_site_replication_state_read_lock(move |state| async move { + let Some(current_peer) = state.peers.get(&fallback_peer.deployment_id).cloned() else { + return Ok(None); + }; + let service_account_secret_key = + match site_replicator_service_account_secret(&state.service_account_access_key).await { + Ok(secret) => secret, + Err(err) => { + let Some(secret) = legacy_site_replicator_state_secret(&state) else { + return Err(err); + }; + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + result = "legacy_state_service_account_secret_fallback", + error = ?err, + "admin site replication state" + ); + secret + } + }; + let result = async { + let transport = PeerTransport::for_runtime_peer(¤t_peer).await?; + PeerAdminRequest::put(&transport.connection, &delivery_path, &state.service_account_access_key) + .with_client(&transport.client) + .send(&service_account_secret_key, &serde_json::json!({})) + .await + } + .await; + Ok(Some((current_peer, result))) + }) + .await; + match delivery { + Ok(Some((current_peer, Ok(_)))) => { + dequeue_observed_site_replication_retry_event(¤t_peer, &observed).await; + None + } + Ok(Some((current_peer, Err(err)))) => { + // Keep the failed deletion operator-visible, but never + // replay it automatically: without a bucket-incarnation + // fence, a delayed delete could erase a recreated bucket. + enqueue_site_replication_retry_event(¤t_peer, &request_path, &err).await; + Some(err) + } + Ok(None) => { + dequeue_observed_site_replication_retry_event(&reservation.peer, &observed).await; + None + } + Err(err) => { + enqueue_site_replication_retry_event(&reservation.peer, &request_path, &err).await; + Some(err) + } + } + } + }); + futures::future::join_all(sends) + .await + .into_iter() + .flatten() + .next() + .map_or(Ok(()), Err) +} + +pub(crate) async fn commit_site_replication_delete_bucket(intent: &SiteReplicationDeleteBucketIntent) -> S3Result<()> { + let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.read().await; + let store = + current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?; + let retry_peers = intent + .reservations + .iter() + .map(|reservation| reservation.peer.clone()) + .collect::>(); + let retry_path = intent.path.clone(); + let delivery_intent = SiteReplicationDeleteBucketIntent { + path: intent.path.clone(), + reservations: intent.reservations.clone(), + displaced: Vec::new(), + }; + match with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move { + broadcast_site_replication_delete_bucket(&delivery_intent).await + }) + .await + { + 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<()> { @@ -515,6 +904,39 @@ pub(crate) fn maybe_time(value: OffsetDateTime) -> Option { (value != OffsetDateTime::UNIX_EPOCH).then_some(value) } +async fn populate_sr_bucket_info_from_metadata(entry: &mut SRBucketInfo, metadata: &BucketMetadata) { + entry.policy = raw_config_to_string(&metadata.policy_config_json).and_then(|raw| serde_json::from_str(&raw).ok()); + entry.versioning = raw_config_to_base64(&metadata.versioning_config_xml); + entry.tags = raw_config_to_base64(&metadata.tagging_config_xml); + entry.object_lock_config = raw_config_to_base64(&metadata.object_lock_config_xml); + entry.sse_config = raw_config_to_base64(&metadata.encryption_config_xml); + entry.replication_config = raw_config_to_base64(&metadata.replication_config_xml); + entry.quota_config = raw_config_to_base64(&metadata.quota_config_json); + // Expiry subset only: this entry feeds both the bootstrap/repair plan + // (peers must not receive transition rules) and cross-site consistency + // views (transition rules are site-local and would read as false + // mismatches). A deleted expiry state is a `None` value with the + // deletion's axis so repair can converge peers that missed the live + // delete. + let expiry_statement = lifecycle_expiry_statement(metadata); + entry.expiry_lc_config = expiry_statement.as_ref().and_then(|(subset, _)| subset.clone()); + entry.cors_config = raw_config_to_base64(&metadata.cors_config_xml); + entry.policy_updated_at = maybe_time(metadata.policy_config_updated_at); + entry.tag_config_updated_at = maybe_time(metadata.tagging_config_updated_at); + entry.object_lock_config_updated_at = maybe_time(metadata.object_lock_config_updated_at); + entry.sse_config_updated_at = maybe_time(metadata.encryption_config_updated_at); + entry.versioning_config_updated_at = maybe_time(metadata.versioning_config_updated_at); + entry.replication_config_updated_at = maybe_time(metadata.replication_config_updated_at); + entry.quota_config_updated_at = maybe_time(metadata.quota_config_updated_at); + // The expiry axis, not the whole-config write time: local transition-only + // edits inflate the latter, and a repair item stamped with it could + // out-rank a newer real expiry edit on a third site. + entry.expiry_lc_config_updated_at = expiry_statement.map(|(_, axis)| axis); + entry.cors_config_updated_at = maybe_time(metadata.cors_config_updated_at); + entry.replication_targets_online = + Some(site_replication_targets_online(&entry.bucket, &metadata.replication_config_xml).await); +} + pub(crate) async fn build_sr_info(state: &SiteReplicationState, local_peer: &PeerInfo) -> S3Result { let Some(store) = current_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); @@ -546,37 +968,7 @@ pub(crate) async fn build_sr_info(state: &SiteReplicationState, local_peer: &Pee }; if let Some(metadata) = metadata { - entry.policy = raw_config_to_string(&metadata.policy_config_json).and_then(|raw| serde_json::from_str(&raw).ok()); - entry.versioning = raw_config_to_base64(&metadata.versioning_config_xml); - entry.tags = raw_config_to_base64(&metadata.tagging_config_xml); - entry.object_lock_config = raw_config_to_base64(&metadata.object_lock_config_xml); - entry.sse_config = raw_config_to_base64(&metadata.encryption_config_xml); - entry.replication_config = raw_config_to_base64(&metadata.replication_config_xml); - entry.quota_config = raw_config_to_base64(&metadata.quota_config_json); - // Expiry subset only: this entry feeds both the bootstrap/repair - // plan (peers must not receive transition rules) and cross-site - // consistency views (transition rules are site-local and would - // read as false mismatches). A deleted expiry state is a `None` - // value with the deletion's axis so repair can converge peers - // that missed the live delete. - let expiry_statement = lifecycle_expiry_statement(&metadata); - entry.expiry_lc_config = expiry_statement.as_ref().and_then(|(subset, _)| subset.clone()); - entry.cors_config = raw_config_to_base64(&metadata.cors_config_xml); - entry.policy_updated_at = maybe_time(metadata.policy_config_updated_at); - entry.tag_config_updated_at = maybe_time(metadata.tagging_config_updated_at); - entry.object_lock_config_updated_at = maybe_time(metadata.object_lock_config_updated_at); - entry.sse_config_updated_at = maybe_time(metadata.encryption_config_updated_at); - entry.versioning_config_updated_at = maybe_time(metadata.versioning_config_updated_at); - entry.replication_config_updated_at = maybe_time(metadata.replication_config_updated_at); - entry.quota_config_updated_at = maybe_time(metadata.quota_config_updated_at); - // The expiry axis, not the whole-config write time: local - // transition-only edits inflate the latter, and a repair item - // stamped with it could out-rank a newer real expiry edit on a - // third site. - entry.expiry_lc_config_updated_at = expiry_statement.map(|(_, axis)| axis); - entry.cors_config_updated_at = maybe_time(metadata.cors_config_updated_at); - entry.replication_targets_online = - Some(site_replication_targets_online(&bucket.name, &metadata.replication_config_xml).await); + populate_sr_bucket_info_from_metadata(&mut entry, &metadata).await; } info.buckets.insert(bucket.name, entry); diff --git a/rustfs/src/site_replication/mod.rs b/rustfs/src/site_replication/mod.rs index 4c3ad07d9..630c50576 100644 --- a/rustfs/src/site_replication/mod.rs +++ b/rustfs/src/site_replication/mod.rs @@ -47,6 +47,7 @@ use self::identity::{ canonical_endpoint, deployment_id_for_endpoint, mark_unknown_peer_sync_enabled, normalize_peer_map_by_identity_with, same_identity_endpoint, }; +pub(crate) use self::state_lock::with_site_replication_state_read_lock; use self::state_lock::{SITE_REPLICATION_STATE_PATH, with_site_replication_state_lock}; use crate::auth::constant_time_eq; use crate::config::get_config_snapshot; @@ -64,12 +65,12 @@ use crate::storage_api::site_replication::s3::{ #[cfg(test)] use crate::storage_api::site_replication::save_config as save_admin_config; use crate::storage_api::site_replication::{ - ARN, BUCKET_REPLICATION_CONFIG, BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, BucketOperations, BucketOptions, BucketTarget, - BucketTargetSys, BucketTargetType, BucketTargets, Credentials, ECStore, OperatorRuleContract, StorageError, - VersioningApi as _, assign_site_replication_rule_priorities, delete_config_no_lock, deserialize, is_site_replication_role, - lock_bucket_targets_metadata, metadata_sys, read_config as read_admin_config, read_config_no_lock, - replication_target_arn_deployment_id, save_config_no_lock, serialize, site_replication_rule_deployment_id, - with_config_object_read_lock, with_config_object_write_lock, + ARN, BUCKET_REPLICATION_CONFIG, BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, BucketMetadata, BucketOperations, + BucketOptions, BucketTarget, BucketTargetSys, BucketTargetType, BucketTargets, Credentials, ECStore, OperatorRuleContract, + StorageError, VersioningApi as _, assign_site_replication_rule_priorities, delete_config_no_lock, deserialize, + is_err_bucket_not_found, is_site_replication_role, lock_bucket_targets_metadata, metadata_sys, + read_config as read_admin_config, read_config_no_lock, replication_target_arn_deployment_id, save_config_no_lock, serialize, + site_replication_rule_deployment_id, with_config_object_read_lock, with_config_object_write_lock, }; use base64_simd::STANDARD as BASE64_STANDARD; use base64_simd::URL_SAFE_NO_PAD; diff --git a/rustfs/src/site_replication/repair.rs b/rustfs/src/site_replication/repair.rs index 5b3c260e3..b2785c0fd 100644 --- a/rustfs/src/site_replication/repair.rs +++ b/rustfs/src/site_replication/repair.rs @@ -649,7 +649,9 @@ pub(crate) async fn persist_site_replication_repair_task( let path = path.to_string(); update_site_replication_state(move |state| { match failure.as_deref() { - Some(error) => upsert_site_replication_retry_event(&mut state.retry_queue, &peer, &path, error, None), + Some(error) => { + upsert_site_replication_retry_event(&mut state.retry_queue, &peer, &path, error, None)?; + } None => { dequeue_site_replication_retry_events_including_escalated(&mut state.retry_queue, &peer, &path); // A repair is the operator's accountability transfer for the diff --git a/rustfs/src/site_replication/retry.rs b/rustfs/src/site_replication/retry.rs index 818fbee40..4b818c752 100644 --- a/rustfs/src/site_replication/retry.rs +++ b/rustfs/src/site_replication/retry.rs @@ -13,6 +13,7 @@ // limitations under the License. use super::*; +use futures::{StreamExt, stream}; pub(crate) const SITE_REPLICATION_RETRY_QUEUE_LIMIT: usize = 256; @@ -41,6 +42,12 @@ pub(crate) struct SiteReplicationRetryEvent { /// [`settle_site_replication_retry_events`]. #[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, 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")] + pub(crate) peer_unreachable: bool, /// Whether every failure folded into this collapsed IAM entry had its /// deletion body (if it was a deletion) recorded in /// [`SiteReplicationState::iam_deletion_replays`]. Only then may a @@ -196,27 +203,69 @@ pub(crate) fn settle_site_replication_retry_events( before.saturating_sub(queue.len()) } +pub(crate) fn settle_observed_site_replication_retry_event( + queue: &mut Vec, + peer: &PeerInfo, + observed: &SiteReplicationRetryEvent, +) -> usize { + let before = queue.len(); + queue.retain(|current| { + !(retry_event_matches(current, peer, &observed.path) + && current.id == observed.id + && current.updated_at == observed.updated_at) + }); + before.saturating_sub(queue.len()) +} + pub(crate) fn upsert_site_replication_retry_event( queue: &mut Vec, peer: &PeerInfo, path: &str, error: &str, generation: Option, -) { +) -> S3Result> { let path = collapsed_retry_queue_path(path).unwrap_or(path); let now = OffsetDateTime::now_utc(); let detail = summarize_peer_error_detail(error); + let peer_unreachable = retry_error_indicates_peer_unreachable(error); if let Some(event) = queue.iter_mut().find(|event| retry_event_matches(event, peer, path)) { + // The id is the event revision used by probe promotion and replay + // settlement. Refresh it on every failure so an older in-flight + // success can never acknowledge the newer observation. + event.id = Uuid::new_v4().to_string(); event.retry_count = event.retry_count.saturating_add(1); event.failed = event.retry_count >= SITE_REPLICATION_RETRY_FAILED_AFTER; event.last_error = detail; event.updated_at = Some(now); + event.peer_unreachable = peer_unreachable; // Keep the newest generation: an older delivery that fails afterwards // must not lower the fence and let its own success settle the event. event.edit_generation = event.edit_generation.max(generation); - return; + return Ok(Vec::new()); } + let slots_needed = queue + .len() + .saturating_add(1) + .saturating_sub(SITE_REPLICATION_RETRY_QUEUE_LIMIT); + let mut evict_indices = queue + .iter() + .enumerate() + .filter_map(|(index, event)| retry_event_is_safely_replayable(event).then_some(index)) + .take(slots_needed) + .collect::>(); + if evict_indices.len() != slots_needed { + return Err(S3Error::with_message( + S3ErrorCode::ServiceUnavailable, + "site replication retry queue is full of non-evictable liabilities; repair them before recording more failures" + .to_string(), + )); + } + let mut evicted = Vec::with_capacity(evict_indices.len()); + while let Some(index) = evict_indices.pop() { + evicted.push(queue.remove(index)); + } + evicted.reverse(); queue.push(SiteReplicationRetryEvent { id: Uuid::new_v4().to_string(), peer_deployment_id: peer.deployment_id.clone(), @@ -227,12 +276,38 @@ pub(crate) fn upsert_site_replication_retry_event( last_error: detail, updated_at: Some(now), edit_generation: generation, + peer_unreachable, deletions_recorded: false, }); - if queue.len() > SITE_REPLICATION_RETRY_QUEUE_LIMIT { - let overflow = queue.len() - SITE_REPLICATION_RETRY_QUEUE_LIMIT; - queue.drain(0..overflow); + Ok(evicted) +} + +pub(crate) fn is_destructive_bucket_retry_path(path: &str) -> bool { + matches!( + retry_bucket_operation(path).as_deref(), + Some("delete-bucket" | "force-delete-bucket" | "purge-deleted-bucket") + ) +} + +fn retry_event_is_safely_replayable(event: &SiteReplicationRetryEvent) -> bool { + if is_destructive_bucket_retry_path(&event.path) { + return false; } + matches!( + classify_site_replication_retry_event(event), + Some(RetryDrainAction::PeerEdit | RetryDrainAction::BucketOpReplay { .. }) + ) +} + +pub(crate) fn retry_error_indicates_peer_unreachable(error: &str) -> bool { + let error = error.to_ascii_lowercase(); + let Some((_, request)) = error.split_once("peer request to ") else { + return false; + }; + let Some((_, failure)) = request.split_once(" failed ") else { + return false; + }; + failure.starts_with("(connect):") || failure.starts_with("(dns resolution):") || failure.starts_with("(tls handshake):") } pub(crate) fn retry_stats_for_state(state: &SiteReplicationState) -> Option { @@ -271,7 +346,7 @@ pub(crate) async fn enqueue_site_replication_retry_event_for_generation( // (remove_sites already pruned them); recording a late failure for it // would only pollute retry_stats until the queue cap evicts it. if state.peers.contains_key(&peer_owned.deployment_id) { - upsert_site_replication_retry_event(&mut state.retry_queue, &peer_owned, &path_owned, &error_text, generation); + upsert_site_replication_retry_event(&mut state.retry_queue, &peer_owned, &path_owned, &error_text, generation)?; } Ok(()) }) @@ -356,12 +431,17 @@ pub(crate) fn iam_item_deletion_entity(item: &SRIAMItem) -> Option { /// live in the same state so the caller commits them in one transaction — a /// retry entry can never exist whose deletion body was lost to a separate /// failed write. -pub(crate) fn record_failed_iam_delivery(state: &mut SiteReplicationState, peer: &PeerInfo, item: &SRIAMItem, error: &str) { +pub(crate) fn record_failed_iam_delivery( + state: &mut SiteReplicationState, + peer: &PeerInfo, + item: &SRIAMItem, + error: &str, +) -> S3Result<()> { let existed = state .retry_queue .iter() .any(|event| retry_event_matches(event, peer, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH)); - upsert_site_replication_retry_event(&mut state.retry_queue, peer, SITE_REPLICATION_PEER_IAM_ITEM_WIRE_PATH, error, None); + upsert_site_replication_retry_event(&mut state.retry_queue, peer, SITE_REPLICATION_PEER_IAM_ITEM_WIRE_PATH, error, None)?; if !existed && let Some(event) = state .retry_queue @@ -375,13 +455,13 @@ pub(crate) fn record_failed_iam_delivery(state: &mut SiteReplicationState, peer: } let Some(entity) = iam_item_deletion_entity(item) else { - return; + return Ok(()); }; let item_value = match serde_json::to_value(item) { Ok(value) => value, Err(_) => { degrade_iam_retry_event_to_escalation(state, peer); - return; + return Ok(()); } }; let now = OffsetDateTime::now_utc(); @@ -390,9 +470,13 @@ pub(crate) fn record_failed_iam_delivery(state: &mut SiteReplicationState, peer: .iter_mut() .find(|record| iam_deletion_replay_matches(record, peer) && record.entity == entity) { + // This id is the replay-record revision. A settlement that sent the + // previous body must not remove a same-entity deletion that failed + // while its snapshot was in flight. + existing.id = Uuid::new_v4().to_string(); existing.item = item_value; existing.recorded_at = Some(now); - return; + return Ok(()); } let per_peer = state @@ -424,6 +508,7 @@ pub(crate) fn record_failed_iam_delivery(state: &mut SiteReplicationState, peer: item: item_value, recorded_at: Some(now), }); + Ok(()) } pub(crate) fn degrade_iam_retry_event_to_escalation(state: &mut SiteReplicationState, peer: &PeerInfo) { @@ -467,7 +552,7 @@ pub(crate) async fn record_failed_site_replication_iam_delivery(peer: &PeerInfo, // A departed peer can never drain its entries again (remove_sites // already pruned them) — mirror enqueue_site_replication_retry_event. if state.peers.contains_key(&peer_owned.deployment_id) { - record_failed_iam_delivery(state, &peer_owned, &item_owned, &error_text); + record_failed_iam_delivery(state, &peer_owned, &item_owned, &error_text)?; } Ok(()) }) @@ -511,15 +596,14 @@ pub(crate) async fn record_failed_site_replication_iam_delivery(peer: &PeerInfo, pub(crate) fn settle_replayed_iam_retry_events( state: &mut SiteReplicationState, peer: &PeerInfo, - path: &str, - snapshot_updated_at: Option, + observed: &SiteReplicationRetryEvent, replayed_record_ids: &[String], ) -> bool { state .iam_deletion_replays .retain(|record| !(iam_deletion_replay_matches(record, peer) && replayed_record_ids.contains(&record.id))); - if collapsed_retry_queue_path(path) != Some(SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH) { + if collapsed_retry_queue_path(&observed.path) != Some(SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH) { return false; } let Some(index) = state @@ -530,11 +614,7 @@ pub(crate) fn settle_replayed_iam_retry_events( return false; }; let event = &state.retry_queue[index]; - let newer_failure_recorded = match (event.updated_at, snapshot_updated_at) { - (Some(current), Some(seen)) => current > seen, - (Some(_), None) => true, - (None, _) => false, - }; + let newer_failure_recorded = event.id != observed.id || event.updated_at != observed.updated_at; if newer_failure_recorded && event.last_error != SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER { // The newer failure's own deletion (if any) has its own record; the // next drain pass replays it. @@ -548,24 +628,22 @@ pub(crate) fn settle_replayed_iam_retry_events( state.retry_queue.remove(index); return true; } - escalate_site_replication_retry_events_up_to(&mut state.retry_queue, peer, path, snapshot_updated_at); + escalate_site_replication_retry_events_up_to(&mut state.retry_queue, peer, &observed.path, observed.updated_at); false } pub(crate) async fn settle_replayed_site_replication_iam_retry_event( peer: &PeerInfo, - path: &str, - snapshot_updated_at: Option, + observed: &SiteReplicationRetryEvent, replayed_record_ids: Vec, ) { let peer_owned = peer.clone(); - let path_owned = path.to_string(); + let observed_owned = observed.clone(); let result = update_site_replication_state(move |state| { Ok(settle_replayed_iam_retry_events( state, &peer_owned, - &path_owned, - snapshot_updated_at, + &observed_owned, &replayed_record_ids, )) }) @@ -591,7 +669,7 @@ pub(crate) async fn settle_replayed_site_replication_iam_retry_event( event = EVENT_ADMIN_SITE_REPLICATION_STATE, peer = %peer.endpoint, deployment_id = %peer.deployment_id, - path, + path = %observed.path, error = ?err, "failed to settle replayed site replication IAM retry event" ); @@ -652,8 +730,30 @@ pub(crate) const SITE_REPLICATION_RETRY_DRAIN_BASE_BACKOFF_SECS: i64 = 600; /// converges at the next tick instead of waiting out this ceiling. pub(crate) const SITE_REPLICATION_RETRY_DRAIN_MAX_BACKOFF_SECS: i64 = 86_400; -/// What the background drain may do for one retry event. Everything not -/// representable here is operator territory (manual repair). +/// Bound a background drain round to one complete bucket bootstrap chain: +/// make, at most nine metadata records, then replication configuration. +/// Larger snapshots and topology edits remain queued for operator repair. +pub(crate) const SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER: usize = 11; + +/// Bound sockets for every retry pass. The lightweight pass also admits at +/// most this many peer request chains per round, bounding its lock hold time. +pub(crate) const SITE_REPLICATION_RETRY_DRAIN_PEER_CONCURRENCY: usize = 4; + +/// Rotate the bounded lightweight window instead of always admitting the +/// lexicographically first peers. Advancing by one full window per scheduler +/// round gives every queued peer a turn within `ceil(peer_count / limit)` +/// rounds, even when earlier peers each have a large bucket backlog. +pub(crate) fn lightweight_retry_peer_rotation(peer_count: usize, round: i64) -> usize { + if peer_count == 0 { + return 0; + } + let window = SITE_REPLICATION_RETRY_DRAIN_PEER_CONCURRENCY.min(peer_count); + (round.rem_euclid(peer_count as i64) as usize * window) % peer_count +} + +/// A replay shape the retry machinery can derive from persisted state. The +/// lightweight 30-second scheduler admits only bounded bucket-op chains; +/// snapshot and topology-wide work remains operator territory. #[derive(Debug, Clone, PartialEq, Eq)] pub(crate) enum RetryDrainAction { /// Constant-path IAM item deliveries collapse into one queue entry per @@ -670,6 +770,10 @@ pub(crate) enum RetryDrainAction { PeerEdit, } +pub(crate) fn is_lightweight_retry_drain_action(action: &RetryDrainAction) -> bool { + matches!(action, RetryDrainAction::BucketOpReplay { .. }) +} + #[derive(Clone)] pub(crate) enum RetrySnapshot { Iam(Vec), @@ -724,27 +828,126 @@ 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 body = match task { + SiteReplicationRepairTask::Iam(item) => serde_json::to_value(item), + SiteReplicationRepairTask::BucketMetadata(item) => serde_json::to_value(item), + SiteReplicationRepairTask::BucketMake(_) | SiteReplicationRepairTask::Replication(_) => Ok(serde_json::json!({})), + } + .map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize retry task failed: {err}")))?; + send_retry_request_if_peer_current(peer, transport, task.path(), access_key, secret_key, body).await +} + +pub(crate) async fn send_retry_request_if_peer_current( + peer: &PeerInfo, + transport: &PeerTransport, + path: &str, + access_key: &str, + secret_key: &str, + body: Value, +) -> S3Result { + let peer = peer.clone(); + let transport = transport.clone(); + let path = path.to_string(); + let access_key = access_key.to_string(); + let secret_key = secret_key.to_string(); + with_site_replication_state_read_lock(move |state| async move { + let current = state + .peers + .get(&peer.deployment_id) + .is_some_and(|current| same_identity_endpoint(¤t.endpoint, &peer.endpoint)); + if !current { + return Ok(false); + } + PeerAdminRequest::put(&transport.connection, &path, &access_key) + .with_client(&transport.client) + .send(&secret_key, &body) + .await?; + Ok(true) + }) + .await +} + +async fn send_peer_edit_retry_if_peer_current( + peer: &PeerInfo, + transport: &PeerTransport, + path: &str, + access_key: &str, + secret_key: &str, + body: Value, +) -> S3Result { + let peer_owned = peer.clone(); + let current = with_site_replication_state_read_lock(move |state| async move { + Ok(state + .peers + .get(&peer_owned.deployment_id) + .is_some_and(|current| same_identity_endpoint(¤t.endpoint, &peer_owned.endpoint))) + }) + .await?; + if !current { + return Ok(false); + } + // A peer-edit handler takes its own site's state write lock. Releasing + // this site's read lock before the request prevents simultaneous A -> B + // and B -> A retries from waiting on each other's write lock. The edit + // generation carried by `path` fences a delivery overtaken by a newer + // topology commit after this check. + PeerAdminRequest::put(&transport.connection, path, access_key) + .with_client(&transport.client) + .send(secret_key, &body) + .await?; + Ok(true) +} + #[derive(Hash, PartialEq, Eq)] pub(crate) enum IamSnapshotKey { Policy(String), @@ -872,6 +1075,66 @@ pub(crate) fn retry_bucket_name(path: &str) -> Option { .find_map(|(key, value)| (key == "bucket" && !value.is_empty()).then(|| value.into_owned())) } +pub(crate) fn bucket_op_retry_replay_tasks<'a>( + plan: &'a SiteReplicationBootstrapPlan, + operation: &str, + bucket: &str, +) -> S3Result>> { + let matches_bucket = |path: &&String| retry_bucket_name(path).as_deref() == Some(bucket); + match operation { + SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING => { + let mut tasks = plan + .bucket_make_ops + .iter() + .filter(matches_bucket) + .map(|path| SiteReplicationRepairTask::BucketMake(path.as_str())) + .collect::>(); + if tasks.is_empty() { + return Ok(tasks); + } + let configure_tasks = plan + .bucket_configure_ops + .iter() + .filter(matches_bucket) + .map(|path| SiteReplicationRepairTask::Replication(path.as_str())) + .collect::>(); + if configure_tasks.is_empty() { + return Err(S3Error::with_message( + S3ErrorCode::InternalError, + 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) + } + SITE_REPLICATION_BUCKET_OP_CONFIGURE_REPLICATION => Ok(plan + .bucket_configure_ops + .iter() + .filter(matches_bucket) + .map(|path| SiteReplicationRepairTask::Replication(path.as_str())) + .collect()), + _ => Err(S3Error::with_message( + S3ErrorCode::InvalidArgument, + format!("unsupported site replication retry bucket operation {operation:?}"), + )), + } +} + +pub(crate) fn retry_drain_request_count(action: &RetryDrainAction, plan: Option<&SiteReplicationBootstrapPlan>) -> usize { + match action { + RetryDrainAction::BucketOpReplay { operation, bucket } => plan + .and_then(|plan| bucket_op_retry_replay_tasks(plan, operation, bucket).ok()) + .map_or(0, |tasks| tasks.len()), + RetryDrainAction::IamSnapshot | RetryDrainAction::BucketMetadataSnapshot | RetryDrainAction::PeerEdit => usize::MAX, + } +} + /// A collapsed retry event after a stable snapshot resend is escalated with /// this marker instead of being cleared: the snapshot contains no task for a /// failed deletion, so remote absence remains operator-visible. Collapsed @@ -989,10 +1252,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). The drain probes each such peer with one cheap -/// request per tick and promotes its backlog when the probe succeeds. The -/// base backoff still floors individual re-attempts so a reachable peer that -/// keeps failing a delivery 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, @@ -1005,7 +1268,14 @@ pub(crate) fn deferred_site_replication_retry_events( .filter(|event| !site_replication_retry_backoff_elapsed(event, now)) .filter(|event| { event.updated_at.is_none_or(|updated_at| { - now.unix_timestamp().saturating_sub(updated_at.unix_timestamp()) >= SITE_REPLICATION_RETRY_DRAIN_BASE_BACKOFF_SECS + event.peer_unreachable + // Older binaries do not persist `peer_unreachable`. Parse + // only the locally-produced outer transport-error shape so + // rolling upgrades retain fast recovery without trusting + // an HTTP error body containing the same words. + || retry_error_indicates_peer_unreachable(&event.last_error) + || now.unix_timestamp().saturating_sub(updated_at.unix_timestamp()) + >= SITE_REPLICATION_RETRY_DRAIN_BASE_BACKOFF_SECS }) }) .cloned() @@ -1036,11 +1306,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 { @@ -1054,29 +1344,50 @@ 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 recovered = Vec::new(); + let probe_results = stream::iter(probes) + .buffer_unordered(SITE_REPLICATION_RETRY_DRAIN_PEER_CONCURRENCY) + .collect::>() + .await; + for (deployment_id, peer_endpoint, events, reachable) in probe_results { + if reachable { info!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, component = LOG_COMPONENT_ADMIN, subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, - event = EVENT_ADMIN_SITE_REPLICATION_STATE, - peer = %peer.endpoint, + result = "retry_backoff_probe_promoted", + peer = %peer_endpoint, 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 replay" ); - 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 @@ -1142,6 +1453,78 @@ pub(crate) async fn drain_site_replication_retry_queue() { } } +/// Fast outage-recovery pass used by the 30-second scheduler. It limits work +/// to one bounded bucket-op chain per peer, builds no site-wide snapshot, and +/// runs reachability probes concurrently with replay for other peers. +pub(crate) async fn drain_site_replication_retry_queue_lightweight() { + if let Err(err) = drain_site_replication_retry_queue_lightweight_inner().await { + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + result = "retry_drain_failed", + error = ?err, + "admin site replication state" + ); + } +} + +async fn drain_site_replication_retry_queue_lightweight_inner() -> S3Result<()> { + let Some(runtime) = runtime_site_replication_targets().await? else { + return Ok(()); + }; + log_site_replication_retry_liabilities(&runtime.state); + 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 mut deferred = deferred_site_replication_retry_events(&runtime.state, now); + actionable.retain(|event| { + classify_site_replication_retry_event(event).is_some_and(|action| is_lightweight_retry_drain_action(&action)) + }); + deferred.retain(|event| { + classify_site_replication_retry_event(event).is_some_and(|action| is_lightweight_retry_drain_action(&action)) + }); + if actionable.is_empty() && deferred.is_empty() { + return Ok(()); + } + let Some(store) = current_object_store_handle() else { + return Ok(()); + }; + + // Probes are read-only and can consume the full request timeout. Keep + // them outside repair coordination; promotion is fenced by event id and + // timestamp, and the locked reload below decides what may actually send. + promote_reachable_deferred_retry_events(&runtime, &actionable, deferred).await?; + + with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move { + 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); + actionable.retain(|event| { + classify_site_replication_retry_event(event).is_some_and(|action| is_lightweight_retry_drain_action(&action)) + }); + if actionable.is_empty() { + return Ok(()); + } + drain_site_replication_retry_queue_lightweight_locked(Arc::new(runtime), actionable, now).await + }) + .await + .map_err(ApiError::from)? +} + pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> { let Some(runtime) = runtime_site_replication_targets().await? else { return Ok(()); @@ -1150,7 +1533,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(()); @@ -1167,27 +1550,146 @@ 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 + + // Persist successful recovery probes without monopolizing repair + // coordination. The locked reload below observes those promotions and + // replays them in this same round. + promote_reachable_deferred_retry_events(&runtime, &actionable, deferred).await?; + + // Serialize against operator repair execution. Peer membership is + // re-checked from a distributed state snapshot 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. An operator repair holding the lock makes - // this tick skip after the lock-acquire timeout. + // re-runs the dry-run. The lock elects one server to 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 actionable = actionable_site_replication_retry_events(&runtime.state, now); + if actionable.is_empty() { + return Ok(()); + } drain_site_replication_retry_queue_locked(runtime, actionable).await }) .await .map_err(ApiError::from)? } +async fn drain_site_replication_retry_queue_lightweight_locked( + runtime: Arc, + events: Vec, + now: OffsetDateTime, +) -> S3Result<()> { + let mut events_by_peer: BTreeMap> = BTreeMap::new(); + for event in events { + events_by_peer + .entry(event.peer_deployment_id.clone()) + .or_default() + .push(event); + } + + let mut peer_groups = 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) + { + return None; + } + Some((peer, peer_events)) + }) + .collect::>(); + let interval_secs = crate::site_replication_reconcile::RETRY_DRAIN_INTERVAL.as_secs() as i64; + let round = now.unix_timestamp().div_euclid(interval_secs); + let rotation = lightweight_retry_peer_rotation(peer_groups.len(), round); + peer_groups.rotate_left(rotation); + + let peer_replays = peer_groups + .into_iter() + .map(|(peer, peer_events)| { + let runtime = Arc::clone(&runtime); + async move { + let Some((event, action, bucket)) = peer_events.into_iter().find_map(|event| { + let action = classify_site_replication_retry_event(&event)?; + let bucket = match &action { + RetryDrainAction::BucketOpReplay { bucket, .. } => bucket.clone(), + _ => return None, + }; + Some((event, action, bucket)) + }) else { + return (0, 0); + }; + let plan = match site_replication_bucket_retry_plan( + &bucket, + site_replication_state_replicates_ilm_expiry(&runtime.state), + ) + .await + { + Ok(plan) => plan, + Err(err) => { + enqueue_site_replication_retry_event(&peer, &event.path, &err).await; + return (0, 1); + } + }; + if retry_drain_request_count(&action, Some(&plan)) > SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER { + return (0, 0); + } + let transport = match PeerTransport::for_runtime_peer(&peer).await { + Ok(transport) => transport, + Err(err) => { + enqueue_site_replication_retry_event(&peer, &event.path, &err).await; + return (0, 1); + } + }; + match drain_one_site_replication_retry_event(&runtime, &peer, &transport, &event, action, Some(&plan)).await { + Ok(true) => (1, 0), + Ok(false) => (0, 0), + Err(_) => (0, 1), + } + } + }) + .take(SITE_REPLICATION_RETRY_DRAIN_PEER_CONCURRENCY); + + let mut settled = 0usize; + let mut failures = 0usize; + let replay_results = stream::iter(peer_replays) + .buffer_unordered(SITE_REPLICATION_RETRY_DRAIN_PEER_CONCURRENCY) + .collect::>() + .await; + for (peer_settled, peer_failures) in replay_results { + settled += peer_settled; + failures += peer_failures; + } + if settled > 0 || failures > 0 { + info!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + result = "retry_drain_settled", + settled, + failures, + "admin site replication state" + ); + } + Ok(()) +} + pub(crate) async fn drain_site_replication_retry_queue_locked( runtime: SiteReplicationRuntime, events: Vec, @@ -1212,39 +1714,54 @@ 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; + let replay_results = stream::iter(peer_replays) + .buffer_unordered(SITE_REPLICATION_RETRY_DRAIN_PEER_CONCURRENCY) + .collect::>() + .await; + for (peer_settled, peer_failures) in replay_results { + settled += peer_settled; + failures += peer_failures; } if settled > 0 || failures > 0 { @@ -1292,12 +1809,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()); } @@ -1306,22 +1832,20 @@ 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)?; let fresh_snapshot = RetrySnapshot::from_plan(&action, &fresh_plan).expect("snapshot action has a snapshot"); if fresh_snapshot.fingerprint()? == current_fingerprint { if is_iam { - settle_replayed_site_replication_iam_retry_event( - peer, - &event.path, - event.updated_at, - replayed_record_ids, - ) - .await; + settle_replayed_site_replication_iam_retry_event(peer, event, replayed_record_ids).await; } else { escalate_site_replication_retry_event_up_to(peer, &event.path, event.updated_at).await; } @@ -1339,37 +1863,29 @@ pub(crate) async fn drain_one_site_replication_retry_event( // Replay from the CURRENT plan, never the recorded path: the // recorded query can carry an expired one-shot bootstrap token or // a stale createdAt. - let make_op = operation == SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING; - let paths = if make_op { - &plan.bucket_make_ops - } else { - &plan.bucket_configure_ops - }; - let tasks: Vec> = paths - .iter() - .filter(|path| retry_bucket_name(path).as_deref() == Some(bucket.as_str())) - .map(|path| { - if make_op { - SiteReplicationRepairTask::BucketMake(path) - } else { - SiteReplicationRepairTask::Replication(path) - } - }) - .collect(); - if tasks.is_empty() { - // The bucket left the plan (deleted, or replication no longer - // configured): the recorded intent is stale, settle it. - dequeue_site_replication_retry_event(peer, &event.path).await; - return Ok(true); - } - for task in &tasks { - if let Err(err) = task.send(transport, access_key, secret_key).await { + let tasks = match bucket_op_retry_replay_tasks(plan, &operation, &bucket) { + Ok(tasks) => tasks, + Err(err) => { enqueue_site_replication_retry_event(peer, &event.path, &err).await; return Err(err); } + }; + if tasks.is_empty() { + // The bucket left the plan (deleted, or replication no longer + // configured): the recorded intent is stale, settle it. + return Ok(dequeue_observed_site_replication_retry_event(peer, event).await); } - dequeue_site_replication_retry_event(peer, &event.path).await; - Ok(true) + for task in &tasks { + 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); + } + } + } + Ok(dequeue_observed_site_replication_retry_event(peer, event).await) } RetryDrainAction::PeerEdit => { // The recorded generation is stale by definition — the receiver @@ -1394,19 +1910,22 @@ 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 { - if let Err(err) = PeerAdminRequest::put(&transport.connection, &edit_path, access_key) - .with_client(&transport.client) - .send(secret_key, body) - .await - { - enqueue_site_replication_retry_event_for_generation( - peer, - SITE_REPLICATION_PEER_EDIT_PATH, - &err, - delivery_fence, - ) - .await; - return Err(err); + let body = serde_json::to_value(body).map_err(|err| { + S3Error::with_message(S3ErrorCode::InternalError, format!("serialize retry peer edit failed: {err}")) + })?; + match send_peer_edit_retry_if_peer_current(peer, transport, &edit_path, access_key, secret_key, body).await { + Ok(true) => {} + Ok(false) => return Ok(false), + Err(err) => { + enqueue_site_replication_retry_event_for_generation( + peer, + SITE_REPLICATION_PEER_EDIT_PATH, + &err, + delivery_fence, + ) + .await; + return Err(err); + } } } dequeue_site_replication_retry_event_for_generation(peer, SITE_REPLICATION_PEER_EDIT_PATH, delivery_fence).await; @@ -1422,6 +1941,40 @@ pub(crate) async fn dequeue_site_replication_retry_event(peer: &PeerInfo, path: dequeue_site_replication_retry_event_for_generation(peer, path, None).await } +pub(crate) async fn dequeue_observed_site_replication_retry_event(peer: &PeerInfo, observed: &SiteReplicationRetryEvent) -> bool { + let result = async { + let mut probe = load_site_replication_state().await?; + if settle_observed_site_replication_retry_event(&mut probe.retry_queue, peer, observed) == 0 { + return Ok(false); + } + let peer_owned = peer.clone(); + let observed_owned = observed.clone(); + update_site_replication_state(move |state| { + Ok(settle_observed_site_replication_retry_event(&mut state.retry_queue, &peer_owned, &observed_owned) > 0) + }) + .await + } + .await; + + match result { + Ok(settled) => settled, + Err(err) => { + warn!( + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + result = "retry_event_dequeue_failed", + peer = %peer.endpoint, + deployment_id = %peer.deployment_id, + path = %observed.path, + error = ?err, + "failed to dequeue observed site replication retry event" + ); + false + } + } +} + pub(crate) async fn dequeue_site_replication_retry_event_for_generation(peer: &PeerInfo, path: &str, generation: Option) { let result = async { // Fast path: this sits on every successful hook broadcast, so probe diff --git a/rustfs/src/site_replication/state_lock.rs b/rustfs/src/site_replication/state_lock.rs index 6f9cc7bfd..9eb228882 100644 --- a/rustfs/src/site_replication/state_lock.rs +++ b/rustfs/src/site_replication/state_lock.rs @@ -28,14 +28,18 @@ //! process-local lock must never be reintroduced in front of it as if it //! added protection. All IO inside the closure must use the `*_no_lock` //! config helpers — the locked variants would self-deadlock on the same -//! object lock. Do not perform peer network calls or take other config locks -//! inside the closure. +//! object lock. Write-lock closures must not perform peer network calls or +//! take other config locks. A read-lock closure may carry bounded peer +//! delivery only when the receiver cannot write this state. Peer-edit +//! delivery must run after the read lock is released. //! -//! Lock order: lifecycle -> bucket operation -> repair admission -//! -> state object lock -> per-bucket metadata. +//! Lock order: lifecycle -> bucket-mutation admission -> per-bucket mutation +//! -> bucket operation -> repair admission -> state object lock -> +//! per-bucket metadata. A path may skip levels, but must not acquire an +//! earlier level while holding a later one. -use super::{S3Error, S3ErrorCode, S3Result}; -use crate::storage_api::site_replication::{ECStore, with_config_object_write_lock}; +use super::{S3Error, S3ErrorCode, S3Result, SiteReplicationState, load_site_replication_state_no_lock}; +use crate::storage_api::site_replication::{ECStore, with_config_object_read_lock, with_config_object_write_lock}; use std::sync::Arc; use crate::runtime_sources::current_object_store_handle; @@ -57,6 +61,27 @@ where with_site_replication_state_lock_on(store, operation).await } +/// Hold the distributed state-object read lock while `operation` validates a +/// topology snapshot. The closure may carry a bounded peer delivery only when +/// its receiver cannot write site replication state; peer-edit delivery must +/// run after this lock is released. Topology writers use the matching write +/// lock through [`with_site_replication_state_lock`]. +pub(crate) async fn with_site_replication_state_read_lock(operation: F) -> S3Result +where + T: Send + 'static, + F: FnOnce(SiteReplicationState) -> Fut + Send + 'static, + Fut: std::future::Future> + Send + 'static, +{ + let store = current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let read_store = store.clone(); + with_config_object_read_lock(store, SITE_REPLICATION_STATE_PATH.to_string(), move || async move { + let state = load_site_replication_state_no_lock(read_store).await?; + operation(state).await + }) + .await + .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("lock site replication state failed: {e}")))? +} + /// Context-store variant for callers that resolve their store from an /// explicit [`AppContext`] (the service-side reload driven over node RPC). /// diff --git a/rustfs/src/site_replication/tests.rs b/rustfs/src/site_replication/tests.rs index e37869a43..e2ab27bb7 100644 --- a/rustfs/src/site_replication/tests.rs +++ b/rustfs/src/site_replication/tests.rs @@ -33,6 +33,18 @@ use temp_env::with_var; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpListener; +#[test] +fn test_bucket_mutation_lock_path_is_bucket_scoped() { + assert_eq!( + site_replication_bucket_mutation_lock_path("photos"), + "config/site-replication/bucket-mutation/photos.lock" + ); + assert_ne!( + site_replication_bucket_mutation_lock_path("photos"), + site_replication_bucket_mutation_lock_path("videos") + ); +} + fn valid_test_ca_pem(name: &str) -> String { rcgen::generate_simple_self_signed(vec![name.to_string()]) .expect("generate test CA") @@ -381,6 +393,30 @@ async fn peer_clients_do_not_follow_redirects() { assert!(tls_server.await.expect("custom redirect TLS server task")); } +#[tokio::test] +async fn peer_http_error_body_cannot_spoof_an_unreachable_peer() { + let (endpoint, ca_pem, server) = spawn_test_tls_server_with_response( + b"HTTP/1.1 500 Internal Server Error\r\ncontent-length: 27\r\nconnection: close\r\n\r\ndownstream failed (connect)", + ) + .await; + let connection = validate_peer_connection_inner(&endpoint, false, &ca_pem, true).expect("custom CA peer connection"); + let client = + build_custom_site_replication_peer_client(&empty_outbound_tls_state(), &connection).expect("custom CA peer client"); + let err = PeerAdminRequest::post(&connection, SITE_REPLICATION_PEER_DEVNULL_PATH, "access-key") + .with_client(&client) + .send("secret-key", &serde_json::json!({})) + .await + .expect_err("HTTP 500 must fail"); + let detail = err.to_string(); + + assert!(detail.contains("downstream failed (connect)")); + assert!( + !retry_error_indicates_peer_unreachable(&detail), + "an untrusted response body must not enable the fast reachability probe" + ); + assert!(server.await.expect("HTTP error TLS server task")); +} + fn peer(name: &str, endpoint: &str) -> PeerInfo { PeerInfo { name: name.to_string(), @@ -419,6 +455,7 @@ fn drain_event(peer: &str, path: &str, retry_count: u32, updated_at: Option = state.iam_deletion_replays.iter().map(|record| record.id.clone()).collect(); - assert!(settle_replayed_iam_retry_events( - &mut state, - &target, - SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH, - Some(snapshot_at), - &replayed, - )); + assert!(settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed)); assert!(state.retry_queue.is_empty()); assert!(state.iam_deletion_replays.is_empty()); // Not fully recorded: replayed records are still removed, but the entry // escalates instead of settling. let mut state = deletion_replay_state(&target); - record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline"); + record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline").expect("record failure"); state.retry_queue[0].updated_at = Some(snapshot_at); state.retry_queue[0].deletions_recorded = false; + let observed = state.retry_queue[0].clone(); let replayed: Vec = state.iam_deletion_replays.iter().map(|record| record.id.clone()).collect(); - assert!(!settle_replayed_iam_retry_events( - &mut state, - &target, - SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH, - Some(snapshot_at), - &replayed, - )); + assert!(!settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed)); assert!(state.iam_deletion_replays.is_empty()); assert_eq!(state.retry_queue.len(), 1); assert_eq!(state.retry_queue[0].last_error, SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER); @@ -654,17 +683,13 @@ fn test_settle_replayed_iam_retry_events_settles_or_escalates() { // Newer failure since the snapshot: entry untouched and drain-eligible, // residual (unreplayed) record kept for the next pass. let mut state = deletion_replay_state(&target); - record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline"); + record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline").expect("record failure"); + state.retry_queue[0].updated_at = Some(snapshot_at); + let observed = state.retry_queue[0].clone(); let replayed: Vec = state.iam_deletion_replays.iter().map(|record| record.id.clone()).collect(); - record_failed_iam_delivery(&mut state, &target, &user_delete_item("bob"), "peer offline"); + record_failed_iam_delivery(&mut state, &target, &user_delete_item("bob"), "peer offline").expect("record failure"); state.retry_queue[0].updated_at = Some(snapshot_at + time::Duration::seconds(5)); - assert!(!settle_replayed_iam_retry_events( - &mut state, - &target, - SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH, - Some(snapshot_at), - &replayed, - )); + assert!(!settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed)); assert_eq!(state.retry_queue.len(), 1); assert_ne!(state.retry_queue[0].last_error, SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER); assert!( @@ -673,6 +698,24 @@ fn test_settle_replayed_iam_retry_events_settles_or_escalates() { ); assert_eq!(state.iam_deletion_replays.len(), 1); assert_eq!(state.iam_deletion_replays[0].entity, "iam-user:bob"); + + // A newer deletion of the same entity gets a fresh replay-record id. An + // older settlement therefore removes neither its body nor its queue + // revision, even if the persisted timestamps happen to be equal. + let mut state = deletion_replay_state(&target); + record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "first failure").expect("record failure"); + state.retry_queue[0].updated_at = Some(snapshot_at); + let observed = state.retry_queue[0].clone(); + let replayed = vec![state.iam_deletion_replays[0].id.clone()]; + record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "newer failure").expect("record newer failure"); + state.retry_queue[0].updated_at = Some(snapshot_at); + assert_ne!(state.iam_deletion_replays[0].id, replayed[0]); + + assert!(!settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed)); + assert_eq!(state.retry_queue.len(), 1); + assert_ne!(state.retry_queue[0].id, observed.id); + assert_eq!(state.iam_deletion_replays.len(), 1); + assert_eq!(state.iam_deletion_replays[0].entity, "iam-user:alice"); } /// Merging legacy wire-path rows into the collapsed entry must not launder an @@ -748,6 +791,281 @@ fn test_classify_site_replication_retry_event_actions() { assert_eq!(classify("/rustfs/admin/v3/site-replication/peer/unknown"), None); } +#[test] +fn test_bucket_make_retry_replays_matching_configure_before_settlement() { + let make_photos = + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning".to_string(); + let configure_photos = + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=configure-replication".to_string(); + let plan = SiteReplicationBootstrapPlan { + bucket_make_ops: vec![ + make_photos.clone(), + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=videos&operation=make-with-versioning".to_string(), + ], + bucket_configure_ops: vec![ + "/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() + }; + + let tasks = bucket_op_retry_replay_tasks(&plan, SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING, "photos") + .expect("make retry plan should include its configure follow-up"); + assert_eq!( + tasks.iter().map(SiteReplicationRepairTask::path).collect::>(), + 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::BucketMetadata(item) if item.bucket == "photos")); + assert!(matches!(tasks[2], SiteReplicationRepairTask::Replication(_))); +} + +#[test] +fn test_bucket_make_retry_without_matching_configure_fails_closed() { + let plan = SiteReplicationBootstrapPlan { + bucket_make_ops: vec![ + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning".to_string(), + ], + ..Default::default() + }; + + let err = match bucket_op_retry_replay_tasks(&plan, SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING, "photos") { + Ok(_) => panic!("make retry must not settle without a matching configure operation"), + Err(err) => err, + }; + assert_eq!(err.code(), &S3ErrorCode::InternalError); +} + +#[test] +fn test_retry_drain_bounds_each_peer_round_to_one_small_request_chain() { + let plan = SiteReplicationBootstrapPlan { + iam_items: vec![SRIAMItem::default(); 3], + bucket_make_ops: vec![ + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning".to_string(), + ], + bucket_configure_ops: vec![ + "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=configure-replication".to_string(), + ], + bucket_items: vec![SRBucketMeta { + bucket: "photos".to_string(), + r#type: "tags".to_string(), + ..Default::default() + }], + ..Default::default() + }; + let make = RetryDrainAction::BucketOpReplay { + operation: SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING.to_string(), + bucket: "photos".to_string(), + }; + + assert!(is_lightweight_retry_drain_action(&make)); + assert!(!is_lightweight_retry_drain_action(&RetryDrainAction::IamSnapshot)); + assert!(!is_lightweight_retry_drain_action(&RetryDrainAction::PeerEdit)); + assert_eq!(retry_drain_request_count(&make, Some(&plan)), 3); + assert!(retry_drain_request_count(&make, Some(&plan)) <= SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER); + assert!( + retry_drain_request_count(&RetryDrainAction::IamSnapshot, Some(&plan)) + > SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER + ); + assert!( + retry_drain_request_count(&RetryDrainAction::PeerEdit, Some(&plan)) > SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER + ); +} + +#[test] +fn test_lightweight_retry_peer_rotation_covers_all_queued_peers() { + let limit = SITE_REPLICATION_RETRY_DRAIN_PEER_CONCURRENCY; + for peer_count in 1..=(limit * 3 + 1) { + let rounds = peer_count.div_ceil(limit); + let mut seen = HashSet::new(); + for round in 7..(7 + rounds as i64) { + let start = lightweight_retry_peer_rotation(peer_count, round); + for offset in 0..limit.min(peer_count) { + seen.insert((start + offset) % peer_count); + } + } + assert_eq!( + seen.len(), + peer_count, + "every peer must enter the bounded lightweight window within {rounds} rounds" + ); + } +} + +#[test] +fn test_lightweight_bucket_retry_plan_is_targeted_and_preserves_make_options() { + let created_at = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp"); + let bucket = SRBucketInfo { + bucket: "photos".to_string(), + created_at: Some(created_at), + object_lock_config: Some(String::new()), + tags: Some("dGFncy14bWw=".to_string()), + tag_config_updated_at: Some(created_at), + ..Default::default() + }; + let plan = site_replication_bucket_retry_plan_for(&bucket, false).expect("targeted retry plan"); + + assert!(plan.iam_items.is_empty()); + assert_eq!(plan.bucket_make_ops.len(), 1); + assert!(plan.bucket_make_ops[0].contains("bucket=photos")); + assert!(plan.bucket_make_ops[0].contains("lockEnabled=true")); + assert!(plan.bucket_make_ops[0].contains("createdAt=")); + assert_eq!(plan.bucket_items.len(), 2); + assert_eq!(plan.bucket_items[0].r#type, "tags"); + assert_eq!(plan.bucket_items[1].r#type, "object-lock-config"); + assert_eq!(plan.bucket_configure_ops.len(), 1); + assert!(plan.bucket_configure_ops[0].contains("operation=configure-replication")); + + let tasks = + bucket_op_retry_replay_tasks(&plan, SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING, "photos").expect("retry task chain"); + assert!(matches!(tasks[0], SiteReplicationRepairTask::BucketMake(_))); + assert!(matches!(tasks[1], SiteReplicationRepairTask::BucketMetadata(_))); + assert!(matches!(tasks[2], SiteReplicationRepairTask::BucketMetadata(_))); + assert!(matches!(tasks[3], SiteReplicationRepairTask::Replication(_))); +} + +#[test] +fn test_lightweight_bucket_retry_plan_orders_real_metadata_and_counts_it() { + let versioning = bucket_versioning_xml().expect("canonical versioning config"); + let replication = serialize(&site_repl_config("remote-dep")).expect("derived replication config"); + let bucket = SRBucketInfo { + bucket: "photos".to_string(), + policy: Some(serde_json::json!({"Version":"2012-10-17","Statement":[]})), + tags: Some(BASE64_STANDARD.encode_to_string("")), + versioning: Some(BASE64_STANDARD.encode_to_string(&versioning)), + replication_config: Some(BASE64_STANDARD.encode_to_string(&replication)), + ..Default::default() + }; + let plan = site_replication_bucket_retry_plan_from_info(&bucket, false).expect("targeted retry plan"); + let tasks = bucket_op_retry_replay_tasks(&plan, SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING, "photos") + .expect("bucket replay tasks"); + + assert!(matches!(tasks.first(), Some(SiteReplicationRepairTask::BucketMake(_)))); + assert!(matches!(tasks.last(), Some(SiteReplicationRepairTask::Replication(_)))); + assert!( + tasks[1..tasks.len() - 1] + .iter() + .all(|task| matches!(task, SiteReplicationRepairTask::BucketMetadata(_))) + ); + assert_eq!(tasks.len(), 4, "make + policy + tags + configure must all count against the budget"); + assert!( + tasks.len() <= SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER, + "the complete metadata chain must fit the bounded lightweight replay" + ); + + let mut operator_replication = site_repl_config("remote-dep"); + operator_replication.rules.push(operator_rule("operator-backup")); + let mut bucket_with_operator_rule = bucket; + bucket_with_operator_rule.replication_config = + Some(BASE64_STANDARD.encode_to_string(&serialize(&operator_replication).expect("operator replication config"))); + let plan = site_replication_bucket_retry_plan_from_info(&bucket_with_operator_rule, false).expect("targeted retry plan"); + assert!( + plan.bucket_items.iter().any(|item| item.r#type == "replication-config"), + "operator-authored replication rules cannot be replaced by configure-replication" + ); +} + +#[test] +fn test_delete_bucket_broadcast_fences_target_membership_through_delivery() { + let hooks = include_str!("hooks.rs"); + let delete_broadcast = hooks + .split("async fn broadcast_site_replication_delete_bucket") + .nth(1) + .and_then(|rest| rest.split("pub(crate) async fn commit_site_replication_delete_bucket").next()) + .expect("delete-bucket broadcast should exist"); + assert!( + delete_broadcast.contains("with_site_replication_state_read_lock(move |state| async move {") + && delete_broadcast.contains("state.peers.get(&fallback_peer.deployment_id)") + && delete_broadcast.contains("site_replicator_service_account_secret(&state.service_account_access_key)") + && delete_broadcast + .contains("PeerAdminRequest::put(&transport.connection, &delivery_path, &state.service_account_access_key)"), + "a destructive bucket delivery must resolve current topology and credentials under the distributed state read lock" + ); + assert!( + delete_broadcast.contains("enqueue_site_replication_retry_event(¤t_peer, &request_path, &err).await"), + "a failed destructive delivery must remain visible for operator repair" + ); + + let usecase = include_str!("../app/bucket_usecase.rs"); + let delete = usecase + .split("async fn execute_delete_bucket_inner") + .nth(1) + .and_then(|rest| rest.split("pub async fn execute_head_bucket").next()) + .expect("delete bucket usecase"); + assert!( + delete + .find("prepare_site_replication_delete_bucket") + .expect("durable reservation") + < delete.find(".delete_bucket(").expect("local delete"), + "destructive peer liabilities must be persisted before the local bucket is deleted" + ); +} + +#[test] +fn test_bucket_retry_settlement_preserves_a_newer_same_path_failure() { + let peer = peer("remote", "https://remote.example.com"); + let path = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=configure-replication"; + let observed_at = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp"); + let observed = drain_event("remote", path, 1, Some(observed_at)); + let mut queue = vec![observed.clone()]; + + queue[0].id = "evt-remote-new-revision".to_string(); + queue[0].retry_count += 1; + assert_eq!(settle_observed_site_replication_retry_event(&mut queue, &peer, &observed), 0); + assert_eq!(queue.len(), 1, "a newer same-timestamp failure must survive stale settlement"); + + let current = queue[0].clone(); + assert_eq!(settle_observed_site_replication_retry_event(&mut queue, &peer, ¤t), 1); + assert!(queue.is_empty()); +} + +#[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() + }; + state + .peers + .insert("remote".to_string(), peer("remote", "https://remote.example.com")); + + 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); + assert_eq!( + actionable_site_replication_retry_events(&state, now).len(), + 1, + "a successful probe must make the event replayable in the same drain tick" + ); + + 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 { @@ -828,6 +1146,100 @@ fn test_site_replication_retry_backoff_schedule() { assert!(elapsed(30, 86_401)); } +#[test] +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"; + + upsert_site_replication_retry_event( + &mut queue, + &peer, + bucket_make, + "peer request to https://remote.example.com failed (connect): connection refused", + None, + ) + .expect("upsert retry event"); + 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, + ) + .expect("upsert retry event"); + assert!( + !queue[0].peer_unreachable, + "a whole-request timeout does not prove the peer is unreachable" + ); + + upsert_site_replication_retry_event( + &mut queue, + &peer, + bucket_make, + "peer request to https://remote.example.com failed with 500 Internal Server Error: downstream failed (connect)", + None, + ) + .expect("upsert retry event"); + assert!( + !queue[0].peer_unreachable, + "application failures and their untrusted bodies must keep the normal replay backoff" + ); + + upsert_site_replication_retry_event( + &mut queue, + &peer, + bucket_make, + "peer request to https://remote.example.com failed with 500 Internal Server Error: backend failed (connect): spoofed", + None, + ) + .expect("upsert retry event"); + assert!(!queue[0].peer_unreachable, "peer response bodies must not spoof transport failures"); +} + +#[test] +fn test_connect_timeout_is_classified_as_a_connection_failure() { + assert_eq!(classify_peer_transport_error(true, true, "tcp connect timed out"), "connect"); + assert_eq!(classify_peer_transport_error(false, true, "request timed out"), "timeout"); + assert_eq!( + classify_peer_transport_error(false, true, "request timed out for https://tls-gateway.example"), + "timeout" + ); + assert_eq!(classify_peer_transport_error(true, false, "tls handshake failed"), "tls handshake"); +} + +#[test] +fn test_retry_event_peer_unreachable_is_legacy_serde_default() { + let json = r#"{ + "id":"evt-legacy", + "peer_deployment_id":"remote", + "peer_endpoint":"https://remote.example.com", + "path":"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning", + "retry_count":1, + "failed":false, + "last_error":"peer request to https://remote.example.com failed (connect): connection refused" + }"#; + + let mut event: SiteReplicationRetryEvent = serde_json::from_str(json).expect("legacy retry event decodes"); + assert!(!event.peer_unreachable); + + let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp"); + event.updated_at = Some(now - time::Duration::seconds(30)); + let mut state = SiteReplicationState::default(); + state + .peers + .insert("remote".to_string(), peer("remote", "https://remote.example.com")); + state.retry_queue.push(event); + + assert_eq!( + deferred_site_replication_retry_events(&state, now).len(), + 1, + "rolling-upgrade records must retain fast recovery from their trusted outer error shape" + ); +} + /// The actionable subset respects classification, peer membership and /// backoff; everything else stays untouched in the queue. #[test] @@ -915,6 +1327,51 @@ fn test_deferred_site_replication_retry_events_partition() { assert_eq!(actionable[0].path, "/rustfs/admin/v3/site-replication/peer/bucket-meta"); } +#[test] +fn test_deferred_retry_events_probe_fresh_peer_transport_failures() { + let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp"); + let mut state = SiteReplicationState::default(); + state + .peers + .insert("remote".to_string(), peer("remote", "https://remote.example.com")); + + let bucket_make = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning"; + let mut fresh_transport_failure = drain_event("remote", bucket_make, 1, Some(now - time::Duration::seconds(30))); + fresh_transport_failure.peer_unreachable = true; + state.retry_queue.push(fresh_transport_failure); + + let deferred = deferred_site_replication_retry_events(&state, now); + assert_eq!( + deferred.len(), + 1, + "fresh transport failures must be eligible for a cheap reachability probe" + ); + assert_eq!(deferred[0].path, bucket_make); + + let actionable = actionable_site_replication_retry_events(&state, now); + assert!(actionable.is_empty(), "the event is still protected from direct replay by normal backoff"); +} + +#[test] +fn test_deferred_retry_events_do_not_probe_fresh_application_failures() { + let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp"); + let mut state = SiteReplicationState::default(); + state + .peers + .insert("remote".to_string(), peer("remote", "https://remote.example.com")); + + let bucket_make = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning"; + state + .retry_queue + .push(drain_event("remote", bucket_make, 1, Some(now - time::Duration::seconds(30)))); + + assert!( + deferred_site_replication_retry_events(&state, now).is_empty(), + "reachable peers that reject an operation must keep the base replay backoff" + ); + assert!(actionable_site_replication_retry_events(&state, now).is_empty()); +} + /// The drain settles a peer-edit success under a freshly allocated /// generation; legacy queue entries carry `edit_generation: None` and /// must be cleared by that generation-scoped settlement (`(Some, None)` @@ -996,7 +1453,7 @@ fn test_escalate_up_to_marks_snapshot_replayed_and_keeps_newer_failures() { // successful Bob update on the shared wire path cannot erase it even // before the drain runs. let mut queue = Vec::new(); - upsert_site_replication_retry_event(&mut queue, &target, path, "alice delete failed", None); + upsert_site_replication_retry_event(&mut queue, &target, path, "alice delete failed", None).expect("upsert retry event"); assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH); assert_eq!(dequeue_site_replication_retry_events(&mut queue, &target, path), 0); assert_eq!(queue.len(), 1); @@ -1005,7 +1462,7 @@ fn test_escalate_up_to_marks_snapshot_replayed_and_keeps_newer_failures() { // A later hook failure overwrites the marker and re-arms the drain. let mut queue = vec![drain_event("remote", path, 2, Some(snapshot_at))]; escalate_site_replication_retry_events_up_to(&mut queue, &target, path, Some(snapshot_at)); - upsert_site_replication_retry_event(&mut queue, &target, path, "peer offline", None); + upsert_site_replication_retry_event(&mut queue, &target, path, "peer offline", None).expect("upsert retry event"); assert!(classify_site_replication_retry_event(&queue[0]).is_some()); // Legacy entry without a timestamp: escalated. @@ -1781,17 +2238,85 @@ fn test_retry_event_upsert_marks_repeated_failures() { }; let mut queue = Vec::new(); - upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "first", None); - upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "second", None); - upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "third", None); + upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "first", None) + .expect("upsert retry event"); + let first_revision = queue[0].id.clone(); + upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "second", None) + .expect("upsert retry event"); + let second_revision = queue[0].id.clone(); + upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "third", None) + .expect("upsert retry event"); assert_eq!(queue.len(), 1); + assert_ne!(first_revision, second_revision); + assert_ne!(second_revision, queue[0].id, "each failure must advance the settlement revision"); assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH); assert_eq!(queue[0].retry_count, SITE_REPLICATION_RETRY_FAILED_AFTER); assert!(queue[0].failed); assert_eq!(queue[0].last_error, "third"); } +#[test] +fn retry_queue_capacity_never_evicts_destructive_bucket_liabilities() { + let target = PeerInfo { + deployment_id: "remote-dep".to_string(), + ..peer("remote", "https://remote.example.com") + }; + let destructive = |index: usize| SiteReplicationRetryEvent { + id: format!("delete-{index}"), + peer_deployment_id: target.deployment_id.clone(), + peer_endpoint: target.endpoint.clone(), + path: format!("{SITE_REPLICATION_PEER_BUCKET_OPS_PATH}?bucket=bucket-{index}&operation=delete-bucket"), + ..Default::default() + }; + let mut queue = (0..SITE_REPLICATION_RETRY_QUEUE_LIMIT).map(destructive).collect::>(); + let original_ids = queue.iter().map(|event| event.id.clone()).collect::>(); + let new_path = format!("{SITE_REPLICATION_PEER_BUCKET_OPS_PATH}?bucket=overflow&operation=force-delete-bucket"); + + let err = upsert_site_replication_retry_event(&mut queue, &target, &new_path, "reserve delete", None) + .expect_err("an all-destructive full queue must fail closed"); + assert_eq!(err.code(), &S3ErrorCode::ServiceUnavailable); + assert_eq!(queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT); + assert_eq!(queue.iter().map(|event| event.id.clone()).collect::>(), original_ids); + + queue[0] = SiteReplicationRetryEvent { + id: "iam-snapshot".to_string(), + peer_deployment_id: target.deployment_id.clone(), + peer_endpoint: target.endpoint.clone(), + path: SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH.to_string(), + deletions_recorded: true, + ..Default::default() + }; + upsert_site_replication_retry_event(&mut queue, &target, &new_path, "reserve delete", None) + .expect_err("a collapsed IAM liability may contain a deletion and must not be evicted"); + assert!(queue.iter().any(|event| event.id == "iam-snapshot")); + + queue[0] = SiteReplicationRetryEvent { + id: "rebuildable".to_string(), + peer_deployment_id: target.deployment_id.clone(), + peer_endpoint: target.endpoint.clone(), + path: SITE_REPLICATION_PEER_EDIT_PATH.to_string(), + ..Default::default() + }; + let preserved_delete_ids = queue + .iter() + .filter(|event| is_destructive_bucket_retry_path(&event.path)) + .map(|event| event.id.clone()) + .collect::>(); + let evicted = upsert_site_replication_retry_event(&mut queue, &target, &new_path, "reserve delete", None) + .expect("a rebuildable row may make room for a destructive liability"); + + assert_eq!(evicted.len(), 1); + assert_eq!(evicted[0].id, "rebuildable"); + assert_eq!(queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT); + assert!( + preserved_delete_ids + .iter() + .all(|id| queue.iter().any(|event| &event.id == id)) + ); + assert!(queue.iter().any(|event| event.path == new_path)); +} + /// P1-15 review follow-up: a successful peer-edit delivery only proves the /// peer reached the state THAT delivery carried. Settling it must not /// erase a retry event a newer edit left behind, or the local site sits on @@ -1807,7 +2332,8 @@ fn retry_settlement_must_not_erase_a_newer_generation_failure() { // Edit A (generation 5) delivered successfully and is stalled before // settling. Edit B (generation 6) commits meanwhile, fails delivery to // the same peer, and enqueues. - upsert_site_replication_retry_event(&mut queue, &peer, SITE_REPLICATION_PEER_EDIT_PATH, "peer offline", Some(6)); + upsert_site_replication_retry_event(&mut queue, &peer, SITE_REPLICATION_PEER_EDIT_PATH, "peer offline", Some(6)) + .expect("upsert retry event"); // A resumes: its own settlement must leave B's retry alone. assert_eq!( @@ -1818,7 +2344,8 @@ fn retry_settlement_must_not_erase_a_newer_generation_failure() { assert_eq!(queue[0].edit_generation, Some(6)); // An even older delivery failing afterwards must not lower the fence. - upsert_site_replication_retry_event(&mut queue, &peer, SITE_REPLICATION_PEER_EDIT_PATH, "still offline", Some(4)); + upsert_site_replication_retry_event(&mut queue, &peer, SITE_REPLICATION_PEER_EDIT_PATH, "still offline", Some(4)) + .expect("upsert retry event"); assert_eq!(queue[0].edit_generation, Some(6)); // B's own delivery succeeding is what clears it. @@ -1831,7 +2358,7 @@ fn retry_settlement_must_not_erase_a_newer_generation_failure() { // Collapsed broadcast failures live under an internal snapshot path; // an unrelated success on their shared wire path cannot settle them. let iam_path = "/rustfs/admin/v3/site-replication/peer/iam-item"; - upsert_site_replication_retry_event(&mut queue, &peer, iam_path, "peer offline", None); + upsert_site_replication_retry_event(&mut queue, &peer, iam_path, "peer offline", None).expect("upsert retry event"); assert_eq!(dequeue_site_replication_retry_events(&mut queue, &peer, iam_path), 0); assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH); } diff --git a/rustfs/src/site_replication/transport.rs b/rustfs/src/site_replication/transport.rs index 30ad28f95..bc941c9f0 100644 --- a/rustfs/src/site_replication/transport.rs +++ b/rustfs/src/site_replication/transport.rs @@ -331,6 +331,7 @@ pub(crate) fn runtime_peer_connection(peer: &PeerInfo) -> S3Result PeerAdminRequest<'a> { } let response = req.send().await.map_err(|e| { - let classify = if e.is_timeout() { - "timeout" - } else if e.is_connect() && e.to_string().to_ascii_lowercase().contains("dns") { - "dns resolution" - } else if e.to_string().to_ascii_lowercase().contains("certificate") - || e.to_string().to_ascii_lowercase().contains("tls") - { - "tls handshake" - } else if e.is_connect() { - "connect" - } else { - "request" - }; + let classify = classify_peer_transport_error(e.is_connect(), e.is_timeout(), &e.to_string()); S3Error::with_message(S3ErrorCode::InternalError, format!("peer request to {url} failed ({classify}): {e}")) })?; @@ -825,6 +814,21 @@ impl<'a> PeerAdminRequest<'a> { } } +pub(crate) fn classify_peer_transport_error(is_connect: bool, is_timeout: bool, detail: &str) -> &'static str { + let detail = detail.to_ascii_lowercase(); + if is_connect && detail.contains("dns") { + "dns resolution" + } else if is_connect && (detail.contains("certificate") || detail.contains("tls")) { + "tls handshake" + } else if is_connect { + "connect" + } else if is_timeout { + "timeout" + } else { + "request" + } +} + pub(crate) fn peer_error_may_be_secret_mismatch(detail: &str) -> bool { let detail = detail.to_ascii_lowercase(); detail.contains("signaturedoesnotmatch") diff --git a/rustfs/src/site_replication_reconcile.rs b/rustfs/src/site_replication_reconcile.rs index e327e3669..b4ea52c93 100644 --- a/rustfs/src/site_replication_reconcile.rs +++ b/rustfs/src/site_replication_reconcile.rs @@ -28,16 +28,19 @@ use std::pin::Pin; use std::sync::OnceLock; use std::time::Duration; +use tokio::time::Instant; use tokio_util::sync::CancellationToken; use tracing::warn; const RECONCILE_INTERVAL: Duration = Duration::from_secs(600); +pub(crate) const RETRY_DRAIN_INTERVAL: Duration = Duration::from_secs(30); /// A reconciler reports its own failures; the outcome carries no value because neither /// caller can act on one — a site that cannot repair its replication wiring still serves S3. type ReconcileHook = fn() -> Pin + Send>>; static RECONCILER: OnceLock = OnceLock::new(); +static RETRY_DRAINER: OnceLock = OnceLock::new(); /// Install the admin layer's reconciler. Idempotent: a second call is ignored, which keeps /// repeated router construction (tests, the embedded server) from panicking. @@ -45,6 +48,12 @@ pub(crate) fn register_site_replication_reconciler(reconcile: ReconcileHook) { let _ = RECONCILER.set(reconcile); } +/// Install the admin layer's lightweight retry drain. Idempotent for the same +/// reason as [`register_site_replication_reconciler`]. +pub(crate) fn register_site_replication_retry_drainer(drain: ReconcileHook) { + let _ = RETRY_DRAINER.set(drain); +} + /// Repair drifted site-replication wiring, immediately and then on a timer. /// /// The first pass runs inside the spawned task rather than on the caller's path: it walks @@ -62,16 +71,38 @@ pub(crate) fn spawn_site_replication_reconcile_task(ctx: CancellationToken) { return; } + spawn_reconcile_loop(ctx.clone(), RECONCILE_INTERVAL, &RECONCILER, true); + + if RETRY_DRAINER.get().is_none() { + warn!("site replication retry drainer is not registered; periodic retry drain disabled"); + return; + } + spawn_reconcile_loop(ctx, RETRY_DRAIN_INTERVAL, &RETRY_DRAINER, false); +} + +fn spawn_reconcile_loop( + ctx: CancellationToken, + interval: Duration, + hook: &'static OnceLock, + run_immediately: bool, +) { tokio::spawn(async move { - let mut ticker = tokio::time::interval(RECONCILE_INTERVAL); + let first_tick = if run_immediately { + Instant::now() + } else { + Instant::now() + interval + }; + let mut ticker = tokio::time::interval_at(first_tick, interval); ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { tokio::select! { _ = ctx.cancelled() => break, - // The first tick fires immediately, which is the startup repair pass. + // The heavy reconciler owns the startup repair pass. The lightweight + // retry drain starts on its normal cadence so it cannot steal that + // first lifecycle lock and defer bucket/IAM repair for a full interval. _ = ticker.tick() => { - if let Some(reconcile) = RECONCILER.get() { + if let Some(reconcile) = hook.get() { reconcile().await; } } @@ -79,3 +110,14 @@ pub(crate) fn spawn_site_replication_reconcile_task(ctx: CancellationToken) { } }); } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn retry_drain_runs_faster_than_heavy_reconcile() { + assert!(RETRY_DRAIN_INTERVAL < RECONCILE_INTERVAL); + assert!(RETRY_DRAIN_INTERVAL <= Duration::from_secs(60)); + } +} diff --git a/rustfs/src/storage_api.rs b/rustfs/src/storage_api.rs index e74e88d02..db459b72d 100644 --- a/rustfs/src/storage_api.rs +++ b/rustfs/src/storage_api.rs @@ -244,8 +244,8 @@ pub(crate) mod site_replication { pub(crate) use crate::storage::storage_api::{Endpoint, Endpoints, PoolEndpoints}; pub(crate) use crate::storage::storage_api::{ - ECStore, EndpointServerPools, StorageError, delete_config_no_lock, lock_bucket_targets_metadata, read_config, - read_config_no_lock, save_config_no_lock, with_config_object_read_lock, with_config_object_write_lock, + ECStore, EndpointServerPools, StorageError, delete_config_no_lock, is_err_bucket_not_found, lock_bucket_targets_metadata, + read_config, read_config_no_lock, save_config_no_lock, with_config_object_read_lock, with_config_object_write_lock, }; pub(crate) mod metadata_sys {