diff --git a/docs/architecture/global-state-inventory.md b/docs/architecture/global-state-inventory.md index a594bc569..6c06b10c8 100644 --- a/docs/architecture/global-state-inventory.md +++ b/docs/architecture/global-state-inventory.md @@ -111,7 +111,7 @@ inventory. Generic function-local names such as `CACHE`, `LOCK`, `INIT`, and | `GET_OBJECT_BUFFER_THRESHOLD_WARNED`, `GET_READER_STREAM_BUFFER_SIZE_OVERRIDE`, function-local `ENABLED`, `OBJECT_SEEK_SUPPORT_THRESHOLD`, `OBJECT_SEEK_SUPPORT_CONCURRENCY_THRESHOLDS` | `rustfs/src/app/object_usecase.rs` | Cache or constant / owner-local cache | Object GET/seek tuning caches and warning guards stay private to object usecase helpers. | | `SUPPORTED_HEADERS` | `rustfs/src/storage/options.rs` | Cache or constant / owner-local constant | Supported-header lookup state stays private to storage option parsing. | | `AUDIT_TARGET_SPECS`, `NOTIFICATION_TARGET_SPECS` | `rustfs/src/admin/handlers/audit.rs`, `rustfs/src/admin/handlers/event.rs`, `rustfs/src/admin/handlers/plugins_instances.rs` | Cache or constant / owner-local constant | Admin target descriptor tables stay private to their handler owners. | -| `SITE_REPLICATION_PEER_CLIENT`, `SITE_REPLICATION_STATE_LOCK` | `rustfs/src/admin/handlers/site_replication.rs` | Process-global owner-local cache / guard | Site-replication peer client cache and state lock stay private to site-replication handlers. | +| `SITE_REPLICATION_PEER_CLIENT` | `rustfs/src/admin/handlers/site_replication.rs` | Process-global owner-local cache | Site-replication peer client cache stays private to site-replication handlers. The state RMW transaction holds no process-local mutex — see `rustfs/src/admin/site_replication_state.rs`. | | `AUDIT_MODULE_ENABLED`, `NOTIFY_MODULE_ENABLED`, `PERSISTED_NOTIFY_MODULE_ENABLED`, `PERSISTED_AUDIT_MODULE_ENABLED`, `PERSISTED_MODULE_SWITCH_CONFIGURED` | `rustfs/src/server/audit.rs`, `rustfs/src/server/event.rs`, `rustfs/src/server/module_switch.rs` | Process-global owner-local toggles | Audit/notify module snapshots stay private to the server module switch owners. | | `DELETE_TAIL_TOTAL`, `DELETE_CLEANUP_TOTAL`, `DELETE_REPLICATION_TOTAL`, `DELETE_NOTIFY_TOTAL` | `rustfs/src/delete_tail_activity.rs` | Process-global owner-local counters | Delete-tail activity counters stay private behind delete-tail activity helpers. | | `EMBEDDED_SERVER_STARTED` | `rustfs/src/startup_lifecycle.rs` | Process-global owner-local guard | Embedded startup single-start protection stays private to startup lifecycle. | diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index bf946c0c4..be2220d64 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -35,7 +35,9 @@ use crate::admin::storage_api::bucket::target::{ARN, BucketTarget, BucketTargetT use crate::admin::storage_api::bucket::target_sys::BucketTargetSys; use crate::admin::storage_api::bucket::utils::{deserialize, serialize}; use crate::admin::storage_api::bucket::{AdminReplicationConfigExt as _, AdminVersioningConfigExt as _}; -use crate::admin::storage_api::config::{delete_admin_config, read_admin_config, save_admin_config}; +use crate::admin::storage_api::config::read_admin_config; +#[cfg(test)] +use crate::admin::storage_api::config::save_admin_config; use crate::admin::storage_api::contract::bucket::{ BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp, }; @@ -112,13 +114,13 @@ const LOG_COMPONENT_ADMIN: &str = "admin"; const LOG_SUBSYSTEM_SITE_REPLICATION: &str = "site_replication"; const EVENT_ADMIN_SITE_REPLICATION_STATE: &str = "admin_site_replication_state"; const SERVICE_ACCOUNT_ENVELOPE_VERSION: u64 = 2; -#[cfg(test)] -use crate::admin::site_replication_state::with_site_replication_state_object_lock; -use crate::admin::site_replication_state::{ - SITE_REPLICATION_STATE_PATH, site_replication_state_process_guard, with_site_replication_state_lock, -}; +use crate::admin::site_replication_state::{SITE_REPLICATION_STATE_PATH, with_site_replication_state_lock}; const SITE_REPLICATION_REPAIR_STATE_PATH: &str = "config/site-replication/repair-state.json"; const SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH: &str = "config/site-replication/repair-execution.lock"; +// Serializes peer-join admission (staleness check -> IAM upsert -> state +// commit) across every node of this site; see admit_peer_join. Never an +// actual object — only a namespace-lock key, like the repair execution lock. +const SITE_REPLICATION_JOIN_ADMISSION_LOCK_PATH: &str = "config/site-replication/join-admission.lock"; const SITE_REPL_ADD_SUCCESS: &str = "Requested sites were configured for replication successfully."; const SITE_REPL_EDIT_SUCCESS: &str = "Requested site was updated successfully."; const SITE_REPL_REMOVE_SUCCESS: &str = "Requested site(s) were removed from cluster replication successfully."; @@ -354,8 +356,11 @@ impl TryFrom<&PeerSite> for PeerConnection { static SITE_REPLICATION_PEER_CLIENT: LazyLock>> = LazyLock::new(|| Mutex::new(None)); // Lock order: lifecycle -> bucket operation -> repair admission -> state -> per-bucket metadata. -// The state mutex lives in crate::admin::site_replication_state and is -// taken through site_replication_state_process_guard (P1-15). +// "state" is the distributed state-object lock in +// crate::admin::site_replication_state, entered through +// update_site_replication_state (P1-15). There is no process-local state +// mutex any more: it could not order two nodes of one site, and the call +// sites that needed ordering carry a generation fence instead. static SITE_REPLICATION_LIFECYCLE_LOCK: LazyLock> = LazyLock::new(|| Mutex::new(())); static SITE_REPLICATION_BUCKET_OP_LOCK: LazyLock> = LazyLock::new(|| RwLock::new(())); static SITE_REPLICATION_ADD_BOOTSTRAP: LazyLock>> = @@ -469,11 +474,10 @@ struct SiteReplicationState { #[serde(default)] sync_state_initialized: bool, /// Fencing token for peer-edit delivery, allocated inside the state - /// transaction (process mutex + distributed state-object lock). Two nodes - /// of THIS site that accept admin edits concurrently therefore get - /// strictly ordered generations even though the process mutex cannot - /// serialize them, and a delivery that stalls can be recognised as stale - /// by the receiving site. + /// transaction (the distributed state-object lock). Two nodes of THIS + /// site that accept admin edits concurrently therefore get strictly + /// ordered generations, and a delivery that stalls can be recognised as + /// stale by the receiving site. #[serde(default)] edit_generation: u64, /// Per-origin high-water mark of the peer edits already applied here, @@ -1123,44 +1127,49 @@ async fn persist_site_replication_state_no_lock(store: Arc, mut state: } } -/// A second node's view of the same transaction: identical production code -/// path minus the process mutex, which is per process and therefore cannot -/// serialize anything across nodes. Used by the separate-nodes regression -/// test so that removing the distributed lock breaks it. -#[cfg(test)] -async fn update_site_replication_state_as_separate_node(update: F) -> S3Result -where - T: Send + 'static, - F: FnOnce(&mut SiteReplicationState) -> S3Result + Send + 'static, -{ - let store = - current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?; - let lock_store = store.clone(); - with_site_replication_state_object_lock(lock_store, move || async move { - let mut state = load_site_replication_state_no_lock(store.clone()).await?; - let result = update(&mut state)?; - persist_site_replication_state_no_lock(store, state).await?; - Ok(result) - }) - .await +/// What a state transaction closure decided to do with the state it was +/// handed. `Unchanged` skips the write entirely: the ack markers and the +/// pending-clearing paths run on every retry and mostly find their pending id +/// already gone, and the retry queue shares this object — rewriting it byte +/// for byte only makes those misses contend with the writers that do have +/// something to say. +enum StateCommit { + Changed(T), + Unchanged(T), } /// The site-replication state RMW transaction: load, mutate, persist — all -/// under the process mutex plus the distributed state-object write lock -/// (see crate::admin::site_replication_state). No peer network calls and no -/// other config locks inside `update`. +/// under the distributed state-object write lock (see +/// crate::admin::site_replication_state). No peer network calls and no other +/// config locks inside `update`; anything that has to talk to a peer belongs +/// between two transactions, with the precondition re-checked inside the +/// second one. async fn update_site_replication_state(update: F) -> S3Result where T: Send + 'static, F: FnOnce(&mut SiteReplicationState) -> S3Result + Send + 'static, +{ + update_site_replication_state_when_changed(move |state| update(state).map(StateCommit::Changed)).await +} + +/// [`update_site_replication_state`] for closures that may find nothing to +/// do — see [`StateCommit`]. +async fn update_site_replication_state_when_changed(update: F) -> S3Result +where + T: Send + 'static, + F: FnOnce(&mut SiteReplicationState) -> S3Result> + Send + 'static, { with_site_replication_state_lock(move || async move { let store = current_object_store_handle() .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?; let mut state = load_site_replication_state_no_lock(store.clone()).await?; - let result = update(&mut state)?; - persist_site_replication_state_no_lock(store, state).await?; - Ok(result) + match update(&mut state)? { + StateCommit::Changed(result) => { + persist_site_replication_state_no_lock(store, state).await?; + Ok(result) + } + StateCommit::Unchanged(result) => Ok(result), + } }) .await } @@ -1216,6 +1225,11 @@ where .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("lock repair state failed: {e}")))? } +/// Test-only seeding of the state object. Every production write goes through +/// [`update_site_replication_state`] — this helper is `cfg(test)` so a new +/// call site cannot reintroduce the pre-P1-15 shape (load through one object +/// lock, save through another, with the mutation in between unprotected). +#[cfg(test)] async fn save_site_replication_state(state: &SiteReplicationState) -> S3Result<()> { let Some(store) = current_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); @@ -1232,27 +1246,6 @@ async fn save_site_replication_state(state: &SiteReplicationState) -> S3Result<( Ok(()) } -async fn clear_site_replication_state() -> S3Result<()> { - let Some(store) = current_object_store_handle() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - match delete_admin_config(store, SITE_REPLICATION_STATE_PATH).await { - Ok(()) | Err(StorageError::ConfigNotFound) => Ok(()), - Err(err) => Err(S3Error::with_message(S3ErrorCode::InternalError, format!("clear state failed: {err}"))), - } -} - -async fn persist_site_replication_state(state: &SiteReplicationState) -> S3Result<()> { - let mut normalized = state.clone(); - normalized.peers = normalize_peer_map_by_identity(normalized.peers); - if normalized.peers.len() <= 1 && normalized.pending_rotation.is_none() && normalized.pending_remove.is_none() { - clear_site_replication_state().await - } else { - save_site_replication_state(&normalized).await - } -} - fn build_site_replication_peer_client(outbound_tls: &GlobalPublishedOutboundTlsState) -> S3Result { build_site_replication_peer_client_with_resolver(outbound_tls, PeerDnsResolver::new(loopback_replication_targets_allowed())) } @@ -1668,7 +1661,14 @@ fn stored_peer_tls_settings(stored_peer: Option<&PeerInfo>) -> (bool, String) { } fn current_local_peer(req: &S3Request, state: &SiteReplicationState) -> PeerInfo { - let endpoint = site_replication_local_endpoint(&req.uri, &req.headers); + local_peer_at_endpoint(site_replication_local_endpoint(&req.uri, &req.headers), state) +} + +/// The local peer record as the given state describes it. Split out of +/// [`current_local_peer`] so a state transaction can rebuild it against the +/// state it just loaded: the request the endpoint came from cannot cross into +/// the transaction closure, but the endpoint itself can. +fn local_peer_at_endpoint(endpoint: String, state: &SiteReplicationState) -> PeerInfo { let deployment_id = current_deployment_id().unwrap_or_else(|| deployment_id_for_endpoint(&endpoint)); let stored_peer = state.peers.get(&deployment_id); let (skip_tls_verify, ca_cert_pem) = stored_peer_tls_settings(stored_peer); @@ -1695,30 +1695,7 @@ fn current_local_peer(req: &S3Request, state: &SiteReplicationState) -> Pe } fn current_local_runtime_peer(state: &SiteReplicationState) -> PeerInfo { - let endpoint = current_local_runtime_endpoint(); - let deployment_id = current_deployment_id().unwrap_or_else(|| deployment_id_for_endpoint(&endpoint)); - let stored_peer = state.peers.get(&deployment_id); - let (skip_tls_verify, ca_cert_pem) = stored_peer_tls_settings(stored_peer); - - PeerInfo { - endpoint: endpoint.clone(), - name: if state.name.is_empty() { - stored_peer - .map(|peer| peer.name.clone()) - .filter(|name| !name.is_empty()) - .unwrap_or_else(|| infer_site_name(&endpoint)) - } else { - state.name.clone() - }, - deployment_id, - sync_state: stored_peer.map(|peer| peer.sync_state.clone()).unwrap_or(SyncStatus::Unknown), - default_bandwidth: stored_peer.map(|peer| peer.default_bandwidth.clone()).unwrap_or_default(), - replicate_ilm_expiry: stored_peer.is_some_and(|peer| peer.replicate_ilm_expiry), - object_naming_mode: stored_peer.map(|peer| peer.object_naming_mode.clone()).unwrap_or_default(), - skip_tls_verify, - ca_cert_pem, - api_version: Some(SITE_REPL_API_VERSION.to_string()), - } + local_peer_at_endpoint(current_local_runtime_endpoint(), state) } fn normalize_peer_map_by_identity(peers: BTreeMap) -> BTreeMap { @@ -2535,6 +2512,51 @@ fn initialize_join_peer_sync_state(peers: &mut BTreeMap, defer } } +/// Whether an incoming peer join carries a snapshot this site has already +/// moved past — an unstamped join against a configured site, or one whose +/// `updated_at` is not newer. Applying it would roll the local view back to +/// the older topology, so the join is answered as a no-op (MinIO-compatible +/// behaviour, kept verbatim from the pre-transaction handler). +fn join_request_is_superseded(state: &SiteReplicationState, incoming_updated_at: Option) -> bool { + let Some(current_updated_at) = state.updated_at else { + return false; + }; + incoming_updated_at.is_none_or(|incoming_updated_at| incoming_updated_at <= current_updated_at) +} + +/// Adopt an accepted peer join: the sending site's snapshot replaces the local +/// topology wholesale. +/// +/// The peer-edit high-water marks are deliberately KEPT. Wiping them here +/// would reopen the exact window the fence closes: every join fan-out (adds +/// AND service-account rotations deliver `SRPeerJoin` to existing peers) +/// would discard live marks, letting a stalled older edit from a peer that +/// never left roll a record back. The one case a kept mark misfences — a +/// site removed while unreachable rejoining with a restarted generation +/// counter — already misfences its ordinary edits identically (pre-existing +/// since the fence landed) and needs an epoch in the fence to fix, not a +/// blanket reset. Marks of origins that left AND were observed leaving are +/// dropped on load by `parse_site_replication_state`. +fn apply_peer_join( + state: &mut SiteReplicationState, + local_peer: &PeerInfo, + join_req: SRPeerJoinReq, + defer_sync_state_enable: bool, +) { + state.service_account_access_key = join_req.svc_acct_access_key; + state.service_account_parent = join_req.svc_acct_parent; + state.updated_at = join_req.updated_at.or_else(|| Some(OffsetDateTime::now_utc())); + state.peers = normalize_join_peers_for_local(local_peer, join_req.peers); + initialize_join_peer_sync_state(&mut state.peers, defer_sync_state_enable); + state.sync_state_initialized = true; + state.name = state + .peers + .get(&local_peer.deployment_id) + .map(|peer| peer.name.clone()) + .filter(|name| !name.is_empty()) + .unwrap_or_else(|| local_peer.name.clone()); +} + fn reconcile_peer_with_actual_identity(mut state: SiteReplicationState, actual_peer: PeerInfo) -> SiteReplicationState { let mut actual_peer = normalize_peer_info(actual_peer); if let Some(requested_peer) = state @@ -2736,7 +2758,8 @@ fn is_missing_service_account_error(err: &rustfs_iam::error::Error) -> bool { /// disables every control-plane push while `replicate info` still reports the site enabled. /// Both used to require deleting and recreating the account by hand. async fn reconcile_site_replicator_service_account() -> S3Result<()> { - let _state_guard = site_replication_state_process_guard().await; + // Read-only against the state: `load_site_replication_state` takes the + // object read lock on its own, and everything after it is IAM work. let state = load_site_replication_state().await?; if !state.enabled() || state.service_account_access_key != SITE_REPLICATOR_SERVICE_ACCOUNT { return Ok(()); @@ -3908,29 +3931,30 @@ async fn persist_site_replication_repair_task( ) -> S3Result<()> { persist_site_replication_repair_operation(operation).await?; - let _state_guard = site_replication_state_process_guard().await; - let mut latest = load_site_replication_state().await?; let family_status = operation .sites .get(&peer.deployment_id) .and_then(|site| site.families.get(family)) .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "repair task status is missing".to_string()))?; - if family_status.failed > 0 { - upsert_site_replication_retry_event( - &mut latest.retry_queue, - peer, - path, - family_status - .errors - .first() - .map(String::as_str) - .unwrap_or("remote-operation-failed"), - None, - ); - } else { - dequeue_site_replication_retry_events(&mut latest.retry_queue, peer, path); - } - persist_site_replication_state(&latest).await + let failure = (family_status.failed > 0).then(|| { + family_status + .errors + .first() + .cloned() + .unwrap_or_else(|| "remote-operation-failed".to_string()) + }); + let peer = peer.clone(); + 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), + None => { + dequeue_site_replication_retry_events(&mut state.retry_queue, &peer, &path); + } + } + Ok(()) + }) + .await } fn admit_site_replication_repair_operation( @@ -3986,14 +4010,10 @@ async fn execute_site_replication_repair( async fn execute_site_replication_repair_locked( request: SiteReplicationRepairExecutionRequest, ) -> S3Result> { - let state = { - let _state_guard = site_replication_state_process_guard().await; - let state = load_site_replication_state().await?; - if !state.enabled() || state.service_account_access_key.is_empty() { - return Err(s3_error!(InvalidRequest, "site replication is not configured")); - } - state - }; + let state = load_site_replication_state().await?; + if !state.enabled() || state.service_account_access_key.is_empty() { + return Err(s3_error!(InvalidRequest, "site replication is not configured")); + } let info = build_sr_info(&state, &request.local_peer).await?; let plan = site_replication_bootstrap_plan(&info)?; let plan_token = site_replication_repair_plan_token(&state, &plan)?; @@ -4110,7 +4130,11 @@ async fn execute_site_replication_repair_locked( 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 = { - let _state_guard = site_replication_state_process_guard().await; + // The bucket-op lock is what orders this against add/remove. The + // state is only read here (through the runtime snapshot), and the + // bucket setup below writes bucket metadata, never the state object — + // holding the state transaction across it would put local metadata + // IO inside a distributed lock for nothing. let Some(runtime) = runtime_site_replication_targets().await? else { return Ok(()); }; @@ -5306,6 +5330,29 @@ fn internal_endpoint_refresh_already_committed(state: &SiteReplicationState, inc .is_some_and(|committed| peer_connection_settings_match(committed, incoming)) } +/// An admin add/edit's precondition, re-evaluated inside the transaction that +/// is about to commit: the topology must still be the one the operation was +/// planned against, and the endpoint refresh must still be the same one (or +/// still absent). The planning snapshot is taken before peer probes and +/// fan-outs, none of which may hold the state-object lock, so only the check +/// inside the committing closure binds — the same check between network +/// stages is advisory, fencing the common race off the side-effect path. +/// `stage` names what was in flight for the operator; a rejected commit is +/// safe to re-run. +fn ensure_edit_precondition( + state: &SiteReplicationState, + expected_updated_at: Option, + expected_pending_id: Option<&String>, + stage: &str, +) -> S3Result<()> { + if state.updated_at != expected_updated_at + || pending_endpoint_refresh(state).as_ref().map(|pending| &pending.id) != expected_pending_id + { + return Err(s3_error!(InvalidRequest, "site replication state changed during {stage}")); + } + Ok(()) +} + fn set_pending_endpoint_refresh(state: &mut SiteReplicationState, pending: PendingEndpointRefresh) -> S3Result<()> { state .retry_queue @@ -5924,9 +5971,9 @@ fn peer_edit_fence(queries: &HashMap) -> Option<(String, u64)> { } /// True when a strictly newer edit from the same origin site already landed -/// here. The process mutex on the sending node cannot order deliveries issued -/// by two nodes of that site, so ordering is decided here, on the generation -/// the sender allocated under the distributed lock. Equal generations are NOT +/// here. No lock on the sending side can order deliveries issued by two +/// nodes of that site, so ordering is decided here, on the generation the +/// sender allocated under the distributed lock. Equal generations are NOT /// stale: one edit legitimately fans out several deliveries under a single /// generation (the ILM-expiry edit sends every peer's record), and a replay of /// an applied delivery re-applies the same edit idempotently. @@ -6206,15 +6253,15 @@ async fn record_pending_rotation_secret_candidate(rotation_id: &str, secret: Str return Ok(()); } - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - if let Some(pending) = state.pending_rotation.as_mut() - && pending.id == rotation_id - { + let rotation_id = rotation_id.to_string(); + update_site_replication_state_when_changed(move |state| { + let Some(pending) = state.pending_rotation.as_mut().filter(|pending| pending.id == rotation_id) else { + return Ok(StateCommit::Unchanged(())); + }; push_unique_secret_candidate(&mut pending.secret_candidates, secret); - save_site_replication_state(&state).await?; - } - Ok(()) + Ok(StateCommit::Changed(())) + }) + .await } async fn record_pending_remove_secret_candidate(remove_id: &str, secret: String) -> S3Result<()> { @@ -6222,27 +6269,26 @@ async fn record_pending_remove_secret_candidate(remove_id: &str, secret: String) return Ok(()); } - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - if let Some(pending) = state.pending_remove.as_mut() - && pending.id == remove_id - { + let remove_id = remove_id.to_string(); + update_site_replication_state_when_changed(move |state| { + let Some(pending) = state.pending_remove.as_mut().filter(|pending| pending.id == remove_id) else { + return Ok(StateCommit::Unchanged(())); + }; push_unique_secret_candidate(&mut pending.secret_candidates, secret); - save_site_replication_state(&state).await?; - } - Ok(()) + Ok(StateCommit::Changed(())) + }) + .await } async fn mark_pending_rotation_peer_acked(rotation_id: &str, deployment_id: &str) -> S3Result<()> { let rotation_id = rotation_id.to_string(); let deployment_id = deployment_id.to_string(); - update_site_replication_state(move |state| { - if let Some(pending) = state.pending_rotation.as_mut() - && pending.id == rotation_id - { - pending.acked_deployment_ids.insert(deployment_id); - } - Ok(()) + update_site_replication_state_when_changed(move |state| { + let Some(pending) = state.pending_rotation.as_mut().filter(|pending| pending.id == rotation_id) else { + return Ok(StateCommit::Unchanged(())); + }; + pending.acked_deployment_ids.insert(deployment_id); + Ok(StateCommit::Changed(())) }) .await } @@ -6250,37 +6296,37 @@ async fn mark_pending_rotation_peer_acked(rotation_id: &str, deployment_id: &str async fn mark_pending_remove_peer_acked(remove_id: &str, deployment_id: &str) -> S3Result<()> { let remove_id = remove_id.to_string(); let deployment_id = deployment_id.to_string(); - update_site_replication_state(move |state| { - if let Some(pending) = state.pending_remove.as_mut() - && pending.id == remove_id - { - pending.acked_deployment_ids.insert(deployment_id); - } - Ok(()) + update_site_replication_state_when_changed(move |state| { + let Some(pending) = state.pending_remove.as_mut().filter(|pending| pending.id == remove_id) else { + return Ok(StateCommit::Unchanged(())); + }; + pending.acked_deployment_ids.insert(deployment_id); + Ok(StateCommit::Changed(())) }) .await } async fn finalize_pending_rotation_if_complete(rotation_id: &str, local_peer: &PeerInfo) -> S3Result { - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - let Some(pending) = state.pending_rotation.as_ref() else { - return Ok(true); - }; - if pending.id != rotation_id { - return Ok(false); - } - if !pending_all_remote_peers_acked(&pending.peers, local_peer, &pending.acked_deployment_ids) { - return Ok(false); - } + let rotation_id = rotation_id.to_string(); + let local_peer = local_peer.clone(); + update_site_replication_state_when_changed(move |state| { + let Some(pending) = state.pending_rotation.as_ref() else { + return Ok(StateCommit::Unchanged(true)); + }; + if pending.id != rotation_id { + return Ok(StateCommit::Unchanged(false)); + } + if !pending_all_remote_peers_acked(&pending.peers, &local_peer, &pending.acked_deployment_ids) { + return Ok(StateCommit::Unchanged(false)); + } - state.pending_rotation = None; - persist_site_replication_state(&state).await?; - Ok(true) + state.pending_rotation = None; + Ok(StateCommit::Changed(true)) + }) + .await } async fn pending_remove_ready_to_finalize(remove_id: &str, local_peer: &PeerInfo) -> S3Result> { - let _state_guard = site_replication_state_process_guard().await; let state = load_site_replication_state().await?; let Some(pending) = state.pending_remove.as_ref() else { return Ok(None); @@ -6296,18 +6342,15 @@ async fn pending_remove_ready_to_finalize(remove_id: &str, local_peer: &PeerInfo } async fn clear_pending_remove(remove_id: &str) -> S3Result<()> { - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - if state - .pending_remove - .as_ref() - .map(|pending| pending.id.as_str() == remove_id) - .unwrap_or(false) - { + let remove_id = remove_id.to_string(); + update_site_replication_state_when_changed(move |state| { + if state.pending_remove.as_ref().is_none_or(|pending| pending.id != remove_id) { + return Ok(StateCommit::Unchanged(())); + } state.pending_remove = None; - persist_site_replication_state(&state).await?; - } - Ok(()) + Ok(StateCommit::Changed(())) + }) + .await } fn removed_deployment_ids_for_pending_remove(pending: &PendingRemove, local_peer: &PeerInfo) -> HashSet { @@ -7375,7 +7418,8 @@ async fn refresh_bucket_targets_after_endpoint_edit(pending_id: &str, service_ac let expected_incarnation_id = metadata_sys::capture_bucket_metadata_incarnation(&bucket.name) .await .map_err(ApiError::from)?; - let _state_guard = site_replication_state_process_guard().await; + // Read-only per bucket: the pending refresh is re-read (and re-checked) + // every round, and the writes below are bucket metadata, not state. let state = load_site_replication_state().await?; let Some(pending) = pending_endpoint_refresh(&state).filter(|pending| pending.id == pending_id) else { return Err(s3_error!(InvalidRequest, "endpoint target refresh state changed during update")); @@ -7696,27 +7740,36 @@ async fn refresh_site_resync_status(mut status: SRResyncOpStatus, peer: &PeerInf } async fn persist_site_resync_status(peer_id: &str, status: &SRResyncOpStatus) -> S3Result<()> { - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - if state - .resync_status - .get(peer_id) - .is_some_and(|current| current.resync_id != status.resync_id || current.generation != status.generation) - { - return Err(s3_error!(InvalidRequest, "site replication resync state changed")); - } - state.resync_status.insert(peer_id.to_string(), status.clone()); - save_site_replication_state(&state).await + let peer_id = peer_id.to_string(); + let status = status.clone(); + update_site_replication_state(move |state| { + // The run identity is checked inside the transaction: a cancel or a + // newer run that committed while this progress snapshot was being + // built must not be overwritten by it. + if state + .resync_status + .get(&peer_id) + .is_some_and(|current| current.resync_id != status.resync_id || current.generation != status.generation) + { + return Err(s3_error!(InvalidRequest, "site replication resync state changed")); + } + state.resync_status.insert(peer_id, status); + Ok(()) + }) + .await } async fn persist_new_site_resync_status(peer_id: &str, status: &SRResyncOpStatus) -> S3Result<()> { - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - if state.resync_status.get(peer_id).is_some_and(site_resync_is_active) { - return Err(s3_error!(InvalidRequest, "site replication resync is already active")); - } - state.resync_status.insert(peer_id.to_string(), status.clone()); - save_site_replication_state(&state).await + let peer_id = peer_id.to_string(); + let status = status.clone(); + update_site_replication_state(move |state| { + if state.resync_status.get(&peer_id).is_some_and(site_resync_is_active) { + return Err(s3_error!(InvalidRequest, "site replication resync is already active")); + } + state.resync_status.insert(peer_id, status); + Ok(()) + }) + .await } fn apply_state_edit_req(mut state: SiteReplicationState, body: SRStateEditReq) -> SiteReplicationState { @@ -8352,7 +8405,10 @@ impl Operation for SiteReplicationAddHandler { reject_site_replicator_on_public_admin(&cred)?; let replicate_ilm_expiry = sr_add_replicate_ilm_expiry(&req.uri); let lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; - let state_guard = site_replication_state_process_guard().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")); @@ -8373,13 +8429,14 @@ impl Operation for SiteReplicationAddHandler { } validate_add_preflight_topology(&preflight_infos, &local_peer)?; let expected_updated_at = current_state.updated_at; - drop(state_guard); require_add_peer_tls_capability(&sites, &local_peer).await?; - let _state_guard = site_replication_state_process_guard().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?; - if latest_state.updated_at != expected_updated_at || pending_endpoint_refresh(&latest_state).is_some() { - return Err(s3_error!(InvalidRequest, "site replication state changed during capability probe")); - } + 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?; @@ -8454,13 +8511,63 @@ impl Operation for SiteReplicationAddHandler { } mark_unknown_peer_sync_enabled(&mut state.peers); - persist_site_replication_state(&state).await?; - // The finalize fan-out below delivers peer-edit payloads, so it stays - // under the state guard: ordering against a concurrent edit matters - // and the peer edit handler has no generation fence. It uses the - // plain transport (no retry-event bookkeeping), so nothing re-enters - // the state transaction while the guard is held. + // 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 only the fields this add computed. Everything else is + // owned by writers that commit without touching `updated_at` + // (retry events, peer-edit generations, resync progress, the + // acks/clears of an already pending rotation or removal), so the + // CAS above cannot vouch for them — they keep the freshly loaded + // value. The exhaustive destructure makes adding a state field a + // compile error here until it is classified. + let SiteReplicationState { + name, + service_account_access_key, + service_account_secret_key: _, + service_account_parent, + peers, + updated_at, + resync_status: _, + pending_rotation: _, + pending_remove: _, + pending_endpoint_refresh: _, + retry_queue: _, + sync_state_initialized, + edit_generation: _, + applied_edit_generations: _, + } = next_state; + state.name = name; + state.service_account_access_key = service_account_access_key; + state.service_account_parent = service_account_parent; + state.peers = peers; + state.updated_at = updated_at; + state.sync_state_initialized = sync_state_initialized; + let edit_generation = next_peer_edit_generation(state); + Ok((state.clone(), edit_generation)) + }) + .await?; + + // The finalize fan-out delivers peer-edit payloads, so it carries the + // generation allocated in the commit above: the receiving site orders + // it against any edit that follows instead of applying whichever + // delivery happens to arrive last. It runs outside the transaction — + // holding the state-object lock across peer traffic would block every + // node of this site, including this add's own retry bookkeeping. + let local_deployment_id = current_deployment_id(); + let finalize_edit_path = peer_edit_path_with_fence(local_deployment_id.as_deref(), edit_generation); for target in state.peers.values() { if target.deployment_id == local_peer.deployment_id || same_identity_endpoint(&target.endpoint, &local_peer.endpoint) { @@ -8477,7 +8584,7 @@ impl Operation for SiteReplicationAddHandler { if let Err(err) = send_peer_admin_request_with_client( &transport.client, &transport.connection, - SITE_REPLICATION_PEER_EDIT_PATH, + &finalize_edit_path, &state.service_account_access_key, &service_account_secret_key, peer, @@ -8490,12 +8597,6 @@ impl Operation for SiteReplicationAddHandler { } } - // Bootstrap and back-fill send bucket-ops, not peer edits, so their - // ordering is not state-sensitive — and their transports do record - // retry events, which re-enter the state transaction. Release the - // guard before them. - drop(_state_guard); - initial_sync_errors.extend(bootstrap_existing_metadata_after_add(&state, &local_peer, &service_account_secret_key).await); // Fix 1: back-fill pre-existing buckets so objects created before `replicate add` @@ -8521,39 +8622,49 @@ impl Operation for SiteReplicationRemoveHandler { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationRemoveAction).await?; reject_site_replicator_on_public_admin(&cred)?; let _lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; + // The request body is read before the bucket-op guard and the state + // transaction: a client that stalls mid-body must hold neither the + // state-object lock nor the write half of the bucket-op RwLock (which + // would starve every bucket-operation hook in the meantime). + let local_endpoint = site_replication_local_endpoint(&req.uri, &req.headers); + let remove_req: SRRemoveReq = read_site_replication_json(req, "", false).await?; let (pending_remove, local_peer) = { let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.write().await; - let _state_guard = site_replication_state_process_guard().await; - 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")); - } - if current_state.pending_rotation.is_some() { - return Err(s3_error!(InvalidRequest, "service account rotation is pending")); - } - let local_peer = current_local_peer(&req, ¤t_state); - let remove_req: SRRemoveReq = read_site_replication_json(req, "", false).await?; + update_site_replication_state_when_changed(move |state| { + if pending_endpoint_refresh(state).is_some() { + return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); + } + if state.pending_rotation.is_some() { + return Err(s3_error!(InvalidRequest, "service account rotation is pending")); + } + let local_peer = local_peer_at_endpoint(local_endpoint, state); - if let Some(pending) = current_state.pending_remove.clone() { - (pending, local_peer) - } else { - validate_remove_sites_req(¤t_state, &remove_req)?; - let mut next_state = remove_sites(current_state.clone(), remove_req.clone()); - let mut peer_remove_req = remove_req; + // Resuming: the peers were already told about this pending + // removal, so re-persisting the same record buys nothing. + if let Some(pending) = state.pending_remove.clone() { + return Ok(StateCommit::Unchanged((pending, local_peer))); + } + + validate_remove_sites_req(state, &remove_req)?; + let service_account_access_key = state.service_account_access_key.clone(); + let secret_candidates = legacy_site_replicator_state_secret(state).into_iter().collect(); + let original_peers = state.peers.clone(); + let mut peer_remove_req = remove_req.clone(); peer_remove_req.requesting_dep_id = local_peer.deployment_id.clone(); + *state = remove_sites(std::mem::take(state), remove_req); let pending = PendingRemove { id: Uuid::new_v4().to_string(), req: peer_remove_req, - service_account_access_key: current_state.service_account_access_key.clone(), - secret_candidates: legacy_site_replicator_state_secret(¤t_state).into_iter().collect(), - original_peers: current_state.peers.clone(), + service_account_access_key, + secret_candidates, + original_peers, acked_deployment_ids: BTreeSet::new(), - updated_at: next_state.updated_at, + updated_at: state.updated_at, }; - next_state.pending_remove = Some(pending.clone()); - persist_site_replication_state(&next_state).await?; - (pending, local_peer) - } + state.pending_remove = Some(pending.clone()); + Ok(StateCommit::Changed((pending, local_peer))) + }) + .await? }; let mut peer_errors = Vec::new(); @@ -8724,102 +8835,186 @@ impl Operation for SiteReplicationNetPerfHandler { pub struct SRPeerJoinHandler {} +/// What the join admission decided about an incoming peer join. The verdict — +/// and the committed state the back-fill afterwards needs — travel out of +/// [`admit_peer_join`] instead of being answered where they are decided. +enum PeerJoinOutcome { + Applied(Box, PeerInfo), + /// A newer join already landed here; the sender is answered with the local + /// peer record and nothing is written. + Superseded(PeerInfo), +} + +/// The serialized half of an accepted peer join: staleness check, IAM apply, +/// state commit. +/// +/// Two locks, two scopes. The lifecycle guard (process-local) keeps the +/// admission mutually exclusive with this node's add / remove / rotate / +/// reconciler. The distributed join-admission lock then serializes the +/// admission CLUSTER-WIDE — the IAM write and the state commit cannot share +/// a transaction, so without it two joins accepted by different nodes of +/// this site interleave as "A checks for older T1, B applies secret B and +/// commits newer T2, A overwrites IAM with secret A, A's commit is refused +/// as superseded" — leaving the persisted state advertising B's contract +/// while IAM only accepts A's secret. Under the admission lock the +/// staleness check runs against a load taken INSIDE the lock, before +/// `apply_iam` changes anything, so a superseded join exits without +/// touching IAM at all. Crash safety is the lock subsystem's lease expiry +/// (same pattern as the repair execution lock); the closing transaction +/// still re-checks staleness for defence in depth and for old-version nodes +/// that do not take the admission lock during a rolling upgrade. +/// +/// Lock order: lifecycle -> join admission -> state object lock (the repair +/// path nests config-object locks the same way: repair execution -> state). +/// +/// `apply_iam` is injected so the interleaving regression tests can gate it +/// mid-flight; production passes the real service-account upsert. +async fn admit_peer_join( + local_endpoint: String, + join_req: SRPeerJoinReq, + defer_sync_state_enable: bool, + apply_iam: F, +) -> S3Result +where + F: FnOnce(SRPeerJoinReq) -> Fut + Send + 'static, + Fut: std::future::Future> + Send + 'static, +{ + let _lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; + admit_peer_join_across_nodes(local_endpoint, join_req, defer_sync_state_enable, apply_iam).await +} + +/// [`admit_peer_join`] minus the process-local lifecycle guard: the +/// distributed admission lock plus the fenced sequence under it. This is +/// exactly what a second node of this site runs concurrently — the lifecycle +/// guard cannot reach it — so the separate-nodes regression test drives this +/// function directly, and removing the admission lock breaks it. +async fn admit_peer_join_across_nodes( + local_endpoint: String, + join_req: SRPeerJoinReq, + defer_sync_state_enable: bool, + apply_iam: F, +) -> S3Result +where + F: FnOnce(SRPeerJoinReq) -> 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".to_string()))?; + with_config_object_write_lock(store, SITE_REPLICATION_JOIN_ADMISSION_LOCK_PATH.to_string(), move || async move { + let fresh = load_site_replication_state().await?; + let fresh_local_peer = local_peer_at_endpoint(local_endpoint.clone(), &fresh); + if join_request_is_superseded(&fresh, join_req.updated_at) { + let peer = fresh + .peers + .get(&fresh_local_peer.deployment_id) + .cloned() + .unwrap_or(fresh_local_peer); + return Ok(PeerJoinOutcome::Superseded(peer)); + } + + apply_iam(join_req.clone()).await?; + + let incoming_updated_at = join_req.updated_at; + update_site_replication_state_when_changed(move |state| { + let local_peer = local_peer_at_endpoint(local_endpoint, state); + if join_request_is_superseded(state, incoming_updated_at) { + let peer = state.peers.get(&local_peer.deployment_id).cloned().unwrap_or(local_peer); + return Ok(StateCommit::Unchanged(PeerJoinOutcome::Superseded(peer))); + } + apply_peer_join(state, &local_peer, join_req, defer_sync_state_enable); + Ok(StateCommit::Changed(PeerJoinOutcome::Applied(Box::new(state.clone()), local_peer))) + }) + .await + }) + .await + .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("lock site replication join admission failed: {e}")))? +} + +/// Upsert the replication service account a peer join carries. No-op when the +/// join brings no credentials. +async fn apply_peer_join_service_account(join_req: SRPeerJoinReq) -> S3Result<()> { + if join_req.svc_acct_access_key.is_empty() || join_req.svc_acct_secret_key.is_empty() { + return Ok(()); + } + let Some(iam_sys) = current_iam_handle() else { + return Err(s3_error!(InvalidRequest, "iam not init")); + }; + + if iam_sys.get_service_account(&join_req.svc_acct_access_key).await.is_ok() { + iam_sys + .update_service_account( + &join_req.svc_acct_access_key, + UpdateServiceAccountOpts { + session_policy: if join_req.svc_acct_access_key == SITE_REPLICATOR_SERVICE_ACCOUNT { + Some(site_replicator_service_account_policy()?) + } else { + None + }, + secret_key: Some(join_req.svc_acct_secret_key.clone()), + name: None, + description: None, + expiration: None, + status: None, + parent_user: None, + allow_site_replicator_account: join_req.svc_acct_access_key == SITE_REPLICATOR_SERVICE_ACCOUNT, + }, + ) + .await + .map_err(ApiError::from)?; + } else { + iam_sys + .new_service_account( + &join_req.svc_acct_parent, + None, + NewServiceAccountOpts { + session_policy: if join_req.svc_acct_access_key == SITE_REPLICATOR_SERVICE_ACCOUNT { + Some(site_replicator_service_account_policy()?) + } else { + None + }, + access_key: join_req.svc_acct_access_key.clone(), + secret_key: join_req.svc_acct_secret_key.clone(), + name: None, + description: None, + expiration: None, + allow_site_replicator_account: join_req.svc_acct_access_key == SITE_REPLICATOR_SERVICE_ACCOUNT, + claims: None, + }, + ) + .await + .map_err(ApiError::from)?; + } + Ok(()) +} + #[async_trait::async_trait] impl Operation for SRPeerJoinHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationAddAction).await?; let bootstrap_token = site_replication_bootstrap_token(&req.uri); - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - let local_peer = current_local_peer(&req, &state); + let local_endpoint = site_replication_local_endpoint(&req.uri, &req.headers); + // The body is fully read before the admission takes the lifecycle + // guard: a sender that stalls mid-body must not block this node's + // add/remove/rotate/reconciler. let join_envelope: SRPeerJoinEnvelope = read_site_replication_json(req, &cred.secret_key, true).await?; let defer_sync_state_enable = join_envelope.defer_sync_state_enable; let join_req = join_envelope.request; validate_join_peer_snapshot(&join_req.peers)?; - if let Some(current_updated_at) = state.updated_at { - let Some(incoming_updated_at) = join_req.updated_at else { + let committed = + admit_peer_join(local_endpoint, join_req, defer_sync_state_enable, apply_peer_join_service_account).await?; + // Committed; the reverse-reachability probe and the bucket back-fill + // run outside the transaction — their transport helpers' retry-event + // bookkeeping re-enters it (P1-15). + let (state, local_peer) = match committed { + PeerJoinOutcome::Applied(state, local_peer) => (*state, local_peer), + PeerJoinOutcome::Superseded(peer) => { return json_response(&SRPeerJoinResponse { - peer: state.peers.get(&local_peer.deployment_id).cloned().unwrap_or(local_peer), - ..Default::default() - }); - }; - if incoming_updated_at <= current_updated_at { - return json_response(&SRPeerJoinResponse { - peer: state.peers.get(&local_peer.deployment_id).cloned().unwrap_or(local_peer), + peer, ..Default::default() }); } - } - - if !join_req.svc_acct_access_key.is_empty() && !join_req.svc_acct_secret_key.is_empty() { - let Some(iam_sys) = current_iam_handle() else { - return Err(s3_error!(InvalidRequest, "iam not init")); - }; - - if iam_sys.get_service_account(&join_req.svc_acct_access_key).await.is_ok() { - iam_sys - .update_service_account( - &join_req.svc_acct_access_key, - UpdateServiceAccountOpts { - session_policy: if join_req.svc_acct_access_key == SITE_REPLICATOR_SERVICE_ACCOUNT { - Some(site_replicator_service_account_policy()?) - } else { - None - }, - secret_key: Some(join_req.svc_acct_secret_key.clone()), - name: None, - description: None, - expiration: None, - status: None, - parent_user: None, - allow_site_replicator_account: join_req.svc_acct_access_key == SITE_REPLICATOR_SERVICE_ACCOUNT, - }, - ) - .await - .map_err(ApiError::from)?; - } else { - iam_sys - .new_service_account( - &join_req.svc_acct_parent, - None, - NewServiceAccountOpts { - session_policy: if join_req.svc_acct_access_key == SITE_REPLICATOR_SERVICE_ACCOUNT { - Some(site_replicator_service_account_policy()?) - } else { - None - }, - access_key: join_req.svc_acct_access_key.clone(), - secret_key: join_req.svc_acct_secret_key.clone(), - name: None, - description: None, - expiration: None, - allow_site_replicator_account: join_req.svc_acct_access_key == SITE_REPLICATOR_SERVICE_ACCOUNT, - claims: None, - }, - ) - .await - .map_err(ApiError::from)?; - } - } - - state.service_account_access_key = join_req.svc_acct_access_key; - state.service_account_parent = join_req.svc_acct_parent; - state.updated_at = join_req.updated_at.or_else(|| Some(OffsetDateTime::now_utc())); - state.peers = normalize_join_peers_for_local(&local_peer, join_req.peers); - initialize_join_peer_sync_state(&mut state.peers, defer_sync_state_enable); - state.sync_state_initialized = true; - state.name = state - .peers - .get(&local_peer.deployment_id) - .map(|peer| peer.name.clone()) - .filter(|name| !name.is_empty()) - .unwrap_or_else(|| local_peer.name.clone()); - persist_site_replication_state(&state).await?; - // Committed; release the state lock before the reverse-reachability - // probe and bucket back-fill — their transport helpers' retry-event - // bookkeeping re-enters the state transaction (P1-15). - drop(_state_guard); + }; // Fix 1 (receiving side): ensure the joining peer also sets up replication for any // buckets it already owns so the reverse direction works from the start. Per-bucket // failures are logged (BUG2) so a reverse-direction back-fill gap is observable. @@ -9005,7 +9200,10 @@ impl Operation for SiteReplicationEditHandler { let ilm_expiry_override = sr_edit_ilm_expiry_override(&req.uri); let body = read_site_replication_body(req, &cred.secret_key, true).await?; let (mut incoming, tls_presence) = parse_public_peer_edit(&body)?; - let mut state_guard = Some(site_replication_state_process_guard().await); + // Planning snapshot: every commit below re-loads the state inside its + // transaction and re-checks the `updated_at` / pending-refresh + // precondition there, because the peer probes and fan-outs in between + // must not run under the state-object lock. let current_state = load_site_replication_state().await?; apply_public_peer_edit_tls_presence(¤t_state, &mut incoming, tls_presence); if !incoming.deployment_id.is_empty() || !incoming.endpoint.is_empty() || !incoming.name.is_empty() { @@ -9027,6 +9225,11 @@ impl Operation for SiteReplicationEditHandler { acked_deployment_ids: BTreeSet::new(), }) }); + // The precondition every commit below re-checks: the topology this + // edit was planned against, and the endpoint refresh it either + // continues or requires the absence of. + let expected_updated_at = current_state.updated_at; + let expected_pending_id = persisted_pending.as_ref().map(|pending| pending.id.clone()); let local_peer = current_local_runtime_peer(¤t_state); let existing_peer = existing_peer_for_edit(¤t_state, &incoming); let tls_capability_required = edit_peer_tls_capability_required(existing_peer, &incoming); @@ -9037,8 +9240,6 @@ impl Operation for SiteReplicationEditHandler { return Err(s3_error!(InvalidRequest, "site replication service account is not configured")); } let secret = site_replicator_service_account_secret(¤t_state.service_account_access_key).await?; - let expected_updated_at = current_state.updated_at; - drop(state_guard.take()); if tls_capability_required { require_edit_peer_tls_capability( ¤t_state, @@ -9052,39 +9253,30 @@ impl Operation for SiteReplicationEditHandler { if tls_transport_probe_required { probe_proposed_peer_tls_transport(&incoming, ¤t_state.service_account_access_key, &secret).await?; } - state_guard = Some(site_replication_state_process_guard().await); + // Early exit on a state that moved under the probe. Advisory only: + // the binding check is the CAS inside whichever commit follows. let latest_state = load_site_replication_state().await?; - if latest_state.updated_at != expected_updated_at - || pending_endpoint_refresh(&latest_state).as_ref().map(|pending| &pending.id) - != persisted_pending.as_ref().map(|pending| &pending.id) - { - return Err(s3_error!(InvalidRequest, "site replication state changed during capability probe")); - } + ensure_edit_precondition(&latest_state, expected_updated_at, expected_pending_id.as_ref(), "capability probe")?; service_account_secret_key = Some(secret); } - let mut state = if endpoint_refresh_requested { - current_state.clone() - } else { - edit_state(current_state.clone(), incoming.clone(), ilm_expiry_override) - }; - if endpoint_refresh_requested && current_state.service_account_access_key.is_empty() { return Err(s3_error!(InvalidRequest, "site replication service account is not configured")); } if current_state.service_account_access_key.is_empty() { - save_site_replication_state(&state).await?; + // No peers to notify: the edit is the whole operation, so it is + // computed and committed in one transaction. + let incoming = incoming.clone(); + update_site_replication_state(move |state| { + ensure_edit_precondition(state, expected_updated_at, expected_pending_id.as_ref(), "the edit")?; + *state = edit_state(std::mem::take(state), incoming, ilm_expiry_override); + Ok(()) + }) + .await?; } else { let service_account_secret_key = match service_account_secret_key { Some(secret) => secret, None => site_replicator_service_account_secret(¤t_state.service_account_access_key).await?, }; - let peers_to_send: Vec = if let Some(pending) = pending.as_ref() { - vec![pending.peer.clone()] - } else if ilm_expiry_override.is_some() { - state.peers.values().cloned().collect() - } else { - vec![normalize_peer_info(incoming)] - }; let routing_peers = pending .as_ref() .map(|pending| &pending.remote_peers) @@ -9096,8 +9288,6 @@ impl Operation for SiteReplicationEditHandler { let pending = pending.clone().ok_or_else(|| { S3Error::with_message(S3ErrorCode::InternalError, "endpoint refresh state is missing".to_string()) })?; - let expected_updated_at = current_state.updated_at; - drop(state_guard.take()); let probes = futures::future::join_all(remote_targets.iter().map(|target| { send_endpoint_refresh_admin_request_raw( target, @@ -9124,24 +9314,23 @@ impl Operation for SiteReplicationEditHandler { } } - state_guard = Some(site_replication_state_process_guard().await); - let latest_state = load_site_replication_state().await?; - if latest_state.updated_at != expected_updated_at - || pending_endpoint_refresh(&latest_state).as_ref().map(|pending| &pending.id) - != persisted_pending.as_ref().map(|pending| &pending.id) - { - return Err(s3_error!(InvalidRequest, "site replication state changed during capability probe")); - } - state = latest_state; let pending_id = pending.id.clone(); let refresh_request = EndpointRefreshRequest { id: pending.id.clone(), peer: pending.peer.clone(), }; - let pending = merge_pending_endpoint_refresh(&state, &pending, std::iter::empty::())?; - set_pending_endpoint_refresh(&mut state, pending.clone())?; - save_site_replication_state(&state).await?; - drop(state_guard.take()); + // Announce the pending refresh. The CAS sits in the same + // transaction as the write it guards, so a topology change + // that landed during the capability probes above cannot be + // overwritten by this snapshot. + let expected_pending_id = expected_pending_id.clone(); + let pending = update_site_replication_state(move |state| { + ensure_edit_precondition(state, expected_updated_at, expected_pending_id.as_ref(), "capability probe")?; + let pending = merge_pending_endpoint_refresh(state, &pending, std::iter::empty::())?; + set_pending_endpoint_refresh(state, pending.clone())?; + Ok(pending) + }) + .await?; let responses = futures::future::join_all(remote_targets.iter().map(|target| async { if legacy_deployment_ids.contains(&target.deployment_id) { refresh_legacy_peer_bucket_targets( @@ -9177,51 +9366,63 @@ impl Operation for SiteReplicationEditHandler { } } - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - let Some(pending) = pending_endpoint_refresh(&state).filter(|pending| pending.id == pending_id) else { - return Err(s3_error!(InvalidRequest, "endpoint target refresh state changed during update")); - }; - let pending = merge_pending_endpoint_refresh(&state, &pending, acked_deployment_ids)?; - set_pending_endpoint_refresh(&mut state, pending)?; - save_site_replication_state(&state).await?; + let acked_pending_id = pending_id.clone(); + let service_account_access_key = update_site_replication_state(move |state| { + let Some(pending) = pending_endpoint_refresh(state).filter(|pending| pending.id == acked_pending_id) else { + return Err(s3_error!(InvalidRequest, "endpoint target refresh state changed during update")); + }; + let pending = merge_pending_endpoint_refresh(state, &pending, acked_deployment_ids)?; + set_pending_endpoint_refresh(state, pending)?; + Ok(state.service_account_access_key.clone()) + }) + .await?; if let Some(err) = refresh_error { return Err(err); } - let service_account_secret_key = - site_replicator_service_account_secret(&state.service_account_access_key).await?; - drop(_state_guard); + let service_account_secret_key = site_replicator_service_account_secret(&service_account_access_key).await?; refresh_bucket_targets_after_endpoint_edit(&pending_id, &service_account_secret_key).await?; - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - let Some(pending) = pending_endpoint_refresh(&state).filter(|pending| pending.id == pending_id) else { - return Err(s3_error!(InvalidRequest, "endpoint target refresh state changed during update")); - }; - state = edit_state(state, pending.peer, ilm_expiry_override); - clear_pending_endpoint_refresh(&mut state); - save_site_replication_state(&state).await?; + update_site_replication_state(move |state| { + let Some(pending) = pending_endpoint_refresh(state).filter(|pending| pending.id == pending_id) else { + return Err(s3_error!(InvalidRequest, "endpoint target refresh state changed during update")); + }; + *state = edit_state(std::mem::take(state), pending.peer, ilm_expiry_override); + clear_pending_endpoint_refresh(state); + Ok(()) + }) + .await?; } else { // Commit before the peer fan-out (mirrors the add/join // handlers): a failed notification is recorded as a retry // event and converges from the committed local state — // fanning out first meant the retry event pointed at a state - // the local site had not saved. The generation is allocated in - // the same commit, i.e. under the state-object lock, so it - // orders this edit against one another node accepts - // concurrently — the process guard below cannot. - let edit_generation = next_peer_edit_generation(&mut state); - save_site_replication_state(&state).await?; + // the local site had not saved. The edit itself is applied to + // the state the transaction loads, under the CAS, so a + // topology change that slipped past the planning snapshot + // fails the edit instead of being overwritten by it. The + // generation is allocated in that same commit, i.e. under the + // state-object lock, so it orders this edit against one + // another node of this site accepts concurrently. + let incoming = incoming.clone(); + let (edit_generation, peers_to_send) = update_site_replication_state(move |state| { + ensure_edit_precondition(state, expected_updated_at, expected_pending_id.as_ref(), "the edit")?; + *state = edit_state(std::mem::take(state), incoming.clone(), ilm_expiry_override); + let peers_to_send: Vec = if ilm_expiry_override.is_some() { + state.peers.values().cloned().collect() + } else { + vec![normalize_peer_info(incoming)] + }; + Ok((next_peer_edit_generation(state), peers_to_send)) + }) + .await?; let edit_path = peer_edit_path_with_fence(local_deployment_id.as_deref(), edit_generation); let delivery_fence = local_deployment_id.is_some().then_some(edit_generation); - // The fan-out stays UNDER the state guard so deliveries issued - // by THIS node keep their commit order; the generation fence - // above is what covers the cross-node case, where the peer - // rejects a delivery an older edit is still trying to make. - // The retry-event bookkeeping is what cannot run under the - // guard — it re-enters the state transaction (P1-15) — so - // deliver with the plain transport here and settle the retry - // queue after the guard is released. + // The fan-out runs outside the transaction — peer traffic + // under the state-object lock would stall every writer of this + // site, and the retry bookkeeping below re-enters it (P1-15). + // Ordering is the generation fence's job: a delivery this + // fan-out is still retrying is rejected by the receiver once a + // newer generation from this site has landed there. let mut delivered: Vec = Vec::new(); let mut failure: Option<(PeerInfo, S3Error)> = None; 'fanout: for target in remote_targets { @@ -9243,7 +9444,6 @@ impl Operation for SiteReplicationEditHandler { } delivered.push(target.clone()); } - drop(state_guard.take()); // Settle only what this generation is entitled to: a newer // edit that committed and failed its own delivery while this @@ -9293,6 +9493,20 @@ impl Operation for SRPeerEditCapabilitiesHandler { pub struct SRPeerEditHandler {} +/// What the peer-edit transaction decided about an incoming delivery. The +/// checks and the write share one transaction, so the verdict has to travel +/// out of the closure instead of being answered where it is taken. +enum PeerEditOutcome { + /// Applied; carries the service account access key the follow-up + /// endpoint-refresh work needs from the committed state. + Applied(String), + /// Nothing to do — a superseded delivery or one this site already + /// committed. Answered as success so the sender stops retrying. + Acked, + /// Refused, with the detail the sender is told. + Rejected(&'static str), +} + #[async_trait::async_trait] impl Operation for SRPeerEditHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { @@ -9300,101 +9514,107 @@ impl Operation for SRPeerEditHandler { let queries = query_pairs(&req.uri); let ilm_expiry_override = sr_edit_ilm_expiry_override(&req.uri); let endpoint_refresh_requested = queries.get("refresh-targets").is_some_and(|value| value == "true"); - let delivery_fence = peer_edit_fence(&queries); - let state_guard = site_replication_state_process_guard().await; - let state = load_site_replication_state().await?; - // 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 roll the - // peer back to the older edit. Ack it — the newer edit already landed, - // so the sender has nothing to retry. - if let Some((origin, generation)) = delivery_fence.as_ref() - && peer_edit_delivery_is_stale(&state, origin, *generation) - { - return json_response(&ReplicateEditStatus { - success: true, - status: SITE_REPL_EDIT_SUCCESS.to_string(), - api_version: Some(SITE_REPL_API_VERSION.to_string()), - ..Default::default() - }); - } - if endpoint_refresh_requested && (state.pending_rotation.is_some() || state.pending_remove.is_some()) { - return json_response(&ReplicateEditStatus { - success: false, - status: SITE_REPL_EDIT_SUCCESS.to_string(), - err_detail: "another site replication operation is pending".to_string(), - api_version: Some(SITE_REPL_API_VERSION.to_string()), - }); - } - let local_peer = current_local_peer(&req, &state); - let (refresh_id, mut incoming) = if endpoint_refresh_requested { + let commit_fence = peer_edit_fence(&queries); + let local_endpoint = site_replication_local_endpoint(&req.uri, &req.headers); + let (refresh_id, incoming) = if endpoint_refresh_requested { let refresh: EndpointRefreshRequest = read_site_replication_json(req, "", false).await?; (Some(refresh.id), refresh.peer) } else { (None, read_site_replication_json(req, "", false).await?) }; - if same_identity_endpoint(&incoming.endpoint, &local_peer.endpoint) { - incoming.deployment_id = local_peer.deployment_id.clone(); - if incoming.name.is_empty() { - incoming.name = local_peer.name.clone(); - } - } - align_peer_edit_deployment_id(&state, &mut incoming); - if endpoint_refresh_requested - && pending_endpoint_refresh(&state).is_some_and(|pending| refresh_id.as_deref() != Some(&pending.id)) - { - return json_response(&ReplicateEditStatus { - success: false, - status: SITE_REPL_EDIT_SUCCESS.to_string(), - err_detail: "another endpoint target refresh is pending".to_string(), - api_version: Some(SITE_REPL_API_VERSION.to_string()), - }); - } - if endpoint_refresh_requested - && (refresh_id.as_ref().is_none_or(String::is_empty) || !peer_endpoint_edit_requested(&state, &incoming)) - { - return json_response(&ReplicateEditStatus { - success: false, - status: SITE_REPL_EDIT_SUCCESS.to_string(), - err_detail: "peer endpoint was not found".to_string(), - api_version: Some(SITE_REPL_API_VERSION.to_string()), - }); - } - if endpoint_refresh_requested && internal_endpoint_refresh_already_committed(&state, &incoming) { - return json_response(&ReplicateEditStatus { - success: true, - status: SITE_REPL_EDIT_SUCCESS.to_string(), - api_version: Some(SITE_REPL_API_VERSION.to_string()), - ..Default::default() - }); - } - let mut state = if endpoint_refresh_requested { - validate_proposed_peer(&incoming)?; - state - } else { - apply_internal_peer_edit(state, &local_peer, incoming.clone(), ilm_expiry_override)? + // Everything the delivery is checked against — the fence, the pending + // operations, the peer it names — is read inside the transaction that + // applies it. Checking against a state loaded before the lock would + // let the check pass on one snapshot and the write land on another. + let commit_endpoint = local_endpoint.clone(); + let commit_refresh_id = refresh_id.clone(); + let outcome = update_site_replication_state_when_changed(move |state| { + let mut incoming = incoming; + let local_peer = local_peer_at_endpoint(commit_endpoint, state); + // 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 + // roll the peer back to the older edit. Ack it — the newer edit + // already landed, so the sender has nothing to retry. + if let Some((origin, generation)) = commit_fence.as_ref() + && peer_edit_delivery_is_stale(state, origin, *generation) + { + return Ok(StateCommit::Unchanged(PeerEditOutcome::Acked)); + } + if endpoint_refresh_requested && (state.pending_rotation.is_some() || state.pending_remove.is_some()) { + return Ok(StateCommit::Unchanged(PeerEditOutcome::Rejected( + "another site replication operation is pending", + ))); + } + if same_identity_endpoint(&incoming.endpoint, &local_peer.endpoint) { + incoming.deployment_id = local_peer.deployment_id.clone(); + if incoming.name.is_empty() { + incoming.name = local_peer.name.clone(); + } + } + align_peer_edit_deployment_id(state, &mut incoming); + if endpoint_refresh_requested + && pending_endpoint_refresh(state).is_some_and(|pending| commit_refresh_id.as_deref() != Some(&pending.id)) + { + return Ok(StateCommit::Unchanged(PeerEditOutcome::Rejected( + "another endpoint target refresh is pending", + ))); + } + if endpoint_refresh_requested + && (commit_refresh_id.as_ref().is_none_or(String::is_empty) || !peer_endpoint_edit_requested(state, &incoming)) + { + return Ok(StateCommit::Unchanged(PeerEditOutcome::Rejected("peer endpoint was not found"))); + } + if endpoint_refresh_requested && internal_endpoint_refresh_already_committed(state, &incoming) { + return Ok(StateCommit::Unchanged(PeerEditOutcome::Acked)); + } + + if endpoint_refresh_requested { + validate_proposed_peer(&incoming)?; + set_pending_endpoint_refresh( + state, + PendingEndpointRefresh { + id: commit_refresh_id.unwrap_or_default(), + peer: incoming, + remote_peers: BTreeMap::new(), + acked_deployment_ids: BTreeSet::new(), + }, + )?; + } else { + *state = apply_internal_peer_edit(std::mem::take(state), &local_peer, incoming, ilm_expiry_override)?; + } + // Raise the origin's high-water mark in the same commit as the + // edit it fences: a crash between the two would let the superseded + // delivery apply on the next attempt. + if let Some((origin, generation)) = commit_fence.as_ref() { + record_applied_peer_edit_generation(state, origin, *generation); + } + Ok(StateCommit::Changed(PeerEditOutcome::Applied(state.service_account_access_key.clone()))) + }) + .await?; + + let service_account_access_key = match outcome { + PeerEditOutcome::Applied(service_account_access_key) => service_account_access_key, + PeerEditOutcome::Acked => { + return json_response(&ReplicateEditStatus { + success: true, + status: SITE_REPL_EDIT_SUCCESS.to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + ..Default::default() + }); + } + PeerEditOutcome::Rejected(err_detail) => { + return json_response(&ReplicateEditStatus { + success: false, + status: SITE_REPL_EDIT_SUCCESS.to_string(), + err_detail: err_detail.to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + }); + } }; if endpoint_refresh_requested { - set_pending_endpoint_refresh( - &mut state, - PendingEndpointRefresh { - id: refresh_id.clone().unwrap_or_default(), - peer: incoming.clone(), - remote_peers: BTreeMap::new(), - acked_deployment_ids: BTreeSet::new(), - }, - )?; - } - // Raise the origin's high-water mark in the same commit as the edit it - // fences: a crash between the two would let the superseded delivery - // apply on the next attempt. - if let Some((origin, generation)) = delivery_fence.as_ref() { - record_applied_peer_edit_generation(&mut state, origin, *generation); - } - save_site_replication_state(&state).await?; - if endpoint_refresh_requested { - if state.service_account_access_key.is_empty() { + if service_account_access_key.is_empty() { return json_response(&ReplicateEditStatus { success: false, status: SITE_REPL_EDIT_SUCCESS.to_string(), @@ -9402,23 +9622,29 @@ impl Operation for SRPeerEditHandler { api_version: Some(SITE_REPL_API_VERSION.to_string()), }); } - let service_account_secret_key = site_replicator_service_account_secret(&state.service_account_access_key).await?; + let service_account_secret_key = site_replicator_service_account_secret(&service_account_access_key).await?; let pending_id = refresh_id.unwrap_or_default(); - drop(state_guard); + // The bucket-target rewrite talks to the store for every bucket; + // it runs between the two transactions, never inside one. refresh_bucket_targets_after_endpoint_edit(&pending_id, &service_account_secret_key).await?; - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; - let Some(pending) = pending_endpoint_refresh(&state).filter(|pending| pending.id == pending_id) else { + let committed = update_site_replication_state_when_changed(move |state| { + let local_peer = local_peer_at_endpoint(local_endpoint, state); + let Some(pending) = pending_endpoint_refresh(state).filter(|pending| pending.id == pending_id) else { + return Ok(StateCommit::Unchanged(false)); + }; + *state = apply_internal_peer_edit(std::mem::take(state), &local_peer, pending.peer, ilm_expiry_override)?; + clear_pending_endpoint_refresh(state); + Ok(StateCommit::Changed(true)) + }) + .await?; + if !committed { return json_response(&ReplicateEditStatus { success: false, status: SITE_REPL_EDIT_SUCCESS.to_string(), err_detail: "endpoint target refresh state changed during update".to_string(), api_version: Some(SITE_REPL_API_VERSION.to_string()), }); - }; - state = apply_internal_peer_edit(state, &local_peer, pending.peer, ilm_expiry_override)?; - clear_pending_endpoint_refresh(&mut state); - save_site_replication_state(&state).await?; + } return json_response(&ReplicateEditStatus { success: true, status: SITE_REPL_EDIT_SUCCESS.to_string(), @@ -9439,19 +9665,19 @@ impl Operation for SRPeerRemoveHandler { let remove_req: SRRemoveReq = read_site_replication_json(req, "", false).await?; let _lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.write().await; - let _state_guard = site_replication_state_process_guard().await; - 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")); - } - if current_state.pending_rotation.is_some() { - return Err(s3_error!(InvalidRequest, "service account rotation is pending")); - } + let removed_deployment_ids = update_site_replication_state(move |state| { + if pending_endpoint_refresh(state).is_some() { + return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); + } + if state.pending_rotation.is_some() { + return Err(s3_error!(InvalidRequest, "service account rotation is pending")); + } - let removed_deployment_ids = removed_deployment_ids_for_remove_req(¤t_state, &remove_req); - - let state = remove_sites(current_state, remove_req); - persist_site_replication_state(&state).await?; + let removed_deployment_ids = removed_deployment_ids_for_remove_req(state, &remove_req); + *state = remove_sites(std::mem::take(state), remove_req); + Ok(removed_deployment_ids) + }) + .await?; // Clean up bucket targets and replication rules that referenced removed peers. if !removed_deployment_ids.is_empty() @@ -9483,7 +9709,6 @@ impl Operation for SiteReplicationResyncOpHandler { let requested_peer: PeerInfo = read_site_replication_json(req, "", false).await?; let _lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; let (peer, existing_status) = { - let _state_guard = site_replication_state_process_guard().await; let state = load_site_replication_state().await?; let local_peer = current_local_runtime_peer(&state); let requested_peer = normalize_peer_info(requested_peer); @@ -9644,9 +9869,11 @@ impl Operation for SRStateEditHandler { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?; reject_site_replicator_on_public_admin(&cred)?; let body: SRStateEditReq = read_site_replication_json(req, "", false).await?; - let _state_guard = site_replication_state_process_guard().await; - let state = apply_state_edit_req(load_site_replication_state().await?, body); - save_site_replication_state(&state).await?; + update_site_replication_state(move |state| { + *state = apply_state_edit_req(std::mem::take(state), body); + Ok(()) + }) + .await?; Ok(empty_response(StatusCode::OK)) } } @@ -9658,15 +9885,11 @@ impl Operation for SiteReplicationRepairHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?; reject_site_replicator_on_public_admin(&cred)?; - let (state, local_peer) = { - let _state_guard = site_replication_state_process_guard().await; - let state = load_site_replication_state().await?; - if !state.enabled() || state.service_account_access_key.is_empty() { - return Err(s3_error!(InvalidRequest, "site replication is not configured")); - } - let local_peer = current_local_peer(&req, &state); - (state, local_peer) - }; + let state = load_site_replication_state().await?; + if !state.enabled() || state.service_account_access_key.is_empty() { + return Err(s3_error!(InvalidRequest, "site replication is not configured")); + } + let local_peer = current_local_peer(&req, &state); let body: SiteReplicationRepairRequest = read_site_replication_json(req, "", false).await?; let info = build_sr_info(&state, &local_peer).await?; let plan = site_replication_bootstrap_plan(&info)?; @@ -9769,44 +9992,55 @@ impl Operation for SRRotateServiceAccountHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?; reject_site_replicator_on_public_admin(&cred)?; - let (pending_rotation, local_peer, previous_access_key) = { - let _state_guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await?; + // The lifecycle guard is what keeps the rotation's IAM writes and the + // background service-account reconciler apart: the reconciler runs + // its whole repair under a lifecycle try-acquire, and its + // pending-rotation precheck is only sound if a rotation cannot start + // mid-repair and race its own IAM write against the reconciler's + // stale one. (The removed process mutex used to provide this + // exclusion as a side effect.) + let _lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; + let local_endpoint = site_replication_local_endpoint(&req.uri, &req.headers); + let rotation_parent = cred.access_key.clone(); + let (pending_rotation, local_peer, previous_access_key) = update_site_replication_state_when_changed(move |state| { if !state.enabled() { return Err(s3_error!(InvalidRequest, "site replication is not configured")); } - if pending_endpoint_refresh(&state).is_some() { + if pending_endpoint_refresh(state).is_some() { return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); } if state.pending_remove.is_some() { return Err(s3_error!(InvalidRequest, "site replication remove is pending")); } - let local_peer = current_local_peer(&req, &state); + let local_peer = local_peer_at_endpoint(local_endpoint, state); let previous_access_key = state.service_account_access_key.clone(); + // Resuming a rotation another attempt already recorded must + // not rewrite the state: the pending record is the contract + // the peers were told about. if let Some(pending) = state.pending_rotation.clone() { - (pending, local_peer, previous_access_key) - } else { - let new_secret_key = rustfs_credentials::gen_secret_key(40) - .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("generate secret key failed: {e}")))?; - state.service_account_access_key = SITE_REPLICATOR_SERVICE_ACCOUNT.to_string(); - state.service_account_parent = cred.access_key.clone(); - state.updated_at = Some(OffsetDateTime::now_utc()); - let pending = PendingRotation { - id: Uuid::new_v4().to_string(), - access_key: SITE_REPLICATOR_SERVICE_ACCOUNT.to_string(), - parent: cred.access_key.clone(), - new_secret_key, - secret_candidates: legacy_site_replicator_state_secret(&state).into_iter().collect(), - peers: state.peers.clone(), - acked_deployment_ids: BTreeSet::new(), - updated_at: state.updated_at, - }; - state.pending_rotation = Some(pending.clone()); - save_site_replication_state(&state).await?; - (pending, local_peer, previous_access_key) + return Ok(StateCommit::Unchanged((pending, local_peer, previous_access_key))); } - }; + + let new_secret_key = rustfs_credentials::gen_secret_key(40) + .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("generate secret key failed: {e}")))?; + state.service_account_access_key = SITE_REPLICATOR_SERVICE_ACCOUNT.to_string(); + state.service_account_parent = rotation_parent.clone(); + state.updated_at = Some(OffsetDateTime::now_utc()); + let pending = PendingRotation { + id: Uuid::new_v4().to_string(), + access_key: SITE_REPLICATOR_SERVICE_ACCOUNT.to_string(), + parent: rotation_parent, + new_secret_key, + secret_candidates: legacy_site_replicator_state_secret(state).into_iter().collect(), + peers: state.peers.clone(), + acked_deployment_ids: BTreeSet::new(), + updated_at: state.updated_at, + }; + state.pending_rotation = Some(pending.clone()); + Ok(StateCommit::Changed((pending, local_peer, previous_access_key))) + }) + .await?; if !previous_access_key.is_empty() && let Ok(previous_iam_secret) = site_replicator_service_account_secret(&previous_access_key).await @@ -11356,6 +11590,35 @@ mod tests { assert!(!peer_capability_response_supported(&remote, StatusCode::NOT_FOUND, b"").expect("legacy peer")); } + /// P1-15 PR2: the add's finalize fan-out delivers peer edits after its + /// state transaction has been released — nothing may hold the state-object + /// lock across peer traffic. Ordering therefore rests entirely on the + /// generation allocated in that commit: an unstamped delivery is applied + /// by the receiver in arrival order, which is what the removed process + /// guard used to paper over (and never could across two nodes). + #[test] + fn add_handler_fans_out_peer_edits_under_the_committed_generation() { + 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"); + + assert!( + add.contains("next_peer_edit_generation"), + "the add must allocate a fan-out generation (inside the committing transaction)" + ); + assert!( + add.contains("peer_edit_path_with_fence"), + "the add's finalize fan-out must carry the committed generation fence" + ); + assert!( + !add.contains("SITE_REPLICATION_PEER_EDIT_PATH"), + "the finalize fan-out must not fall back to the unstamped peer-edit path" + ); + } + #[test] fn test_tls_capability_gates_run_before_add_or_edit_state_side_effects() { let src = include_str!("site_replication.rs"); @@ -11384,8 +11647,8 @@ mod tests { ); assert!( edit.find("require_edit_peer_tls_capability").expect("edit capability gate") - < edit.find("save_site_replication_state").expect("state save"), - "edit capability gate must run before state is saved" + < edit.find("update_site_replication_state(").expect("state commit"), + "edit capability gate must run before the state is committed" ); } @@ -11643,13 +11906,25 @@ mod tests { // rejects against — dropping either half silently restores // last-writer-wins between two nodes of the sending site. assert!( - handler_block.contains("peer_edit_delivery_is_stale(&state, origin, *generation)"), + handler_block.contains("peer_edit_delivery_is_stale(state, origin, *generation)"), "SRPeerEditHandler must reject peer edits a newer generation already superseded" ); assert!( - handler_block.contains("record_applied_peer_edit_generation(&mut state, origin, *generation);"), + handler_block.contains("record_applied_peer_edit_generation(state, origin, *generation);"), "SRPeerEditHandler must record the applied generation so later stale deliveries are recognised" ); + // 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 + // another — which is the interleaving the fence exists to reject. + assert!( + handler_block.contains("update_site_replication_state_when_changed(move |state| {"), + "SRPeerEditHandler must take the fence decision inside the state transaction" + ); + assert!( + !handler_block.contains("save_site_replication_state("), + "SRPeerEditHandler must not write the state outside the transaction boundary" + ); let sender_block = src .split("impl Operation for SiteReplicationEditHandler") @@ -11657,7 +11932,7 @@ mod tests { .and_then(|rest| rest.split("pub struct SRPeerEditCapabilitiesHandler").next()) .expect("SiteReplicationEditHandler block should exist"); assert!( - sender_block.contains("let edit_generation = next_peer_edit_generation(&mut state);"), + sender_block.contains("Ok((next_peer_edit_generation(state), peers_to_send))"), "the edit handler must allocate the generation inside the committed state, not outside the lock" ); } @@ -11936,14 +12211,12 @@ mod tests { let lifecycle = SiteReplicationLifecycleGuard::acquire().await; let add_guard = SiteReplicationAddInProgressGuard::start(lifecycle, HashSet::new()).expect("start site replication add guard"); - let add_state = site_replication_state_process_guard().await; let (started_tx, started_rx) = tokio::sync::oneshot::channel(); let (entered_tx, mut entered_rx) = tokio::sync::oneshot::channel(); let remove = tokio::spawn(async move { let _ = started_tx.send(()); let _lifecycle = SiteReplicationLifecycleGuard::acquire().await; let _bucket_op = SITE_REPLICATION_BUCKET_OP_LOCK.write().await; - let _state = site_replication_state_process_guard().await; let _ = entered_tx.send(()); }); started_rx.await.expect("remove task started"); @@ -11954,7 +12227,6 @@ mod tests { assert!(matches!(entered_rx.try_recv(), Err(tokio::sync::oneshot::error::TryRecvError::Empty))); drop(callback); - drop(add_state); drop(add_guard); tokio::time::timeout(Duration::from_millis(500), remove) .await @@ -12785,11 +13057,10 @@ mod tests { assert_eq!(dequeue_site_replication_retry_events(&mut queue, &peer, iam_path), 1); } - /// P1-15 review follow-up: the receiving side of the ordering fence. The - /// sender's process mutex is per node, so two nodes of the same site can - /// fan out in the opposite order to their commits; the receiver decides - /// ordering from the generation the sender allocated under the distributed - /// state lock. + /// P1-15 review follow-up: the receiving side of the ordering fence. Two + /// nodes of the sending site can fan out in the opposite order to their + /// commits; the receiver decides ordering from the generation the sender + /// allocated under the distributed state lock. #[test] fn peer_edit_fence_rejects_a_delivery_the_newer_edit_already_passed() { let mut state = SiteReplicationState::default(); @@ -12821,6 +13092,62 @@ mod tests { assert!(peer_edit_fence(&HashMap::new()).is_none()); } + /// P1-15 PR2: an accepted join must PRESERVE the peer-edit high-water + /// marks of peers that stayed. Join fan-outs are routine (every add and + /// every service-account rotation delivers `SRPeerJoin` to existing + /// peers), so a blanket reset here would let any stalled older edit from + /// a peer that never left land after the join and roll its record back — + /// exactly the interleaving the fence exists to reject. + #[test] + fn peer_join_preserves_live_edit_generation_marks() { + let local = PeerInfo { + deployment_id: "site-b".to_string(), + ..peer("site-b", "https://site-b.example.com") + }; + let remote = PeerInfo { + deployment_id: "site-a".to_string(), + ..peer("site-a", "https://site-a.example.com") + }; + // `remote` is already a peer and has delivered edits up to generation + // 12; the incoming join (say, a rotation fan-out) keeps both sites. + let mut state = SiteReplicationState { + peers: BTreeMap::from([ + (local.deployment_id.clone(), local.clone()), + (remote.deployment_id.clone(), remote.clone()), + ]), + applied_edit_generations: BTreeMap::from([(remote.deployment_id.clone(), 12)]), + ..Default::default() + }; + + apply_peer_join( + &mut state, + &local, + SRPeerJoinReq { + svc_acct_access_key: SITE_REPLICATOR_SERVICE_ACCOUNT.to_string(), + svc_acct_secret_key: "svc-secret".to_string(), + svc_acct_parent: "root".to_string(), + peers: BTreeMap::from([ + (local.deployment_id.clone(), local.clone()), + (remote.deployment_id.clone(), remote.clone()), + ]), + updated_at: Some(OffsetDateTime::now_utc()), + }, + true, + ); + + assert_eq!( + state.applied_edit_generations.get(&remote.deployment_id), + Some(&12), + "the join must keep the live mark for a peer that stayed: {:?}", + state.applied_edit_generations + ); + assert!( + peer_edit_delivery_is_stale(&state, &remote.deployment_id, 11), + "a stalled pre-join delivery must still be fenced out after the join" + ); + assert_eq!(state.peers.len(), 2, "the join snapshot replaces the local topology"); + } + /// One edit fans out one delivery per peer record under a single /// generation (the ILM-expiry edit sends every peer's record). The /// receiver's fenced sequence — staleness check, apply, raise the @@ -15283,71 +15610,353 @@ mod tests { assert!(parse_site_resync_page(&query, &newer).is_err()); } - /// P1-15 review follow-up: isolates the PROCESS guard. One writer uses - /// the legacy shape that the not-yet-migrated call sites still use (take - /// the process mutex, then load / mutate / save through the plain - /// helpers, which take their own per-IO object locks); the other runs the - /// full transaction. They only stay serialized because the transaction - /// also takes the process mutex — drop it there and this test loses one - /// of the two updates. - #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + /// P1-15 PR2: the persist-or-skip side of the transaction. A miss + /// (`StateCommit::Unchanged`) must not write at all: the shared persist + /// helper clears the whole object once a state has ≤1 peer and no pending + /// rotation/removal, so a no-op ack or clear that "harmlessly" persisted + /// would delete the retry queue and every other field along with it. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[serial] - async fn test_transaction_serializes_against_a_process_only_legacy_writer() { + async fn test_missed_pending_clear_must_not_rewrite_the_state_object() { publish_ready_iam_context().await; + // One peer, no pending records: exactly the shape the persist helper's clear + // branch fires on. Only the test-only seeder can write it. let seed = SiteReplicationState { + peers: BTreeMap::from([( + "site-solo".to_string(), + PeerInfo { + deployment_id: "site-solo".to_string(), + ..peer("site-solo", "https://solo.example:9000") + }, + )]), + retry_queue: vec![SiteReplicationRetryEvent { + id: "evt-1".to_string(), + peer_deployment_id: "site-gone".to_string(), + peer_endpoint: "https://gone.example:9000".to_string(), + path: "/rustfs/admin/v3/site-replication/peer/iam-item".to_string(), + retry_count: 2, + failed: false, + last_error: "peer offline".to_string(), + updated_at: Some(OffsetDateTime::now_utc()), + edit_generation: None, + }], + ..Default::default() + }; + save_site_replication_state(&seed).await.expect("seed state"); + + clear_pending_remove("no-such-remove").await.expect("no-op clear"); + mark_pending_rotation_peer_acked("no-such-rotation", "site-x") + .await + .expect("no-op ack"); + record_pending_remove_secret_candidate("no-such-remove", "secret".to_string()) + .await + .expect("no-op candidate"); + + let reloaded = load_site_replication_state().await.expect("reload"); + assert_eq!( + reloaded.retry_queue.len(), + 1, + "a missed pending lookup persisted (and therefore cleared) the state object" + ); + assert_eq!(reloaded.peers.len(), 1, "the peer record must survive the no-op calls"); + } + + /// Review follow-up on P1-15 PR2 (overtrue): two joins accepted by the + /// same node must not interleave their IAM writes with each other's + /// commits. Join A loads a stale snapshot and pauses before its IAM + /// write; join B (newer) applies secret B and commits; A resumes, its + /// IAM write would overwrite secret B, and its commit is then refused as + /// superseded — the persisted state advertises B's contract while IAM + /// holds A's secret. `admit_peer_join` closes this by serializing the + /// whole admission under the lifecycle guard and re-checking staleness + /// BEFORE the IAM step: with the guard, B cannot even start while A is + /// gated mid-IAM. Remove the guard (or move the IAM step ahead of the + /// fresh staleness check) and this test deadlocks or records B's IAM + /// write before A finishes. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[serial] + async fn test_peer_join_admission_serializes_iam_apply_against_a_newer_join() { + publish_ready_iam_context().await; + + // Whole-second timestamps so the RFC3339 round trip through the state + // object cannot lose sub-second precision under the equality asserts. + let now = OffsetDateTime::now_utc().replace_nanosecond(0).expect("truncate nanos"); + let local = PeerInfo { + deployment_id: "site-local".to_string(), + ..peer("site-local", "https://local.example:9000") + }; + let remote = PeerInfo { + deployment_id: "site-remote".to_string(), + ..peer("site-remote", "https://remote.example:9000") + }; + let seed = SiteReplicationState { + peers: BTreeMap::from([ + (local.deployment_id.clone(), local.clone()), + (remote.deployment_id.clone(), remote.clone()), + ]), + updated_at: Some(now - Duration::from_secs(60)), + ..Default::default() + }; + save_site_replication_state(&seed).await.expect("seed state"); + + let join_peers = BTreeMap::from([ + (local.deployment_id.clone(), local.clone()), + (remote.deployment_id.clone(), remote.clone()), + ]); + let join_req = |updated_at: OffsetDateTime, secret: &str| SRPeerJoinReq { + svc_acct_access_key: "svc-join".to_string(), + svc_acct_secret_key: secret.to_string(), + svc_acct_parent: "root".to_string(), + peers: join_peers.clone(), + updated_at: Some(updated_at), + }; + + let iam_log: Arc>> = Arc::new(StdMutex::new(Vec::new())); + let (a_entered_tx, a_entered_rx) = tokio::sync::oneshot::channel(); + let (a_gate_tx, a_gate_rx) = tokio::sync::oneshot::channel::<()>(); + + // Join A (older, T1): pauses inside its IAM step. + let log_a = iam_log.clone(); + let endpoint_a = "https://local.example:9000".to_string(); + let req_a = join_req(now - Duration::from_secs(30), "secret-a"); + let join_a = tokio::spawn(async move { + admit_peer_join(endpoint_a, req_a, true, move |_req| async move { + let _ = a_entered_tx.send(()); + let _ = a_gate_rx.await; + log_a.lock().expect("iam log").push("iam-a"); + Ok(()) + }) + .await + }); + a_entered_rx.await.expect("join A reached its IAM step"); + + // Join B (newer, T2) arrives while A is gated mid-IAM. The lifecycle + // guard must hold it at the door. + let log_b = iam_log.clone(); + let endpoint_b = "https://local.example:9000".to_string(); + let req_b = join_req(now, "secret-b"); + let join_b = tokio::spawn(async move { + admit_peer_join(endpoint_b, req_b, true, move |_req| async move { + log_b.lock().expect("iam log").push("iam-b"); + Ok(()) + }) + .await + }); + tokio::time::sleep(Duration::from_millis(200)).await; + assert!( + iam_log.lock().expect("iam log").is_empty(), + "join B ran its IAM step while join A was still mid-admission: {:?}", + iam_log.lock().expect("iam log") + ); + + a_gate_tx.send(()).expect("release join A"); + let outcome_a = join_a.await.expect("join A task").expect("join A admission"); + let outcome_b = join_b.await.expect("join B task").expect("join B admission"); + assert!(matches!(outcome_a, PeerJoinOutcome::Applied(..)), "join A must commit first"); + assert!( + matches!(outcome_b, PeerJoinOutcome::Applied(..)), + "the newer join B must still apply after A" + ); + assert_eq!( + *iam_log.lock().expect("iam log"), + vec!["iam-a", "iam-b"], + "IAM writes must land in admission order, ending on the committed join's secret" + ); + assert_eq!( + load_site_replication_state().await.expect("reload").updated_at, + Some(now), + "the persisted state must end on join B, matching the last IAM write" + ); + } + + /// Review follow-up on P1-15 PR2 (overtrue, round 2): the same + /// interleaving driven by two SEPARATE NODES, which the process-local + /// lifecycle guard cannot reach. Both admissions run + /// `admit_peer_join_across_nodes` — the production path minus the + /// process-local guard, exactly what a second node executes — so only + /// the distributed join-admission lock keeps join B out while join A is + /// gated mid-IAM. Remove that lock and B's IAM write lands during A's + /// admission: the assertion on the empty IAM log turns red. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[serial] + async fn test_peer_join_admission_serializes_across_separate_nodes() { + publish_ready_iam_context().await; + + let now = OffsetDateTime::now_utc().replace_nanosecond(0).expect("truncate nanos"); + let local = PeerInfo { + deployment_id: "site-local".to_string(), + ..peer("site-local", "https://local.example:9000") + }; + let remote = PeerInfo { + deployment_id: "site-remote".to_string(), + ..peer("site-remote", "https://remote.example:9000") + }; + let seed = SiteReplicationState { + peers: BTreeMap::from([ + (local.deployment_id.clone(), local.clone()), + (remote.deployment_id.clone(), remote.clone()), + ]), + updated_at: Some(now - Duration::from_secs(60)), + ..Default::default() + }; + save_site_replication_state(&seed).await.expect("seed state"); + + let join_peers = BTreeMap::from([ + (local.deployment_id.clone(), local.clone()), + (remote.deployment_id.clone(), remote.clone()), + ]); + let join_req = |updated_at: OffsetDateTime, secret: &str| SRPeerJoinReq { + svc_acct_access_key: "svc-join".to_string(), + svc_acct_secret_key: secret.to_string(), + svc_acct_parent: "root".to_string(), + peers: join_peers.clone(), + updated_at: Some(updated_at), + }; + + let iam_log: Arc>> = Arc::new(StdMutex::new(Vec::new())); + let (a_entered_tx, a_entered_rx) = tokio::sync::oneshot::channel(); + let (a_gate_tx, a_gate_rx) = tokio::sync::oneshot::channel::<()>(); + + // Node A (older join, T1) pauses inside its IAM step while holding + // only the distributed admission lock. + let log_a = iam_log.clone(); + let req_a = join_req(now - Duration::from_secs(30), "secret-a"); + let join_a = tokio::spawn(async move { + admit_peer_join_across_nodes("https://local.example:9000".to_string(), req_a, true, move |_req| async move { + let _ = a_entered_tx.send(()); + let _ = a_gate_rx.await; + log_a.lock().expect("iam log").push("iam-a"); + Ok(()) + }) + .await + }); + a_entered_rx.await.expect("node A reached its IAM step"); + + // Node B (newer join, T2) arrives on "another node": no process-local + // guard applies. The distributed admission lock must hold it. + let log_b = iam_log.clone(); + let req_b = join_req(now, "secret-b"); + let join_b = tokio::spawn(async move { + admit_peer_join_across_nodes("https://local.example:9000".to_string(), req_b, true, move |_req| async move { + log_b.lock().expect("iam log").push("iam-b"); + Ok(()) + }) + .await + }); + tokio::time::sleep(Duration::from_millis(200)).await; + assert!( + iam_log.lock().expect("iam log").is_empty(), + "node B ran its IAM step while node A was still mid-admission: {:?}", + iam_log.lock().expect("iam log") + ); + + a_gate_tx.send(()).expect("release node A"); + let outcome_a = join_a.await.expect("node A task").expect("node A admission"); + let outcome_b = join_b.await.expect("node B task").expect("node B admission"); + assert!(matches!(outcome_a, PeerJoinOutcome::Applied(..)), "node A must commit first"); + assert!( + matches!(outcome_b, PeerJoinOutcome::Applied(..)), + "the newer join B must still apply after A" + ); + assert_eq!( + *iam_log.lock().expect("iam log"), + vec!["iam-a", "iam-b"], + "IAM writes must land in admission order, ending on the committed join's secret" + ); + assert_eq!( + load_site_replication_state().await.expect("reload").updated_at, + Some(now), + "the persisted state must end on join B, matching the last IAM write" + ); + } + + /// P1-15 PR2: the three-way contract of `finalize_pending_rotation_if_complete` + /// — no pending means "already finalized" (true, nothing written), a + /// different or incomplete rotation is left alone (false), and a fully + /// acked rotation is cleared in the same transaction that reports true. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[serial] + async fn test_finalize_pending_rotation_three_way_contract() { + publish_ready_iam_context().await; + + let local_peer = PeerInfo { + deployment_id: "site-local".to_string(), + ..peer("site-local", "https://local.example:9000") + }; + let remote_peer = PeerInfo { + deployment_id: "site-remote".to_string(), + ..peer("site-remote", "https://remote.example:9000") + }; + let seed = SiteReplicationState { + peers: BTreeMap::from([ + (local_peer.deployment_id.clone(), local_peer.clone()), + (remote_peer.deployment_id.clone(), remote_peer.clone()), + ]), pending_rotation: Some(PendingRotation { - id: "rot-legacy".to_string(), + id: "rot-final".to_string(), access_key: "svc-account".to_string(), + peers: BTreeMap::from([ + (local_peer.deployment_id.clone(), local_peer.clone()), + (remote_peer.deployment_id.clone(), remote_peer.clone()), + ]), ..Default::default() }), ..Default::default() }; save_site_replication_state(&seed).await.expect("seed state"); - const ROUNDS: usize = 8; - for round in 0..ROUNDS { - let legacy_id = format!("legacy-{round}"); - let legacy = tokio::spawn(async move { - // Exactly what an unmigrated call site does today. - let _guard = site_replication_state_process_guard().await; - let mut state = load_site_replication_state().await.expect("legacy load"); - if let Some(pending) = state.pending_rotation.as_mut() { - pending.secret_candidates.push(legacy_id); - } - save_site_replication_state(&state).await.expect("legacy save"); - }); - let ack_id = format!("ack-{round}"); - let migrated = tokio::spawn(async move { - mark_pending_rotation_peer_acked("rot-legacy", &ack_id) - .await - .expect("transaction writer"); - }); - legacy.await.expect("legacy task"); - migrated.await.expect("transaction task"); - } + assert!( + !finalize_pending_rotation_if_complete("other-rotation", &local_peer) + .await + .expect("mismatched id"), + "a different rotation id must not finalize" + ); + assert!( + !finalize_pending_rotation_if_complete("rot-final", &local_peer) + .await + .expect("incomplete acks"), + "an un-acked remote peer must block finalization" + ); + assert!( + load_site_replication_state() + .await + .expect("reload") + .pending_rotation + .is_some(), + "the pending rotation must survive both refusals" + ); - let final_state = load_site_replication_state().await.expect("reload"); - let pending = final_state.pending_rotation.expect("pending rotation survives"); - for round in 0..ROUNDS { - assert!( - pending.secret_candidates.contains(&format!("legacy-{round}")), - "legacy writer update {round} was lost; candidates: {:?}", - pending.secret_candidates - ); - assert!( - pending.acked_deployment_ids.contains(&format!("ack-{round}")), - "transaction writer update {round} was lost; acks: {:?}", - pending.acked_deployment_ids - ); - } + mark_pending_rotation_peer_acked("rot-final", &remote_peer.deployment_id) + .await + .expect("ack remote"); + assert!( + finalize_pending_rotation_if_complete("rot-final", &local_peer) + .await + .expect("finalize"), + "a fully acked rotation must finalize" + ); + assert!( + load_site_replication_state() + .await + .expect("reload") + .pending_rotation + .is_none(), + "finalization must clear the pending rotation" + ); + assert!( + finalize_pending_rotation_if_complete("rot-final", &local_peer) + .await + .expect("idempotent"), + "no pending rotation means already finalized" + ); } - /// P1-15 review follow-up: isolates the DISTRIBUTED guard. Both writers - /// bypass the process mutex, which is what two separate nodes do — the - /// mutex is per process and cannot serialize them. Only the state-object - /// write lock keeps their read-modify-write sequences apart; drop it and - /// this test loses an update. + /// P1-15 review follow-up: isolates the DISTRIBUTED guard, which is now + /// the whole boundary. Two "nodes" run the production transaction + /// concurrently; only the state-object write lock keeps their + /// read-modify-write sequences apart — drop it and this test loses an + /// update. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] async fn test_state_object_lock_serializes_writers_from_separate_nodes() { @@ -15363,11 +15972,10 @@ mod tests { }; save_site_replication_state(&seed).await.expect("seed state"); - // Each "node" runs the production transaction minus the process - // mutex — the distributed state-object lock is the only thing left - // to keep them apart. + // Each "node" runs the production transaction; the distributed + // state-object lock is the only thing keeping them apart. fn node_local_update(candidate: String) -> impl std::future::Future> { - update_site_replication_state_as_separate_node(move |state| { + update_site_replication_state(move |state| { if let Some(pending) = state.pending_rotation.as_mut() { pending.secret_candidates.push(candidate); } @@ -15397,10 +16005,9 @@ mod tests { } /// P1-15 review follow-up: the sending side of the peer-edit ordering - /// fence, driven by two separate nodes. Both bypass the process mutex — - /// which is exactly what two nodes of one site do — so the generation is - /// only unique because it is allocated inside the state transaction, under - /// the distributed state-object lock. Two nodes sharing a generation would + /// fence, driven by two separate nodes. The generation is unique only + /// because it is allocated inside the state transaction, under the + /// distributed state-object lock. Two nodes sharing a generation would /// leave the receiver unable to tell which edit is newer. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] @@ -15419,7 +16026,7 @@ mod tests { save_site_replication_state(&seed).await.expect("seed state"); fn node_local_allocate() -> impl std::future::Future> { - update_site_replication_state_as_separate_node(|state| Ok(next_peer_edit_generation(state))) + update_site_replication_state(|state| Ok(next_peer_edit_generation(state))) } const ROUNDS: usize = 8; diff --git a/rustfs/src/admin/site_replication_state.rs b/rustfs/src/admin/site_replication_state.rs index eb0414665..d39618a2b 100644 --- a/rustfs/src/admin/site_replication_state.rs +++ b/rustfs/src/admin/site_replication_state.rs @@ -18,28 +18,21 @@ //! `config/site-replication/state.json` is mutated by read-modify-write //! sequences spread over many call sites: admin handlers, the retry-event //! writers on every hook broadcast path, and the service-side reload driven -//! over node RPC. Historically only some of them held the process-local -//! mutex and none held a distributed lock across the whole RMW, so -//! concurrent writers overwrote each other (single-process for the unlocked -//! writers, cross-node for everyone). +//! over node RPC. //! //! `with_site_replication_state_lock` is the single transaction boundary: -//! it holds the process-local mutex AND the distributed config-object write -//! lock (the pattern proven by the repair state, -//! `update_site_replication_repair_state`) for the duration of the caller's -//! closure. 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 +//! it holds the distributed config-object write lock (the pattern proven by +//! the repair state, `update_site_replication_repair_state`) for the +//! duration of the caller's closure. The object lock is the sole mechanism — +//! it is the only thing that can serialize two nodes of the same site, so a +//! 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. //! -//! The process-local mutex is transitional: call sites still outside this -//! primitive serialize against migrated ones through it. Once every RMW -//! call site goes through here (P1-15 PR2) it will be removed, leaving the -//! object lock as the only mechanism. -//! -//! Lock order (unchanged from the historical comment next to the mutex): -//! lifecycle -> bucket operation -> repair admission -> state (process -//! mutex, then state object lock) -> per-bucket metadata. +//! Lock order: lifecycle -> bucket operation -> repair admission +//! -> state object lock -> per-bucket metadata. use crate::admin::storage_api::runtime::ECStore; use crate::admin::storage_api::s3::{S3Error, S3ErrorCode, S3Result}; @@ -53,24 +46,8 @@ use super::runtime_sources::current_object_store_handle; /// byte-level tolerant reload on the service side. pub(crate) const SITE_REPLICATION_STATE_PATH: &str = "config/site-replication/state.json"; -/// Transitional process-local mutex — see the module docs. Stays private to -/// this module (owner-local static, enforced by -/// `scripts/check_architecture_migration_rules.sh`); callers go through -/// [`site_replication_state_process_guard`]. -static SITE_REPLICATION_STATE_LOCK: std::sync::LazyLock> = - std::sync::LazyLock::new(|| tokio::sync::Mutex::new(())); - -/// Owner helper for the transitional process mutex: the RMW call sites in -/// `handlers::site_replication` that PR2 has not migrated to -/// [`with_site_replication_state_lock`] yet hold this guard so they stay -/// mutually exclusive with the migrated ones. Removed together with the -/// mutex once every call site runs inside the transaction boundary. -pub(crate) async fn site_replication_state_process_guard() -> tokio::sync::MutexGuard<'static, ()> { - SITE_REPLICATION_STATE_LOCK.lock().await -} - /// Run `operation` under the site-replication state transaction boundary: -/// process mutex first, then the distributed state-object write lock. +/// the distributed state-object write lock. pub(crate) async fn with_site_replication_state_lock(operation: F) -> S3Result where T: Send + 'static, @@ -83,21 +60,11 @@ where /// Context-store variant for callers that resolve their store from an /// explicit [`AppContext`] (the service-side reload driven over node RPC). +/// +/// This is the whole boundary: a state-object write lock serializes writers +/// in *different* processes, which is what two nodes of one site are and +/// what a process mutex could never cover. pub(crate) async fn with_site_replication_state_lock_on(store: Arc, operation: F) -> S3Result -where - T: Send + 'static, - F: FnOnce() -> Fut + Send + 'static, - Fut: std::future::Future> + Send + 'static, -{ - let _process_guard = SITE_REPLICATION_STATE_LOCK.lock().await; - with_site_replication_state_object_lock(store, operation).await -} - -/// The distributed half of the boundary on its own: the state-object write -/// lock, without the process mutex. This is the only thing that serializes -/// writers in *different* processes (the mutex cannot), so it is also what -/// the separate-nodes regression test drives. -pub(crate) async fn with_site_replication_state_object_lock(store: Arc, operation: F) -> S3Result where T: Send + 'static, F: FnOnce() -> Fut + Send + 'static, diff --git a/scripts/check_architecture_migration_rules.sh b/scripts/check_architecture_migration_rules.sh index dba51a333..476fbce90 100755 --- a/scripts/check_architecture_migration_rules.sh +++ b/scripts/check_architecture_migration_rules.sh @@ -4033,7 +4033,7 @@ if [[ -s "$ECSTORE_REMOTE_TIER_DELETE_STATE_BYPASS_HITS_FILE" ]]; then report_failure "remote tier delete state access must stay behind ECStore tier sweeper owner helpers: $(paste -sd '; ' "$ECSTORE_REMOTE_TIER_DELETE_STATE_BYPASS_HITS_FILE")" fi -RUSTFS_OWNER_LOCAL_STATIC_NAMES='(KEYSTONE_AUTH|KEYSTONE_MAPPER|KEYSTONE_CONFIG|LICENSE_STATE|LICENSE_VERIFIER|CPU_CONT_GUARD|PROFILING_CANCEL_TOKEN|MEMORY_SYSTEM|DIAL9_TELEMETRY_GUARD|DISPLAY_CONFIG_SNAPSHOT|GLOBAL_CONFIG_SNAPSHOT|BUFFER_CONFIG_SINGLETON|BUFFER_PROFILE_ENABLED|LEGACY_CREDENTIAL_WARNED_KEYS|CONSOLE_CONFIG|ACTIVE_HTTP_REQUESTS|USE_STARSHARD_CACHE|BUCKET_CACHE_SMALL|BUCKET_CACHE_LARGE|GLOBAL_SSE_DEK_PROVIDER|SSE_TEST_LOCK|AUTH_FS|LOCK_STATS|DEADLOCK_DETECTOR|GET_OBJECT_BUFFER_THRESHOLD_WARNED|GET_READER_STREAM_BUFFER_SIZE_OVERRIDE|OBJECT_SEEK_SUPPORT_THRESHOLD|OBJECT_SEEK_SUPPORT_CONCURRENCY_THRESHOLDS|SUPPORTED_HEADERS|SITE_REPLICATION_PEER_CLIENT|SITE_REPLICATION_STATE_LOCK|AUDIT_MODULE_ENABLED|NOTIFY_MODULE_ENABLED|PERSISTED_NOTIFY_MODULE_ENABLED|PERSISTED_AUDIT_MODULE_ENABLED|PERSISTED_MODULE_SWITCH_CONFIGURED|DELETE_TAIL_TOTAL|DELETE_CLEANUP_TOTAL|DELETE_REPLICATION_TOTAL|DELETE_NOTIFY_TOTAL|EMBEDDED_SERVER_STARTED|TEST_OUTBOUND_TLS_GENERATION|TEST_REMAINING_FAILURES|CAPACITY_DIRTY_SCOPE_ENV|CAPACITY_DIRTY_SCOPE_INIT|GLOBAL_ENV)' +RUSTFS_OWNER_LOCAL_STATIC_NAMES='(KEYSTONE_AUTH|KEYSTONE_MAPPER|KEYSTONE_CONFIG|LICENSE_STATE|LICENSE_VERIFIER|CPU_CONT_GUARD|PROFILING_CANCEL_TOKEN|MEMORY_SYSTEM|DIAL9_TELEMETRY_GUARD|DISPLAY_CONFIG_SNAPSHOT|GLOBAL_CONFIG_SNAPSHOT|BUFFER_CONFIG_SINGLETON|BUFFER_PROFILE_ENABLED|LEGACY_CREDENTIAL_WARNED_KEYS|CONSOLE_CONFIG|ACTIVE_HTTP_REQUESTS|USE_STARSHARD_CACHE|BUCKET_CACHE_SMALL|BUCKET_CACHE_LARGE|GLOBAL_SSE_DEK_PROVIDER|SSE_TEST_LOCK|AUTH_FS|LOCK_STATS|DEADLOCK_DETECTOR|GET_OBJECT_BUFFER_THRESHOLD_WARNED|GET_READER_STREAM_BUFFER_SIZE_OVERRIDE|OBJECT_SEEK_SUPPORT_THRESHOLD|OBJECT_SEEK_SUPPORT_CONCURRENCY_THRESHOLDS|SUPPORTED_HEADERS|SITE_REPLICATION_PEER_CLIENT|AUDIT_MODULE_ENABLED|NOTIFY_MODULE_ENABLED|PERSISTED_NOTIFY_MODULE_ENABLED|PERSISTED_AUDIT_MODULE_ENABLED|PERSISTED_MODULE_SWITCH_CONFIGURED|DELETE_TAIL_TOTAL|DELETE_CLEANUP_TOTAL|DELETE_REPLICATION_TOTAL|DELETE_NOTIFY_TOTAL|EMBEDDED_SERVER_STARTED|TEST_OUTBOUND_TLS_GENERATION|TEST_REMAINING_FAILURES|CAPACITY_DIRTY_SCOPE_ENV|CAPACITY_DIRTY_SCOPE_INIT|GLOBAL_ENV)' ( cd "$ROOT_DIR"