diff --git a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs index 59aa7be9c..6d6419e70 100644 --- a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs @@ -18,7 +18,7 @@ use crate::cluster::rpc::client::{ node_service_time_out_client, }; use crate::cluster::rpc::set_tonic_mutation_body_digest; -use crate::core::pools::PoolMeta; +use crate::core::pools::{PoolMeta, PoolMetaWriteState}; use crate::disk::error::DiskError; use crate::disk::error::{Error, Result}; use crate::disk::error_reduce::{BUCKET_OP_IGNORED_ERRS, is_all_buckets_not_found, reduce_write_quorum_errs}; @@ -324,6 +324,14 @@ pub trait PeerS3Client: Debug + Sync + Send + 'static { async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, _fenced_pools: &[usize]) -> Result { self.heal_bucket(bucket, opts).await } + async fn heal_bucket_with_fence_from_movement_guarded_coordinator( + &self, + bucket: &str, + opts: &HealOpts, + fenced_pools: &[usize], + ) -> Result { + self.heal_bucket_with_fence(bucket, opts, fenced_pools).await + } async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()>; async fn list_bucket(&self, opts: &BucketOptions) -> Result>; async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()>; @@ -381,6 +389,25 @@ impl S3PeerSys { } pub async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, fenced_pools: &[usize]) -> Result { + self.heal_bucket_with_fence_inner(bucket, opts, fenced_pools, false).await + } + + pub async fn heal_bucket_with_fence_from_movement_guarded_coordinator( + &self, + bucket: &str, + opts: &HealOpts, + fenced_pools: &[usize], + ) -> Result { + self.heal_bucket_with_fence_inner(bucket, opts, fenced_pools, true).await + } + + async fn heal_bucket_with_fence_inner( + &self, + bucket: &str, + opts: &HealOpts, + fenced_pools: &[usize], + movement_guard_held: bool, + ) -> Result { let mut opts = *opts; let mut futures = Vec::with_capacity(self.clients.len()); for client in self.clients.iter() { @@ -403,7 +430,14 @@ impl S3PeerSys { let opts_clone = opts; let heal_bucket_results_clone = heal_bucket_results.clone(); futures.push(async move { - match client.heal_bucket_with_fence(bucket, &opts_clone, fenced_pools).await { + let result = if movement_guard_held { + client + .heal_bucket_with_fence_from_movement_guarded_coordinator(bucket, &opts_clone, fenced_pools) + .await + } else { + client.heal_bucket_with_fence(bucket, &opts_clone, fenced_pools).await + }; + match result { Ok(res) => { heal_bucket_results_clone.write().await[idx] = res; None @@ -698,6 +732,63 @@ impl LocalPeerS3Client { .filter(|disk| usize::try_from(disk.endpoint().pool_idx).is_ok_and(|pool_idx| pools.contains(&pool_idx))) .collect() } + + async fn heal_bucket_with_fence_inner( + &self, + bucket: &str, + opts: &HealOpts, + fenced_pools: &[usize], + movement_guard_held: bool, + ) -> Result { + let disks = self.local_disks_for_pools().await.into_iter().map(Some).collect(); + let store = runtime_sources::object_store_handle().filter(|store| Arc::ptr_eq(&store.ctx, &self.instance_ctx)); + #[cfg(not(test))] + if store.is_none() { + return Err(Error::other("bucket heal refused: pool metadata is unavailable for this instance")); + } + let movement_gate = store.as_ref().map(|store| store.ctx.data_movement_operation_gate()); + let movement_guard = try_acquire_bucket_heal_movement_guard(movement_gate.as_ref(), movement_guard_held)?; + let save_guard = acquire_bucket_heal_write_guard(store.as_ref().map(|store| &store.pool_meta_save_gate)).await?; + let result = heal_bucket_local_on_disks_with_pool_meta( + bucket, + opts, + disks, + store.as_ref().map(|store| &store.pool_meta), + fenced_pools, + ) + .await; + drop(save_guard); + drop(movement_guard); + result + } +} + +fn try_acquire_bucket_heal_movement_guard<'a>( + gate: Option<&'a Arc>>, + movement_guard_held: bool, +) -> Result>> { + if movement_guard_held { + return Ok(None); + } + let Some(gate) = gate else { + return Ok(None); + }; + // Do not queue a receiver behind a movement writer while its coordinator + // holds another node's read guard; failing fast breaks that cross-node cycle. + gate.try_read() + .map(Some) + .map_err(|_| crate::error::StorageError::SlowDown.into()) +} + +async fn acquire_bucket_heal_write_guard<'a>( + gate: Option<&'a tokio::sync::Mutex>, +) -> Result>> { + let Some(gate) = gate else { + return Ok(None); + }; + let guard = gate.lock().await; + guard.ensure_write_safe("bucket heal cannot run while pool metadata requires recovery")?; + Ok(Some(guard)) } #[async_trait] @@ -711,14 +802,16 @@ impl PeerS3Client for LocalPeerS3Client { } async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, fenced_pools: &[usize]) -> Result { - let disks = self.local_disks_for_pools().await.into_iter().map(Some).collect(); - let store = runtime_sources::object_store_handle().filter(|store| Arc::ptr_eq(&store.ctx, &self.instance_ctx)); - #[cfg(not(test))] - if store.is_none() { - return Err(Error::other("bucket heal refused: pool metadata is unavailable for this instance")); - } - heal_bucket_local_on_disks_with_pool_meta(bucket, opts, disks, store.as_ref().map(|store| &store.pool_meta), fenced_pools) - .await + self.heal_bucket_with_fence_inner(bucket, opts, fenced_pools, false).await + } + + async fn heal_bucket_with_fence_from_movement_guarded_coordinator( + &self, + bucket: &str, + opts: &HealOpts, + fenced_pools: &[usize], + ) -> Result { + self.heal_bucket_with_fence_inner(bucket, opts, fenced_pools, true).await } async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { @@ -1656,7 +1749,7 @@ async fn clone_drives() -> Vec> { #[cfg(test)] mod tests { use super::*; - use crate::core::pools::{PoolDecommissionInfo, PoolStatus}; + use crate::core::pools::{PoolDecommissionInfo, PoolMetaReplicaState, PoolStatus}; use crate::disk::WalkDirOptions; use crate::disk::disk_store::LocalDiskWrapper; use crate::disk::endpoint::Endpoint; @@ -2241,6 +2334,52 @@ mod tests { reset_local_disk_test_state().await; } + #[tokio::test] + async fn local_bucket_heal_refuses_receiver_with_unsafe_pool_metadata() { + let gate = tokio::sync::Mutex::new(PoolMetaWriteState::default()); + gate.lock().await.observe_replicas(PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + + let err = acquire_bucket_heal_write_guard(Some(&gate)) + .await + .expect_err("receiver-side bucket heal must honor the local pool metadata gate"); + + assert!( + err.to_string() + .contains("bucket heal cannot run while pool metadata requires recovery") + ); + } + + #[tokio::test] + async fn receiver_bucket_heal_fails_fast_behind_queued_movement_writer() { + let gate = Arc::new(tokio::sync::RwLock::new(())); + let coordinator_guard = gate.read().await; + let writer_gate = gate.clone(); + let mut writer = tokio::spawn(async move { + let _writer_guard = writer_gate.write().await; + }); + while gate.try_read().is_ok() { + tokio::task::yield_now().await; + } + + let err = try_acquire_bucket_heal_movement_guard(Some(&gate), false) + .expect_err("receiver must not wait behind a queued movement writer"); + assert_eq!(err, crate::error::StorageError::SlowDown.into()); + assert!( + try_acquire_bucket_heal_movement_guard(Some(&gate), true) + .expect("coordinator-owned movement guard should be reused") + .is_none() + ); + + drop(coordinator_guard); + tokio::time::timeout(std::time::Duration::from_secs(1), &mut writer) + .await + .expect("queued movement writer should proceed after the coordinator guard is released") + .expect("movement writer task should not panic"); + } + #[tokio::test] #[serial] async fn heal_bucket_keeps_suspended_pool_volume_on_remove() { diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 3790ba7b5..1942c9514 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -618,6 +618,23 @@ fn spawn_decommission_index_cancelers( async move { store.do_decommission_in_routine(canceler, idx, entry_budget).await } }); if let Err(err) = await_decommission_worker(idx, worker).await { + if let Err(blocked) = store + .ensure_pool_meta_side_effects_safe("decommission paused because pool metadata requires recovery") + .await + { + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + pool_index = idx, + state = "routine_blocked", + error = %blocked, + "Decommission routine paused without changing terminal state" + ); + store.release_decommission_canceler_slot(idx, &canceler).await; + stop_queue = true; + continue; + } error!( event = EVENT_DECOMMISSION_STATE, component = LOG_COMPONENT_ECSTORE, @@ -975,6 +992,7 @@ fn reconcile_decommission_unresolved_entries_for_completion( Ok(()) } +#[cfg(test)] async fn run_decommission_side_effect( rx: &CancellationToken, operation_gate: &Arc>, @@ -1043,10 +1061,6 @@ fn resolve_decommission_preflight_heal_result(bucket: &str, result: Result result.map_err(|err| Error::other(format!("decommission preflight heal failed for bucket {bucket}: {err}"))) } -fn resolve_decommission_bucket_done_save_result(result: Result<()>, idx: usize, bucket: &str) -> Result<()> { - result.map_err(|err| Error::other(format!("decommission metadata save failed for pool {idx} bucket {bucket}: {err}"))) -} - fn resolve_decommission_optional_bucket_config_result(bucket: &str, stage: &str, result: Result) -> Result> { match result { Ok(config) => Ok(Some(config)), @@ -1660,6 +1674,56 @@ pub(crate) fn merge_pool_status_refresh(current: &mut PoolMeta, persisted: PoolM merged_newer } +fn merge_pool_meta_updates_for_save( + persisted: &mut PoolMeta, + current: &PoolMeta, + indices: &[usize], + operation: &str, +) -> Result<()> { + if persisted.pools.is_empty() { + *persisted = current.clone(); + return Ok(()); + } + if persisted.version != current.version { + return Err(Error::other(format!( + "{operation}: pool metadata version changed from {} to {}", + current.version, persisted.version + ))); + } + + for &idx in indices { + let current_pool = current + .pools + .get(idx) + .ok_or_else(|| invalid_decommission_pool_index_error(current.pools.len(), idx))?; + let persisted_count = persisted.pools.len(); + let persisted_pool = persisted + .pools + .get_mut(idx) + .ok_or_else(|| invalid_decommission_pool_index_error(persisted_count, idx))?; + if current_pool.id != idx || persisted_pool.id != idx || current_pool.cmd_line != persisted_pool.cmd_line { + return Err(Error::other(format!("{operation}: pool metadata layout changed for pool {idx}"))); + } + *persisted_pool = current_pool.clone(); + } + + Ok(()) +} + +fn publish_pool_meta_updates(current: &mut PoolMeta, saved: &PoolMeta, indices: &[usize]) { + for &idx in indices { + let Some(saved_pool) = saved.pools.get(idx) else { + continue; + }; + let Some(current_pool) = current.pools.get_mut(idx) else { + continue; + }; + if current_pool.id == idx && saved_pool.id == idx && current_pool.cmd_line == saved_pool.cmd_line { + *current_pool = saved_pool.clone(); + } + } +} + fn resolve_start_decommission_pool_meta_reload_result(result: Result<()>) -> Result<()> { resolve_decommission_pool_meta_reload_result(result, "start_decommission") } @@ -1933,8 +1997,8 @@ pub(crate) fn observe_pool_activation_start_attempt(kind: PoolActivationStartKin } } -fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) { - *pool_meta = previous_pool_meta; +fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: &PoolMeta, indices: &[usize]) { + publish_pool_meta_updates(pool_meta, previous_pool_meta, indices); } #[derive(Debug)] @@ -1979,8 +2043,25 @@ fn commit_decommission_cancel(pool_meta: &mut PoolMeta, idx: usize, commit: Deco Ok(()) } -fn rollback_start_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) { - rollback_decommission_pool_meta(pool_meta, previous_pool_meta); +fn rollback_start_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: &PoolMeta, indices: &[usize]) { + let active_updates = indices + .iter() + .filter_map(|&idx| pool_meta.pools.get(idx).map(|pool| (idx, pool.last_update))) + .collect::>(); + rollback_decommission_pool_meta(pool_meta, previous_pool_meta, indices); + let rollback_at = OffsetDateTime::now_utc(); + for (idx, active_update) in active_updates { + if let Some(pool) = pool_meta.pools.get_mut(idx) { + pool.last_update = std::cmp::max(rollback_at, active_update + Duration::nanoseconds(1)); + } + } +} + +fn ensure_pool_meta_write_fence(guard: &rustfs_lock::NamespaceLockGuard, operation: &str) -> Result<()> { + if guard.is_lock_lost() { + return Err(Error::other(format!("{operation}: pool metadata distributed fence was lost"))); + } + Ok(()) } fn ensure_pool_not_left_in_cmdline_after_decommission(position: usize, cmd_line: &str, completed: bool) -> Result<()> { @@ -2657,6 +2738,34 @@ impl PoolMetaReplicaState { } } +#[derive(Debug, Clone, Copy, Default)] +pub(crate) struct PoolMetaWriteState { + write_blocked: bool, +} + +impl PoolMetaWriteState { + pub(crate) fn observe_replicas(&mut self, replica_state: PoolMetaReplicaState) { + self.write_blocked |= !replica_state.repair_write_safe; + } + + fn block_writes(&mut self) { + self.write_blocked = true; + } + + fn restore_writes(&mut self, was_write_blocked: bool) { + self.write_blocked = was_write_blocked; + } + + pub(crate) fn ensure_write_safe(self, operation: &str) -> Result<()> { + if !self.write_blocked { + return Ok(()); + } + Err(Error::other(format!( + "{operation}: pool metadata writes remain blocked after a recovery-required replica state; restart after all replicas are readable and consistent, with compatible formats" + ))) + } +} + #[derive(Debug)] struct PoolMetaSelection { meta: PoolMeta, @@ -2850,12 +2959,47 @@ fn select_pool_meta_replica(replicas: Vec) -> Result(pools: Vec>, no_lock: bool) -> Vec +where + S: EcstoreObjectIO, +{ + join_all(pools.into_iter().map(|pool| read_pool_meta_replica(pool, no_lock))).await +} + +fn select_pool_meta_replicas_observing( + write_state: &mut PoolMetaWriteState, + replicas: Vec, +) -> Result { + if replicas + .iter() + .any(|replica| matches!(replica, PoolMetaReplica::Unreadable(_))) + { + write_state.block_writes(); + } + let selection = select_pool_meta_replica(replicas); + if selection.is_err() { + write_state.block_writes(); + } + selection +} + async fn load_pool_meta_replicas(pools: Vec>, no_lock: bool) -> Result where S: EcstoreObjectIO, { - let replicas = join_all(pools.into_iter().map(|pool| read_pool_meta_replica(pool, no_lock))).await; - select_pool_meta_replica(replicas) + select_pool_meta_replica(read_pool_meta_replicas(pools, no_lock).await) +} + +async fn load_pool_meta_replicas_observing( + pools: Vec>, + no_lock: bool, + write_state: &mut PoolMetaWriteState, +) -> Result +where + S: EcstoreObjectIO, +{ + let replicas = read_pool_meta_replicas(pools, no_lock).await; + select_pool_meta_replicas_observing(write_state, replicas) } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -3494,9 +3638,11 @@ impl PoolMeta { Ok(()) } - pub async fn load(&mut self, _pool: Arc, pools: Vec>) -> Result<()> { - let selection = load_pool_meta_replicas(pools, false).await?; - *self = selection.meta; + pub async fn load(&mut self, pool: Arc, pools: Vec>) -> Result<()> { + let pool_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?; + let _pool_meta_guard = pool_meta_lock.get_read_lock(get_lock_acquire_timeout()).await?; + let replica_state = self.load_no_lock_from_replicas(pools).await?; + replica_state.ensure_write_safe("pool metadata load failed")?; Ok(()) } @@ -3511,6 +3657,19 @@ impl PoolMeta { Ok(selection.replica_state) } + pub(crate) async fn load_no_lock_from_replicas_observing( + &mut self, + pools: Vec>, + write_state: &mut PoolMetaWriteState, + ) -> Result + where + S: EcstoreObjectIO, + { + let selection = load_pool_meta_replicas_observing(pools, true, write_state).await?; + *self = selection.meta; + Ok(selection.replica_state) + } + fn encode_config_data(&self) -> Result> { self.encode_config_data_for_v2_gate(pool_meta_v2_writer_enabled()) } @@ -3562,15 +3721,15 @@ impl PoolMeta { } pub async fn save(&self, pools: Vec>) -> Result<()> { - let data = self.encode_config_data()?; - if data.is_empty() { - return Ok(()); - } - for pool in pools { - save_config(pool, POOL_META_NAME, data.clone()).await?; - } - - Ok(()) + let pool = pools + .first() + .cloned() + .ok_or_else(|| Error::other("pool metadata save failed: no storage pools available"))?; + let pool_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?; + let pool_meta_guard = pool_meta_lock.get_write_lock(get_lock_acquire_timeout()).await?; + let selection = load_pool_meta_replicas(pools.clone(), true).await?; + selection.replica_state.ensure_write_safe("pool metadata save failed")?; + self.save_no_lock_with_fence(pools, pool_meta_guard.lock_lost_signal()).await } /// Startup has a single elected local writer, so it must not depend on namespace locks here. @@ -3585,31 +3744,13 @@ impl PoolMeta { where S: EcstoreObjectIO, { - let data = self.encode_config_data()?; - if data.is_empty() { - return Ok(()); - } - for pool in pools { - save_config_with_opts( - pool, - POOL_META_NAME, - data.clone(), - &ObjectOptions { - max_parity: true, - no_lock: true, - ..Default::default() - }, - ) - .await?; - } - - Ok(()) + self.save_no_lock_with_fence(pools, None).await } - async fn save_no_lock_with_activation_fence( + async fn save_no_lock_with_fence( &self, pools: Vec>, - activation_fence: PoolRebalanceActivationFence, + lock_lost: Option>, ) -> Result<()> where S: EcstoreObjectIO, @@ -3618,9 +3759,47 @@ impl PoolMeta { if data.is_empty() { return Ok(()); } + for pool in pools { + if lock_lost.as_ref().is_some_and(|signal| signal.is_lost()) { + return Err(Error::other("pool metadata distributed fence was lost before a replica write")); + } + let mut opts = ObjectOptions { + max_parity: true, + no_lock: true, + ..Default::default() + }; + if let Some(signal) = lock_lost.as_ref() { + opts.add_namespace_lock_lost_signal(signal.clone()); + } + save_config_with_opts(pool, POOL_META_NAME, data.clone(), &opts).await?; + if lock_lost.as_ref().is_some_and(|signal| signal.is_lost()) { + return Err(Error::other("pool metadata distributed fence was lost during a replica write")); + } + } + + Ok(()) + } + + async fn save_no_lock_with_activation_fence( + &self, + pools: Vec>, + write_state: &mut PoolMetaWriteState, + activation_fence: PoolRebalanceActivationFence, + ) -> Result + where + S: EcstoreObjectIO, + { + let was_write_blocked = write_state.write_blocked; + // Arm before the first replica write. If this future is dropped, the + // mutex guard is released with the sticky write gate still blocked. + write_state.block_writes(); + let data = self.encode_config_data()?; + if data.is_empty() { + return Ok(was_write_blocked); + } let mut pools = pools.into_iter(); let Some(canonical_pool) = pools.next() else { - return Ok(()); + return Ok(was_write_blocked); }; let mut opts = ObjectOptions { max_parity: true, @@ -3664,6 +3843,33 @@ impl PoolMeta { } } + Ok(was_write_blocked) + } + + async fn save_no_lock_armed( + &self, + pools: Vec>, + write_state: &mut PoolMetaWriteState, + lock_lost: Option>, + ) -> Result + where + S: EcstoreObjectIO, + { + let was_write_blocked = write_state.write_blocked; + // Arm before the first replica write. If this future is dropped, the + // mutex guard is released with the sticky write gate still blocked. + write_state.block_writes(); + self.save_no_lock_with_fence(pools, lock_lost).await?; + Ok(was_write_blocked) + } + + #[cfg(test)] + async fn save_no_lock_observing(&self, pools: Vec>, write_state: &mut PoolMetaWriteState) -> Result<()> + where + S: EcstoreObjectIO, + { + let was_write_blocked = self.save_no_lock_armed(pools, write_state, None).await?; + write_state.restore_writes(was_write_blocked); Ok(()) } @@ -4532,13 +4738,104 @@ impl ECStore { Ok(apply_expiry_rule_in(self.clone(), &event, &LcEventSrc::Scanner, &object_info).await) } - async fn save_current_pool_meta(&self) -> Result<()> { - let _save_guard = self.pool_meta_save_gate.lock().await; - let snapshot = { - let pool_meta = self.pool_meta.read().await; - pool_meta.clone() + async fn acquire_pool_meta_write_guard( + &self, + write_state: &mut PoolMetaWriteState, + operation: &str, + ) -> Result<(rustfs_lock::NamespaceLockGuard, PoolMeta)> { + write_state.ensure_write_safe(operation)?; + let pool = self + .pools + .first() + .cloned() + .ok_or_else(|| Error::other(format!("{operation}: no storage pools available")))?; + let pool_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?; + let pool_meta_guard = pool_meta_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(activation_pool_meta_lock_error)?; + let selection = load_pool_meta_replicas_observing(self.pools.clone(), true, write_state).await?; + write_state.observe_replicas(selection.replica_state); + write_state.ensure_write_safe(operation)?; + Ok((pool_meta_guard, selection.meta)) + } + + pub(crate) async fn ensure_pool_meta_side_effects_safe(&self, operation: &str) -> Result<()> { + self.pool_meta_save_gate.lock().await.ensure_write_safe(operation) + } + + pub(crate) async fn load_runtime_pool_meta(&self, operation: &str) -> Result { + let mut write_state = self.pool_meta_save_gate.lock().await; + write_state.ensure_write_safe(operation)?; + let pool = self + .pools + .first() + .cloned() + .ok_or_else(|| Error::other(format!("{operation}: no storage pools available")))?; + let pool_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?; + let _pool_meta_guard = pool_meta_lock.get_read_lock(get_lock_acquire_timeout()).await?; + let mut pool_meta = PoolMeta::default(); + let replica_state = pool_meta + .load_no_lock_from_replicas_observing(self.pools.clone(), &mut write_state) + .await?; + write_state.observe_replicas(replica_state); + write_state.ensure_write_safe(operation)?; + Ok(pool_meta) + } + + async fn run_guarded_decommission_side_effect( + &self, + rx: &CancellationToken, + operation_gate: &Arc>, + operation: F, + ) -> std::result::Result + where + F: FnOnce() -> Fut, + Fut: std::future::Future>, + E: From, + { + let _operation_guard = tokio::select! { + biased; + _ = rx.cancelled() => return Err(Error::OperationCanceled.into()), + guard = operation_gate.read() => guard, }; - snapshot.save(self.pools.clone()).await + + if rx.is_cancelled() { + return Err(Error::OperationCanceled.into()); + } + self.ensure_pool_meta_side_effects_safe("decommission side effect blocked because pool metadata requires recovery") + .await + .map_err(E::from)?; + if rx.is_cancelled() { + return Err(Error::OperationCanceled.into()); + } + + let result = operation().await; + if rx.is_cancelled() { + return Err(Error::OperationCanceled.into()); + } + result + } + + async fn save_current_pool_meta(&self, indices: &[usize]) -> Result<()> { + let mut save_guard = self.pool_meta_save_gate.lock().await; + let (pool_meta_guard, mut snapshot) = self + .acquire_pool_meta_write_guard(&mut save_guard, "pool metadata save failed") + .await?; + { + let pool_meta = self.pool_meta.read().await; + merge_pool_meta_updates_for_save(&mut snapshot, &pool_meta, indices, "pool metadata save failed")?; + } + let was_write_blocked = snapshot + .save_no_lock_armed(self.pools.clone(), &mut save_guard, pool_meta_guard.lock_lost_signal()) + .await?; + let mut pool_meta = self.pool_meta.write().await; + ensure_pool_meta_write_fence(&pool_meta_guard, "pool metadata save failed")?; + publish_pool_meta_updates(&mut pool_meta, &snapshot, indices); + ensure_pool_meta_write_fence(&pool_meta_guard, "pool metadata save failed")?; + drop(pool_meta); + save_guard.restore_writes(was_write_blocked); + Ok(()) } async fn persist_decommission_unresolved_entry( @@ -4559,7 +4856,10 @@ impl ECStore { async fn save_decommission_progress_checkpoint(&self, idx: usize, generation: OffsetDateTime) -> Result { // Lock order: save gate, then the short pool metadata read/write sections. Peer // reloads are intentionally performed by the caller after both locks are released. - let _save_guard = self.pool_meta_save_gate.lock().await; + let mut save_guard = self.pool_meta_save_gate.lock().await; + let (pool_meta_guard, mut snapshot) = self + .acquire_pool_meta_write_guard(&mut save_guard, "decommission progress save failed") + .await?; let (snapshot, checkpoint) = { let pool_meta = self.pool_meta.read().await; ensure_decommission_generation(&pool_meta, idx, generation)?; @@ -4572,23 +4872,73 @@ impl ECStore { return Ok(false); }; - let mut snapshot = pool_meta.clone(); - let Some(pool) = snapshot.pools.get_mut(idx) else { - return Err(invalid_decommission_pool_index_error(snapshot.pools.len(), idx)); + let mut current = pool_meta.clone(); + let current_count = current.pools.len(); + let Some(pool) = current.pools.get_mut(idx) else { + return Err(invalid_decommission_pool_index_error(current_count, idx)); }; pool.last_update = checkpoint.checkpoint_at; + merge_pool_meta_updates_for_save(&mut snapshot, ¤t, &[idx], "decommission progress save failed")?; (snapshot, checkpoint) }; - if let Err(err) = snapshot.save(self.pools.clone()).await { - let retry_after = OffsetDateTime::now_utc() + DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF; - let mut pool_meta = self.pool_meta.write().await; - pool_meta.defer_decommission_progress_checkpoint(idx, checkpoint, retry_after); - return Err(err); - } + let was_write_blocked = match snapshot + .save_no_lock_armed(self.pools.clone(), &mut save_guard, pool_meta_guard.lock_lost_signal()) + .await + { + Ok(was_write_blocked) => was_write_blocked, + Err(err) => { + let retry_after = OffsetDateTime::now_utc() + DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF; + let mut pool_meta = self.pool_meta.write().await; + pool_meta.defer_decommission_progress_checkpoint(idx, checkpoint, retry_after); + return Err(err); + } + }; let mut pool_meta = self.pool_meta.write().await; - Ok(pool_meta.commit_decommission_progress_checkpoint(idx, checkpoint)) + ensure_pool_meta_write_fence(&pool_meta_guard, "decommission progress save failed")?; + let committed = pool_meta.commit_decommission_progress_checkpoint(idx, checkpoint); + ensure_pool_meta_write_fence(&pool_meta_guard, "decommission progress save failed")?; + drop(pool_meta); + save_guard.restore_writes(was_write_blocked); + Ok(committed) + } + + async fn mark_decommission_bucket_done_and_save(&self, idx: usize, bucket: &DecomBucketInfo) -> Result { + let mut save_guard = self.pool_meta_save_gate.lock().await; + let (pool_meta_guard, mut snapshot) = self + .acquire_pool_meta_write_guard(&mut save_guard, "decommission bucket completion save failed") + .await?; + let changed = { + let mut pool_meta = self.pool_meta.write().await; + let changed = mark_decommission_bucket_done(&mut pool_meta, idx, bucket)?; + if changed { + merge_pool_meta_updates_for_save( + &mut snapshot, + &pool_meta, + &[idx], + "decommission bucket completion save failed", + )?; + } + changed + }; + if !changed { + return Ok(false); + } + + let was_write_blocked = snapshot + .save_no_lock_armed(self.pools.clone(), &mut save_guard, pool_meta_guard.lock_lost_signal()) + .await + .map_err(|err| { + Error::other(format!("decommission metadata save failed for pool {idx} bucket {}: {err}", bucket.name)) + })?; + let mut pool_meta = self.pool_meta.write().await; + ensure_pool_meta_write_fence(&pool_meta_guard, "decommission bucket completion save failed")?; + pool_meta.mark_decommission_progress_saved(); + ensure_pool_meta_write_fence(&pool_meta_guard, "decommission bucket completion save failed")?; + drop(pool_meta); + save_guard.restore_writes(was_write_blocked); + Ok(true) } async fn save_current_pool_meta_for_decommission_start( @@ -4597,7 +4947,8 @@ impl ECStore { space_infos: Vec<(usize, PoolSpaceInfo)>, decom_buckets: Vec, ) -> Result { - let _save_guard = self.pool_meta_save_gate.lock().await; + let mut save_guard = self.pool_meta_save_gate.lock().await; + save_guard.ensure_write_safe("decommission start failed")?; let rebalance_pool = self .pools .first() @@ -4633,8 +4984,11 @@ impl ECStore { pool_meta.clone() }; let mut latest_pool_meta = PoolMeta::default(); - let replica_state = latest_pool_meta.load_no_lock_from_replicas(self.pools.clone()).await?; - replica_state.ensure_write_safe("decommission start failed")?; + let replica_state = latest_pool_meta + .load_no_lock_from_replicas_observing(self.pools.clone(), &mut save_guard) + .await?; + save_guard.observe_replicas(replica_state); + save_guard.ensure_write_safe("decommission start failed")?; if latest_pool_meta.pools.is_empty() { latest_pool_meta = current_pool_meta; } @@ -4653,17 +5007,43 @@ impl ECStore { } activation_fence.ensure_held()?; - latest_pool_meta - .save_no_lock_with_activation_fence(self.pools.clone(), activation_fence) + let was_write_blocked = latest_pool_meta + .save_no_lock_with_activation_fence(self.pools.clone(), &mut save_guard, activation_fence) .await?; { let mut pool_meta = self.pool_meta.write().await; - *pool_meta = latest_pool_meta; + publish_pool_meta_updates(&mut pool_meta, &latest_pool_meta, indices); } + save_guard.restore_writes(was_write_blocked); Ok(previous_pool_meta) } + async fn rollback_decommission_start_after_reload_failure( + &self, + movement_gate: &Arc>, + previous_pool_meta: &PoolMeta, + indices: &[usize], + ) -> Result<()> { + let _movement_guard = movement_gate.write().await; + let mut save_guard = self.pool_meta_save_gate.lock().await; + let (pool_meta_guard, mut snapshot) = self + .acquire_pool_meta_write_guard(&mut save_guard, "decommission start rollback failed") + .await?; + rollback_start_decommission_pool_meta(&mut snapshot, previous_pool_meta, indices); + let was_write_blocked = snapshot + .save_no_lock_armed(self.pools.clone(), &mut save_guard, pool_meta_guard.lock_lost_signal()) + .await?; + let mut pool_meta = self.pool_meta.write().await; + ensure_pool_meta_write_fence(&pool_meta_guard, "decommission start rollback failed")?; + publish_pool_meta_updates(&mut pool_meta, &snapshot, indices); + ensure_pool_meta_write_fence(&pool_meta_guard, "decommission start rollback failed")?; + drop(pool_meta); + save_guard.restore_writes(was_write_blocked); + self.ctx.advance_data_movement_operation_epoch(); + Ok(()) + } + async fn ensure_decommission_rebalance_idle_after_refresh(&self) -> Result<()> { self.load_rebalance_meta().await?; ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await) @@ -4685,15 +5065,9 @@ impl ECStore { #[tracing::instrument(skip_all)] pub async fn refresh_pool_status_meta(&self) -> Result<()> { - let pool = self - .pools - .first() - .cloned() - .ok_or_else(|| Error::other("refresh_pool_status_meta: no pools available"))?; let movement_gate = self.ctx.data_movement_operation_gate(); let _movement_guard = movement_gate.write().await; - let mut persisted = PoolMeta::default(); - persisted.load(pool, self.pools.clone()).await?; + let persisted = self.load_runtime_pool_meta("refresh pool status metadata failed").await?; let active_workers = { let cancelers = self.decommission_cancelers.read().await; @@ -4750,6 +5124,7 @@ impl ECStore { self.decommission_cancel_with_owner(idx, Some(owner)).await } + #[cfg(test)] async fn decommission_cancel_with_owner_and_save( self: &Arc, idx: usize, @@ -4757,14 +5132,14 @@ impl ECStore { save_pool_meta: Save, ) -> Result<()> where - Save: FnOnce(PoolMeta) -> SaveFuture + Send + 'static, + Save: FnOnce(PoolMeta, Option>) -> SaveFuture + Send + 'static, SaveFuture: Future> + Send + 'static, { let store = self.clone(); let owner = owner.cloned(); // Dropping the RPC waiter detaches this task; the transaction retains // the store and exact owner until persistence is resolved. - tokio::spawn(async move { store.decommission_cancel_transaction(idx, owner, save_pool_meta).await }) + tokio::spawn(async move { store.decommission_cancel_transaction(idx, owner, false, save_pool_meta).await }) .await .map_err(|err| Error::other(format!("decommission cancel transaction task join error: {err}")))? } @@ -4773,20 +5148,34 @@ impl ECStore { &self, idx: usize, owner: Option, + acquire_runtime_fence: bool, save_pool_meta: Save, ) -> Result<()> where - Save: FnOnce(PoolMeta) -> SaveFuture, + Save: FnOnce(PoolMeta, Option>) -> SaveFuture, SaveFuture: Future>, { let owner = owner.as_ref(); ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?; let _start_guard = self.start_gate.lock().await; - let save_guard = self.pool_meta_save_gate.lock().await; + let mut save_guard = self.pool_meta_save_gate.lock().await; + let (_pool_meta_guard, mut persisted_pool_meta) = if acquire_runtime_fence { + let (guard, pool_meta) = self + .acquire_pool_meta_write_guard(&mut save_guard, "decommission cancel failed") + .await?; + (Some(guard), Some(pool_meta)) + } else { + save_guard.ensure_write_safe("decommission cancel failed")?; + (None, None) + }; + let pool_meta_fence = _pool_meta_guard + .as_ref() + .and_then(rustfs_lock::NamespaceLockGuard::lock_lost_signal); - // Lock order: start gate, save gate, decommission_cancelers, then - // pool_meta. The state guards stay held across persistence so the - // active generation cannot change before the cancel is published. + // Lock order: start gate, save gate, distributed pool metadata fence, + // decommission_cancelers, then pool_meta. The state guards stay held + // across persistence so the active generation cannot change before + // the cancel is published. let mut cancelers = self.decommission_cancelers.write().await; let mut pool_meta = self.pool_meta.write().await; let (pending, should_reload_pool_meta, already_canceled, terminal_canceler) = { @@ -4828,6 +5217,10 @@ impl ECStore { .get(idx) .cloned() .ok_or_else(|| invalid_decommission_pool_index_error(pool_meta.pools.len(), idx))?; + if let Some(persisted) = persisted_pool_meta.as_mut() { + merge_pool_meta_updates_for_save(persisted, &snapshot, &[idx], "decommission cancel failed")?; + snapshot = persisted.clone(); + } Some(( snapshot, DecommissionCancelCommit { @@ -4867,18 +5260,34 @@ impl ECStore { let changed = pending.is_some(); let commit_result = if let Some((snapshot, commit)) = pending { - save_pool_meta(snapshot).await?; + if let Err(err) = save_pool_meta(snapshot, pool_meta_fence).await { + save_guard.block_writes(); + return Err(err); + } + if let Some(pool_meta_guard) = _pool_meta_guard.as_ref() + && let Err(err) = ensure_pool_meta_write_fence(pool_meta_guard, "decommission cancel failed") + { + save_guard.block_writes(); + return Err(err); + } commit_decommission_cancel(&mut pool_meta, idx, commit) } else { Ok(()) }; + commit_result?; + if let Some(pool_meta_guard) = _pool_meta_guard.as_ref() + && let Err(err) = ensure_pool_meta_write_fence(pool_meta_guard, "decommission cancel failed") + { + save_guard.block_writes(); + return Err(err); + } if let Some(canceler) = terminal_canceler.as_ref() { take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, canceler); } - commit_result?; drop(pool_meta); drop(cancelers); + drop(_pool_meta_guard); drop(save_guard); if changed { @@ -4944,6 +5353,14 @@ impl ECStore { let mut attempt = 0usize; loop { + if self + .ensure_pool_meta_side_effects_safe("decommission cancel retry paused because pool metadata requires recovery") + .await + .is_err() + { + self.release_decommission_canceler_slot(idx, owner).await; + return; + } let Err(err) = self.decommission_cancel_for_operation(idx, owner).await else { return; }; @@ -4969,6 +5386,14 @@ impl ECStore { async fn retry_decommission_failed_for_operation(&self, idx: usize, owner: &DecommissionCanceler) { let mut attempt = 0usize; loop { + if self + .ensure_pool_meta_side_effects_safe("decommission failure retry paused because pool metadata requires recovery") + .await + .is_err() + { + self.release_decommission_canceler_slot(idx, owner).await; + return; + } let Err(err) = self.decommission_failed_for_operation(idx, owner).await else { return; }; @@ -4993,12 +5418,51 @@ impl ECStore { async fn decommission_cancel_with_owner(self: &Arc, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { let pools = self.pools.clone(); - self.decommission_cancel_with_owner_and_save(idx, owner, move |snapshot| async move { snapshot.save(pools).await }) + let store = self.clone(); + let owner = owner.cloned(); + tokio::spawn(async move { + store + .decommission_cancel_transaction(idx, owner, true, move |snapshot, pool_meta_fence| async move { + snapshot.save_no_lock_with_fence(pools, pool_meta_fence).await + }) + .await + }) + .await + .map_err(|err| Error::other(format!("decommission cancel transaction task join error: {err}")))? + } + + #[cfg(test)] + async fn clear_decommission_with_save(self: &Arc, idx: usize, save_pool_meta: Save) -> Result<()> + where + Save: FnOnce() -> SaveFuture + Send + 'static, + SaveFuture: Future> + Send + 'static, + { + let store = self.clone(); + tokio::spawn(async move { store.clear_decommission_transaction(idx, save_pool_meta).await }) .await + .map_err(|err| Error::other(format!("clear decommission transaction task join error: {err}")))? } #[tracing::instrument(skip(self))] - pub async fn clear_decommission(&self, idx: usize) -> Result<()> { + pub async fn clear_decommission(self: &Arc, idx: usize) -> Result<()> { + let store = self.clone(); + let save_store = store.clone(); + // Dropping the RPC waiter detaches this task; once in-memory state can + // change, the transaction must persist or roll it back before ending. + tokio::spawn(async move { + store + .clear_decommission_transaction(idx, move || async move { save_store.save_current_pool_meta(&[idx]).await }) + .await + }) + .await + .map_err(|err| Error::other(format!("clear decommission transaction task join error: {err}")))? + } + + async fn clear_decommission_transaction(&self, idx: usize, save_pool_meta: Save) -> Result<()> + where + Save: FnOnce() -> SaveFuture, + SaveFuture: Future>, + { ensure_decommission_terminal_operation_supported(self.single_pool(), "clear decommission")?; let _start_guard = self.start_gate.lock().await; @@ -5037,10 +5501,10 @@ impl ECStore { (changed, changed.then_some(previous_pool_meta)) }; - if should_reload_pool_meta && let Err(err) = self.save_current_pool_meta().await { + if should_reload_pool_meta && let Err(err) = save_pool_meta().await { if let Some(previous_pool_meta) = previous_pool_meta { let mut pool_meta = self.pool_meta.write().await; - rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta); + rollback_decommission_pool_meta(&mut pool_meta, &previous_pool_meta, &[idx]); } return Err(err); } @@ -5076,6 +5540,10 @@ impl ECStore { let _start_guard = self.start_gate.lock().await; let movement_gate = self.ctx.data_movement_operation_gate(); let _movement_guard = movement_gate.write().await; + let mut save_guard = self.pool_meta_save_gate.lock().await; + let (pool_meta_guard, mut snapshot) = self + .acquire_pool_meta_write_guard(&mut save_guard, "decommission promotion failed") + .await?; let mut pool_meta = self.pool_meta.write().await; if pool_meta.pools.get(idx).is_none() { return Err(Error::other("failed to start decommission: target pool was not found")); @@ -5083,14 +5551,27 @@ impl ECStore { let reconciled = reconcile_decommission_meta_buckets(&mut pool_meta, idx); let promoted = pool_meta.promote_queued_decommission(idx); let changed = reconciled || promoted; + if changed { + merge_pool_meta_updates_for_save(&mut snapshot, &pool_meta, &[idx], "decommission promotion failed")?; + } drop(pool_meta); - let save_error = if changed { - self.save_current_pool_meta().await.err() + let (was_write_blocked, save_error) = if changed { + match snapshot + .save_no_lock_armed(self.pools.clone(), &mut save_guard, pool_meta_guard.lock_lost_signal()) + .await + { + Ok(was_write_blocked) => (Some(was_write_blocked), None), + Err(err) => (None, Some(err)), + } } else { - None + (None, None) }; let generation = self.active_decommission_generation(idx).await?; + ensure_pool_meta_write_fence(&pool_meta_guard, "decommission promotion failed")?; + if let Some(was_write_blocked) = was_write_blocked { + save_guard.restore_writes(was_write_blocked); + } (changed, generation, save_error) }; @@ -5138,7 +5619,7 @@ impl ECStore { }; if changed { - self.save_current_pool_meta().await?; + self.save_current_pool_meta(&[idx]).await?; } Ok(()) @@ -5205,6 +5686,8 @@ impl ECStore { } let _start_guard = self.start_gate.lock().await; + let save_guard = self.pool_meta_save_gate.lock().await; + save_guard.ensure_write_safe("decommission cannot be scheduled while pool metadata requires recovery")?; let indices = { let pool_meta = self.pool_meta.read().await; resumable_decommission_queue_indices(&pool_meta) @@ -5820,7 +6303,7 @@ impl ECStore { let mut capacity_failure = false; for _ in 0..3 { match classify_decommission_free_version_attempt( - run_decommission_side_effect(&rx, &operation_gate, || async { + self.run_guarded_decommission_side_effect(&rx, &operation_gate, || async { self.decommission_tiered_object( bucket.as_str(), &version.name, @@ -5921,21 +6404,22 @@ impl ECStore { continue; } - if run_decommission_side_effect(&rx, &operation_gate, || async { - should_skip_lifecycle_for_data_movement( - self.clone(), - &bucket, - version, - lifecycle_config.as_ref(), - object_lock_config.as_ref(), - true, - &LcEventSrc::Decom, - None, - ) + if self + .run_guarded_decommission_side_effect(&rx, &operation_gate, || async { + should_skip_lifecycle_for_data_movement( + self.clone(), + &bucket, + version, + lifecycle_config.as_ref(), + object_lock_config.as_ref(), + true, + &LcEventSrc::Decom, + None, + ) + .await + }) .await - }) - .await - .map_err(|err| with_decommission_entry_context("lifecycle_expiry", bucket.as_str(), version.name.as_str(), err))? + .map_err(|err| with_decommission_entry_context("lifecycle_expiry", bucket.as_str(), version.name.as_str(), err))? { expired += 1; cleanup_preflight_allowed_missing.push(data_movement::source_cleanup_version_identity(version)); @@ -5967,15 +6451,16 @@ impl ECStore { let mut error = None; if version.deleted { for version_attempt in 1..=DECOMMISSION_VERSION_COPY_ATTEMPTS { - let result = run_decommission_side_effect(&rx, &operation_gate, || async { - self.delete_object( - bucket.as_str(), - &version.name, - decommission_delete_marker_opts(version, version_id.clone(), idx, expected_bucket_incarnation_id), - ) - .await - }) - .await; + let result = self + .run_guarded_decommission_side_effect(&rx, &operation_gate, || async { + self.delete_object( + bucket.as_str(), + &version.name, + decommission_delete_marker_opts(version, version_id.clone(), idx, expected_bucket_incarnation_id), + ) + .await + }) + .await; #[cfg(test)] let result = decommission_test_wrap_result( "delete_marker_copy", @@ -6113,16 +6598,22 @@ impl ECStore { for version_attempt in 1..=DECOMMISSION_VERSION_COPY_ATTEMPTS { if version.is_remote() { - let result = run_decommission_side_effect(&rx, &operation_gate, || async { - self.decommission_tiered_object( - bucket.as_str(), - &version.name, - version, - &decommission_remote_tiered_opts(version, version_id.clone(), idx, expected_bucket_incarnation_id), - ) - .await - }) - .await; + let result = self + .run_guarded_decommission_side_effect(&rx, &operation_gate, || async { + self.decommission_tiered_object( + bucket.as_str(), + &version.name, + version, + &decommission_remote_tiered_opts( + version, + version_id.clone(), + idx, + expected_bucket_incarnation_id, + ), + ) + .await + }) + .await; #[cfg(test)] let result = decommission_test_wrap_result( "decommission_tiered_object", @@ -6279,12 +6770,13 @@ impl ECStore { ) .await?; - let migrate_result = run_decommission_side_effect(&rx, &operation_gate, || async { - self.clone() - .decommission_object(idx, bucket, rd, expected_bucket_incarnation_id) - .await - }) - .await; + let migrate_result = self + .run_guarded_decommission_side_effect(&rx, &operation_gate, || async { + self.clone() + .decommission_object(idx, bucket, rd, expected_bucket_incarnation_id) + .await + }) + .await; #[cfg(test)] let migrate_result = decommission_test_wrap_result( DECOMMISSION_STAGE_MIGRATE_OBJECT, @@ -6458,26 +6950,27 @@ impl ECStore { let source_cleanup_mutation_fence = self .acquire_decommission_source_cleanup_fence(bucket.as_str(), entry.name.as_str(), set.as_ref()) .await?; - let cleanup_result = run_decommission_side_effect(&rx, &operation_gate, || async { - data_movement::cleanup_source_entry_if_unchanged( - set.clone(), - bucket.as_str(), - entry.name.as_str(), - &fivs, - &cleanup_preflight_allowed_missing, - data_movement::SourceCleanupBucketFence { - expected_incarnation_id: expected_bucket_incarnation_id, - lifecycle_guard: bucket_incarnation_fence - .as_ref() - .and_then(|guard| guard.namespace_lock_guard()), - namespace_lock_lost_signal: None, - object_mutation_fence: Some(&source_cleanup_mutation_fence), - }, - "decommission", - ) - .await - }) - .await; + let cleanup_result = self + .run_guarded_decommission_side_effect(&rx, &operation_gate, || async { + data_movement::cleanup_source_entry_if_unchanged( + set.clone(), + bucket.as_str(), + entry.name.as_str(), + &fivs, + &cleanup_preflight_allowed_missing, + data_movement::SourceCleanupBucketFence { + expected_incarnation_id: expected_bucket_incarnation_id, + lifecycle_guard: bucket_incarnation_fence + .as_ref() + .and_then(|guard| guard.namespace_lock_guard()), + namespace_lock_lost_signal: None, + object_mutation_fence: Some(&source_cleanup_mutation_fence), + }, + "decommission", + ) + .await + }) + .await; match cleanup_result { Ok(_) => {} Err(data_movement::SourceCleanupError::Storage(err)) => { @@ -6884,6 +7377,8 @@ impl ECStore { canceler: &DecommissionCanceler, entry_budget: Arc, ) -> Result<()> { + self.ensure_pool_meta_side_effects_safe("decommission cannot run while pool metadata requires recovery") + .await?; let generation = match self.promote_queued_decommission(idx, canceler).await { Ok(generation) => generation, Err(Error::OperationCanceled) => return Ok(()), @@ -7103,7 +7598,7 @@ impl ECStore { } async fn decommission_failed_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { - self.decommission_failed_with_owner_and_save(idx, owner, self.save_current_pool_meta()) + self.decommission_failed_with_owner_and_save(idx, owner, self.save_current_pool_meta(&[idx])) .await } @@ -7146,7 +7641,7 @@ impl ECStore { if should_reload_pool_meta && let Err(err) = save_pool_meta.await { if let Some(previous_pool_meta) = previous_pool_meta { let mut pool_meta = self.pool_meta.write().await; - rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta); + rollback_decommission_pool_meta(&mut pool_meta, &previous_pool_meta, &[idx]); } return Err(err); } @@ -7272,10 +7767,10 @@ impl ECStore { (changed, completed, changed.then_some(previous_pool_meta), terminal_canceler) }; - if should_reload_pool_meta && let Err(err) = self.save_current_pool_meta().await { + if should_reload_pool_meta && let Err(err) = self.save_current_pool_meta(&[idx]).await { if let Some(previous_pool_meta) = previous_pool_meta { let mut pool_meta = self.pool_meta.write().await; - rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta); + rollback_decommission_pool_meta(&mut pool_meta, &previous_pool_meta, &[idx]); } return Err(err); } @@ -7359,17 +7854,7 @@ impl ECStore { if is_decommissioned { warn!("decommission: already done, moving on {}", bucket.to_string()); - let bucket_done = { - let mut pool_meta = self.pool_meta.write().await; - mark_decommission_bucket_done(&mut pool_meta, idx, &bucket)? - }; - if bucket_done { - resolve_decommission_bucket_done_save_result(self.save_current_pool_meta().await, idx, bucket.name.as_str())?; - { - let mut pool_meta = self.pool_meta.write().await; - pool_meta.mark_decommission_progress_saved(); - } - } + self.mark_decommission_bucket_done_and_save(idx, &bucket).await?; return Ok(()); } @@ -7390,15 +7875,7 @@ impl ECStore { return Err(err); } - let bucket_done = { - let mut pool_meta = self.pool_meta.write().await; - mark_decommission_bucket_done(&mut pool_meta, idx, &bucket)? - }; - if bucket_done { - resolve_decommission_bucket_done_save_result(self.save_current_pool_meta().await, idx, bucket.name.as_str())?; - let mut pool_meta = self.pool_meta.write().await; - pool_meta.mark_decommission_progress_saved(); - } + self.mark_decommission_bucket_done_and_save(idx, &bucket).await?; warn!("decommission: decommission_pool bucket_done {}", &bucket.name); @@ -7550,19 +8027,9 @@ impl ECStore { "Decommission start failed after pool metadata save" ); - let rollback_result = { - let movement_guard = movement_gate.write().await; - { - let mut pool_meta = self.pool_meta.write().await; - rollback_start_decommission_pool_meta(&mut pool_meta, previous_pool_meta.clone()); - } - let rollback_result = self.save_current_pool_meta().await; - if rollback_result.is_ok() { - self.ctx.advance_data_movement_operation_epoch(); - } - drop(movement_guard); - rollback_result - }; + let rollback_result = self + .rollback_decommission_start_after_reload_failure(&movement_gate, &previous_pool_meta, &indices) + .await; if let Err(rollback_save_err) = rollback_result { error!( event = EVENT_DECOMMISSION_STATE, @@ -8568,7 +9035,10 @@ impl ECStore { ) -> Result> { self.ensure_decommission_generation_current(idx, generation).await?; let operation_gate = self.ctx.data_movement_operation_gate(); - run_decommission_side_effect(rx, &operation_gate, || self.check_after_decommission_unfenced(idx, generation)).await + self.run_guarded_decommission_side_effect(rx, &operation_gate, || { + self.check_after_decommission_unfenced(idx, generation) + }) + .await } async fn check_after_decommission_unfenced( @@ -9375,6 +9845,65 @@ mod tests { assert!(!selection.replica_state.repair_write_safe); } + #[test] + fn pool_meta_write_state_remains_blocked_after_unreadable_replica() { + let mut write_state = PoolMetaWriteState::default(); + write_state.observe_replicas(PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + write_state.observe_replicas(PoolMetaReplicaState { + needs_repair: false, + repair_write_safe: true, + }); + + let err = write_state + .ensure_write_safe("pool metadata save failed") + .expect_err("a later clean read must not clear the startup write block"); + assert!( + err.to_string() + .contains("restart after all replicas are readable and consistent") + ); + } + + #[test] + fn pool_meta_write_state_blocks_when_selection_has_no_valid_replica() { + let replicas = vec![ + PoolMetaReplica::Unreadable("pool 0 read quorum unavailable".to_string()), + PoolMetaReplica::Unreadable("pool 1 read quorum unavailable".to_string()), + ]; + let mut write_state = PoolMetaWriteState::default(); + + select_pool_meta_replicas_observing(&mut write_state, replicas).expect_err("all unreadable replicas must fail selection"); + + let err = write_state + .ensure_write_safe("pool metadata save failed") + .expect_err("an all-unreadable runtime read must latch the write block"); + assert!( + err.to_string() + .contains("restart after all replicas are readable and consistent") + ); + } + + #[test] + fn pool_meta_write_state_blocks_on_any_recovery_required_selection() { + fn assert_selection_blocks(replicas: Vec) { + let mut write_state = PoolMetaWriteState::default(); + select_pool_meta_replicas_observing(&mut write_state, replicas) + .expect_err("recovery-required replicas must fail selection"); + write_state + .ensure_write_safe("pool metadata save failed") + .expect_err("a recovery-required selection must latch the write block"); + } + + assert_selection_blocks(vec![PoolMetaReplica::Corrupt("truncated".to_string())]); + assert_selection_blocks(vec![PoolMetaReplica::Incompatible("future format".to_string())]); + assert_selection_blocks(vec![ + decode_pool_meta_replica(pool_meta_replica_test_data("pool-old")), + decode_pool_meta_replica(pool_meta_replica_test_data("pool-new")), + ]); + } + #[test] fn pool_meta_replica_selection_rejects_same_version_tuple_extension() { #[derive(Serialize)] @@ -10581,14 +11110,15 @@ mod pools_tests { ensure_decommission_start_pool_states, ensure_decommission_start_rebalance_meta_allowed, ensure_decommission_start_target_capacity, ensure_decommission_terminal_operation_supported, ensure_decommission_unresolved_verification_disk_count, ensure_local_decommission_pool_leaders, - ensure_valid_decommission_pool_index, get_by_index, guard_decommission_cancelers, has_active_decommission_canceler, - is_decommission_active, is_decommission_cancel_requested, load_decommission_entry_versions, - local_decommission_queue_prefix, mark_decommission_bucket_done, merge_decommission_durable_ilm_receipts, - merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result, - pool_meta_has_active_decommission, reconcile_decommission_meta_buckets, + ensure_pool_meta_write_fence, ensure_valid_decommission_pool_index, get_by_index, guard_decommission_cancelers, + has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested, + load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done, + merge_decommission_durable_ilm_receipts, merge_pool_meta_updates_for_save, merge_pool_status_refresh, + missing_decommission_worker_prefix, observe_decommission_terminal_reload_result, pool_meta_has_active_decommission, + publish_pool_meta_updates, reconcile_decommission_meta_buckets, reconcile_decommission_unresolved_entries_for_completion, record_decommission_unresolved_entry, - require_decommission_store, reserve_decommission_start_cancelers, resolve_decommission_bucket_done_save_result, - resolve_decommission_bucket_state, resolve_decommission_check_after_list_result, + require_decommission_store, + reserve_decommission_start_cancelers, resolve_decommission_bucket_state, resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result, resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result, resolve_decommission_pool_meta_reload_result, @@ -10620,10 +11150,12 @@ mod pools_tests { use crate::runtime::instance::InstanceContext; use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; use crate::storage_api_contracts::bucket::MakeBucketOptions; + use crate::storage_api_contracts::{object::ObjectIO, range::HTTPRangeSpec}; use crate::store::ECStore; use byteorder::{ByteOrder, LittleEndian}; use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo}; use rustfs_filemeta::{MetaCacheEntries, MetadataResolutionParams}; + use rustfs_lock::{GlobalLockManager, LocalClient, LockRequest, LockType, NamespaceLock, ObjectKey}; use rustfs_rio::Index; use std::future::Future; use std::sync::{ @@ -10680,12 +11212,372 @@ mod pools_tests { rebalance_meta: tokio::sync::RwLock::new(None), decommission_cancelers: tokio::sync::RwLock::new(cancelers), start_gate: tokio::sync::Mutex::new(()), - pool_meta_save_gate: tokio::sync::Mutex::new(()), + pool_meta_save_gate: tokio::sync::Mutex::default(), ctx, bucket_fence_registry: Arc::default(), }) } + #[derive(Debug)] + struct PartialPoolMetaWriteStorage { + fail_write: bool, + pending_write: bool, + write_started: tokio::sync::Notify, + wrote: AtomicBool, + } + + #[async_trait::async_trait] + impl ObjectIO for PartialPoolMetaWriteStorage { + type Error = Error; + type RangeSpec = HTTPRangeSpec; + type HeaderMap = http::HeaderMap; + type ObjectOptions = crate::object_api::ObjectOptions; + type ObjectInfo = crate::object_api::ObjectInfo; + type GetObjectReader = crate::object_api::GetObjectReader; + type PutObjectReader = crate::object_api::PutObjReader; + + async fn get_object_reader( + &self, + _bucket: &str, + _object: &str, + _range: Option, + _h: Self::HeaderMap, + _opts: &Self::ObjectOptions, + ) -> std::result::Result { + Err(Error::FileNotFound) + } + + async fn put_object( + &self, + _bucket: &str, + _object: &str, + _data: &mut Self::PutObjectReader, + _opts: &Self::ObjectOptions, + ) -> std::result::Result { + self.write_started.notify_one(); + if self.pending_write { + std::future::pending().await + } + if self.fail_write { + return Err(Error::Timeout); + } + self.wrote.store(true, Ordering::SeqCst); + Ok(crate::object_api::ObjectInfo::default()) + } + } + + #[tokio::test] + async fn test_partial_pool_meta_save_failure_blocks_following_side_effect() { + let store = decommission_worker_test_store(PoolMeta::default(), Vec::new()); + let committed = Arc::new(PartialPoolMetaWriteStorage { + fail_write: false, + pending_write: false, + write_started: tokio::sync::Notify::new(), + wrote: AtomicBool::new(false), + }); + let failed = Arc::new(PartialPoolMetaWriteStorage { + fail_write: true, + pending_write: false, + write_started: tokio::sync::Notify::new(), + wrote: AtomicBool::new(false), + }); + let snapshot = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status(0, None)], + ..Default::default() + }; + + { + let mut save_guard = store.pool_meta_save_gate.lock().await; + snapshot + .save_no_lock_observing(vec![committed.clone(), failed], &mut save_guard) + .await + .expect_err("the second replica write should fail after the first commits"); + } + assert!(committed.wrote.load(Ordering::SeqCst)); + + let ran = Arc::new(AtomicBool::new(false)); + let ran_by_operation = ran.clone(); + let movement_gate = store.ctx.data_movement_operation_gate(); + let result: std::result::Result<(), Error> = store + .run_guarded_decommission_side_effect(&CancellationToken::new(), &movement_gate, move || async move { + ran_by_operation.store(true, Ordering::SeqCst); + Ok(()) + }) + .await; + + assert!(result.is_err(), "a partial pool metadata save must latch the sticky safety gate"); + assert!(!ran.load(Ordering::SeqCst), "the side effect must not run after a partial save"); + } + + #[tokio::test] + async fn test_cancelled_pool_meta_save_blocks_following_side_effect() { + let store = decommission_worker_test_store(PoolMeta::default(), Vec::new()); + let committed = Arc::new(PartialPoolMetaWriteStorage { + fail_write: false, + pending_write: false, + write_started: tokio::sync::Notify::new(), + wrote: AtomicBool::new(false), + }); + let pending = Arc::new(PartialPoolMetaWriteStorage { + fail_write: false, + pending_write: true, + write_started: tokio::sync::Notify::new(), + wrote: AtomicBool::new(false), + }); + let snapshot = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status(0, None)], + ..Default::default() + }; + let save_store = store.clone(); + let save_committed = committed.clone(); + let save_pending = pending.clone(); + let save_task = tokio::spawn(async move { + let mut save_guard = save_store.pool_meta_save_gate.lock().await; + snapshot + .save_no_lock_observing(vec![save_committed, save_pending], &mut save_guard) + .await + }); + + tokio::time::timeout(StdDuration::from_secs(1), pending.write_started.notified()) + .await + .expect("the second replica write should start"); + assert!(committed.wrote.load(Ordering::SeqCst)); + save_task.abort(); + assert!( + save_task.await.expect_err("the save task should be cancelled").is_cancelled(), + "the pending replica write should be aborted" + ); + + { + let save_guard = store.pool_meta_save_gate.lock().await; + save_guard + .ensure_write_safe("cancelled pool metadata save") + .expect_err("a cancelled replica update must latch the sticky safety gate"); + } + + let ran = Arc::new(AtomicBool::new(false)); + let ran_by_operation = ran.clone(); + let movement_gate = store.ctx.data_movement_operation_gate(); + let result: std::result::Result<(), Error> = store + .run_guarded_decommission_side_effect(&CancellationToken::new(), &movement_gate, move || async move { + ran_by_operation.store(true, Ordering::SeqCst); + Ok(()) + }) + .await; + + assert!(result.is_err(), "a cancelled pool metadata save must block following side effects"); + assert!(!ran.load(Ordering::SeqCst), "the side effect must not run after a cancelled save"); + } + + #[tokio::test] + async fn test_cancelled_pool_meta_publish_keeps_write_gate_blocked() { + let store = decommission_worker_test_store(PoolMeta::default(), Vec::new()); + let committed = Arc::new(PartialPoolMetaWriteStorage { + fail_write: false, + pending_write: false, + write_started: tokio::sync::Notify::new(), + wrote: AtomicBool::new(false), + }); + let snapshot = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status(0, None)], + ..Default::default() + }; + let publish_started = Arc::new(tokio::sync::Notify::new()); + let publish_release = Arc::new(tokio::sync::Notify::new()); + let save_store = store.clone(); + let save_committed = committed.clone(); + let task_publish_started = publish_started.clone(); + let task_publish_release = publish_release.clone(); + let save_task = tokio::spawn(async move { + let mut save_guard = save_store.pool_meta_save_gate.lock().await; + let was_write_blocked = snapshot + .save_no_lock_armed(vec![save_committed], &mut save_guard, None) + .await?; + task_publish_started.notify_one(); + task_publish_release.notified().await; + save_guard.restore_writes(was_write_blocked); + Ok::<(), Error>(()) + }); + + tokio::time::timeout(StdDuration::from_secs(1), publish_started.notified()) + .await + .expect("the replica save should complete before publication"); + assert!(committed.wrote.load(Ordering::SeqCst)); + save_task.abort(); + assert!( + save_task + .await + .expect_err("the publication task should be cancelled") + .is_cancelled(), + "the task should be aborted while publication is pending" + ); + + let save_guard = store.pool_meta_save_gate.lock().await; + save_guard + .ensure_write_safe("cancelled pool metadata publication") + .expect_err("cancellation after replica save but before publication must keep writes blocked"); + } + + #[tokio::test] + async fn test_lost_pool_meta_fence_rejects_replica_write() { + let client = Arc::new(LocalClient::with_manager(Arc::new(GlobalLockManager::new()))); + let lock = NamespaceLock::with_clients_and_quorum("pool-meta-fence-loss".to_string(), vec![client], 1); + let request = LockRequest::new( + ObjectKey::new(super::RUSTFS_META_BUCKET, super::POOL_META_NAME), + LockType::Exclusive, + "stale-writer", + ) + .with_acquire_timeout(StdDuration::from_secs(1)) + .with_ttl(StdDuration::from_millis(50)) + .with_refresh_interval(StdDuration::from_millis(50)); + let guard = lock + .acquire_guard(&request) + .await + .expect("pool metadata fence acquisition should not error") + .expect("the stale writer should acquire the pool metadata fence"); + tokio::time::timeout(StdDuration::from_secs(2), guard.lock_lost_notified()) + .await + .expect("the stale writer lease should expire"); + + let storage = Arc::new(PartialPoolMetaWriteStorage { + fail_write: false, + pending_write: false, + write_started: tokio::sync::Notify::new(), + wrote: AtomicBool::new(false), + }); + let snapshot = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status(0, None)], + ..Default::default() + }; + let mut write_state = super::PoolMetaWriteState::default(); + let err = snapshot + .save_no_lock_armed(vec![storage.clone()], &mut write_state, guard.lock_lost_signal()) + .await + .expect_err("a writer must not persist after losing the distributed pool metadata fence"); + + assert!(err.to_string().contains("distributed fence was lost")); + assert!(!storage.wrote.load(Ordering::SeqCst), "the stale writer must not reach replica storage"); + write_state + .ensure_write_safe("lost pool metadata fence") + .expect_err("a lost distributed fence must latch the sticky write gate"); + } + + #[tokio::test] + async fn test_pool_meta_fence_loss_after_publish_keeps_write_gate_blocked() { + let client = Arc::new(LocalClient::with_manager(Arc::new(GlobalLockManager::new()))); + let lock = NamespaceLock::with_clients_and_quorum("pool-meta-publish-fence-loss".to_string(), vec![client], 1); + let request = LockRequest::new( + ObjectKey::new(super::RUSTFS_META_BUCKET, super::POOL_META_NAME), + LockType::Exclusive, + "publishing-writer", + ) + .with_acquire_timeout(StdDuration::from_secs(1)) + .with_ttl(StdDuration::from_secs(1)) + .with_refresh_interval(StdDuration::from_secs(1)); + let guard = lock + .acquire_guard(&request) + .await + .expect("pool metadata fence acquisition should not error") + .expect("the publishing writer should acquire the pool metadata fence"); + let storage = Arc::new(PartialPoolMetaWriteStorage { + fail_write: false, + pending_write: false, + write_started: tokio::sync::Notify::new(), + wrote: AtomicBool::new(false), + }); + let mut saved = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status(0, None)], + ..Default::default() + }; + saved.pools[0].last_update = OffsetDateTime::UNIX_EPOCH + Duration::seconds(1); + let mut current = saved.clone(); + current.pools[0].last_update = OffsetDateTime::UNIX_EPOCH; + let mut write_state = super::PoolMetaWriteState::default(); + let was_write_blocked = saved + .save_no_lock_armed(vec![storage], &mut write_state, guard.lock_lost_signal()) + .await + .expect("the replica save should finish while the fence is valid"); + + ensure_pool_meta_write_fence(&guard, "test pool metadata publish") + .expect("the fence should remain valid before publication"); + publish_pool_meta_updates(&mut current, &saved, &[0]); + tokio::time::timeout(StdDuration::from_secs(2), guard.lock_lost_notified()) + .await + .expect("the fence should expire after publication"); + ensure_pool_meta_write_fence(&guard, "test pool metadata publish") + .expect_err("the post-publication fence check must observe the loss"); + + assert_eq!(current.pools[0].last_update, saved.pools[0].last_update); + assert!(!was_write_blocked); + write_state + .ensure_write_safe("lost pool metadata publish fence") + .expect_err("the sticky write gate must remain armed after post-publication fence loss"); + } + + #[tokio::test] + async fn test_clear_decommission_transaction_survives_caller_abort() { + let pool_meta = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + failed: true, + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![None]); + let save_started = Arc::new(tokio::sync::Notify::new()); + let save_release = Arc::new(tokio::sync::Notify::new()); + let save_completed = Arc::new(AtomicBool::new(false)); + let caller_store = store.clone(); + let caller_save_started = save_started.clone(); + let caller_save_release = save_release.clone(); + let caller_save_completed = save_completed.clone(); + let caller_task = tokio::spawn(async move { + caller_store + .clear_decommission_with_save(0, move || async move { + caller_save_started.notify_one(); + caller_save_release.notified().await; + caller_save_completed.store(true, Ordering::SeqCst); + Ok(()) + }) + .await + }); + + tokio::time::timeout(StdDuration::from_secs(1), save_started.notified()) + .await + .expect("the clear transaction should reach persistence"); + caller_task.abort(); + assert!( + caller_task + .await + .expect_err("the RPC waiter should be cancelled") + .is_cancelled(), + "the clear caller should be aborted while persistence is pending" + ); + save_release.notify_one(); + + let _start_guard = tokio::time::timeout(StdDuration::from_secs(1), store.start_gate.lock()) + .await + .expect("the detached clear transaction should finish"); + assert!( + save_completed.load(Ordering::SeqCst), + "the detached transaction must finish persistence after the caller is aborted" + ); + let pool_meta = store.pool_meta.read().await; + assert!( + pool_meta.pools[0].decommission.is_none(), + "the detached transaction should publish the persisted clear" + ); + } + fn decommission_test_pool_endpoint(idx: usize, is_local: bool) -> PoolEndpoints { let port = 9000usize + idx; let mut endpoint = @@ -11085,6 +11977,82 @@ mod pools_tests { assert!(current.pools[0].decommission.is_none()); } + #[test] + fn test_pool_meta_save_merge_preserves_newer_untouched_pool() { + let older = OffsetDateTime::from_unix_timestamp(1_000).expect("test timestamp should be valid"); + let newer = OffsetDateTime::from_unix_timestamp(2_000).expect("test timestamp should be valid"); + let mut current = PoolMeta { + pools: vec![ + PoolStatus { + id: 0, + cmd_line: "pool-0".to_string(), + last_update: older, + decommission: Some(PoolDecommissionInfo { + complete: true, + ..Default::default() + }), + }, + PoolStatus { + id: 1, + cmd_line: "pool-1".to_string(), + last_update: older, + decommission: None, + }, + ], + ..Default::default() + }; + let mut persisted = current.clone(); + persisted.pools[1].last_update = newer; + persisted.pools[1].decommission = Some(PoolDecommissionInfo { + failed: true, + ..Default::default() + }); + + assert!(current.clear_decommission(0).expect("terminal decommission should clear")); + merge_pool_meta_updates_for_save(&mut persisted, ¤t, &[0], "clear decommission") + .expect("the target pool update should merge into the latest snapshot"); + + assert!(persisted.pools[0].decommission.is_none()); + assert_eq!(persisted.pools[1].last_update, newer); + assert!(persisted.pools[1].decommission.as_ref().is_some_and(|info| info.failed)); + } + + #[test] + fn test_pool_meta_publish_preserves_untouched_runtime_progress() { + let mut current = PoolMeta { + pools: vec![ + decommission_test_pool_status(0, None), + decommission_test_pool_status( + 1, + Some(PoolDecommissionInfo { + items_decommissioned: 10, + bytes_done: 1_024, + ..Default::default() + }), + ), + ], + ..Default::default() + }; + let mut saved = current.clone(); + saved.pools[0].decommission = Some(PoolDecommissionInfo { + complete: true, + ..Default::default() + }); + let saved_pool_1 = saved.pools[1].decommission.as_mut().expect("pool 1 progress should exist"); + saved_pool_1.items_decommissioned = 1; + saved_pool_1.bytes_done = 128; + + publish_pool_meta_updates(&mut current, &saved, &[0]); + + assert!(current.pools[0].decommission.as_ref().is_some_and(|info| info.complete)); + let pool_1 = current.pools[1] + .decommission + .as_ref() + .expect("pool 1 progress should remain present"); + assert_eq!(pool_1.items_decommissioned, 10); + assert_eq!(pool_1.bytes_done, 1_024); + } + #[test] fn test_dedup_indices_removes_duplicates_preserving_order() { assert_eq!(dedup_indices(&[0, 2, 1, 2, 3, 0]), vec![0, 2, 1, 3]); @@ -11511,6 +12479,76 @@ mod pools_tests { assert!(!called.load(Ordering::SeqCst)); } + #[tokio::test] + async fn test_decommission_side_effect_stops_after_pool_meta_write_block() { + let store = decommission_worker_test_store(PoolMeta::default(), Vec::new()); + store + .pool_meta_save_gate + .lock() + .await + .observe_replicas(super::PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + let called = Arc::new(AtomicBool::new(false)); + let operation_gate = store.ctx.data_movement_operation_gate(); + + let result = store + .run_guarded_decommission_side_effect(&CancellationToken::new(), &operation_gate, { + let called = called.clone(); + move || async move { + called.store(true, Ordering::SeqCst); + Ok::<_, Error>(()) + } + }) + .await; + + assert!( + result + .expect_err("sticky pool metadata state must block new movement") + .to_string() + .contains("restart after all replicas are readable and consistent") + ); + assert!(!called.load(Ordering::SeqCst)); + } + + #[tokio::test] + async fn test_decommission_reservation_stops_before_canceler_slot_when_pool_meta_is_blocked() { + let generation = OffsetDateTime::UNIX_EPOCH; + let pool_meta = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![None]); + store + .pool_meta_save_gate + .lock() + .await + .observe_replicas(super::PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + + let err = store + .reserve_decommission_routines(&CancellationToken::new(), &[0]) + .await + .err() + .expect("sticky pool metadata state must block worker reservation"); + + assert!( + err.to_string() + .contains("restart after all replicas are readable and consistent") + ); + assert!(store.decommission_cancelers.read().await[0].is_none()); + } + #[tokio::test] async fn test_decommission_transition_waits_without_registered_canceler() { let store = decommission_worker_test_store(PoolMeta::default(), vec![None]); @@ -11637,10 +12675,13 @@ mod pools_tests { ..Default::default() }; let mut active = previous.clone(); + let active_update = OffsetDateTime::UNIX_EPOCH + Duration::seconds(1); + active.pools[0].last_update = active_update; active.pools[0].decommission = Some(PoolDecommissionInfo { start_time: Some(OffsetDateTime::UNIX_EPOCH), ..Default::default() }); + let mut peer = active.clone(); assert!(active.is_suspended(0)); assert_eq!( @@ -11648,10 +12689,13 @@ mod pools_tests { DecommissionStartPoolState::Decommissioning ); - rollback_start_decommission_pool_meta(&mut active, previous); + rollback_start_decommission_pool_meta(&mut active, &previous, &[0]); assert!(!active.is_suspended(0)); assert_eq!(decommission_start_pool_state(active.pools.first()), DecommissionStartPoolState::Active); + assert!(active.pools[0].last_update > active_update); + assert!(merge_pool_status_refresh(&mut peer, active, &[false])); + assert!(peer.pools[0].decommission.is_none()); } #[test] @@ -11971,21 +13015,6 @@ mod pools_tests { ); } - #[test] - fn test_resolve_decommission_bucket_done_save_result_passthrough_ok() { - assert!(resolve_decommission_bucket_done_save_result(Ok(()), 1, "bucket-a").is_ok()); - } - - #[test] - fn test_resolve_decommission_bucket_done_save_result_wraps_error_context() { - let err = resolve_decommission_bucket_done_save_result(Err(Error::SlowDown), 2, "bucket-a") - .expect_err("metadata save failure should carry pool/bucket context"); - assert!( - err.to_string() - .contains("decommission metadata save failed for pool 2 bucket bucket-a") - ); - } - #[test] fn test_resolve_decommission_optional_bucket_config_result_passthrough() { let result = resolve_decommission_optional_bucket_config_result("bucket-a", "replication", Ok(42_u8)) @@ -14772,7 +15801,148 @@ mod pools_tests { } #[tokio::test] - async fn test_decommission_cancel_save_failure_preserves_generation_and_token_until_retry() { + async fn test_decommission_supervisor_releases_slot_without_terminal_retry_when_pool_meta_is_blocked() { + let generation = OffsetDateTime::UNIX_EPOCH; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let pool_meta = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); + store + .pool_meta_save_gate + .lock() + .await + .observe_replicas(super::PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + let guards = guard_decommission_cancelers(vec![(0, canceler.clone())]); + + tokio::time::timeout( + StdDuration::from_secs(1), + spawn_decommission_index_cancelers(store.clone(), CancellationToken::new(), guards, Arc::new(Semaphore::new(1))), + ) + .await + .expect("blocked supervisor must not enter terminal retry") + .expect("blocked supervisor task should not panic"); + + assert!(store.decommission_cancelers.read().await[0].is_none()); + assert!(!canceler.is_active()); + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("blocked decommission metadata should remain active"); + assert!(!info.failed); + assert!(!info.canceled); + assert!(!info.complete); + assert_eq!(info.start_time, Some(generation)); + } + + #[tokio::test] + async fn test_decommission_promotion_rechecks_sticky_gate_after_start_wait() { + let pool_meta = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + queued: true, + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![None]); + let start_guard = store.start_gate.lock().await; + let mut promotion = tokio::spawn({ + let store = store.clone(); + async move { store.promote_queued_decommission_for_test(0).await } + }); + tokio::task::yield_now().await; + assert!(!promotion.is_finished(), "promotion should be waiting for the start gate"); + + store + .pool_meta_save_gate + .lock() + .await + .observe_replicas(super::PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + drop(start_guard); + + let err = tokio::time::timeout(StdDuration::from_secs(1), &mut promotion) + .await + .expect("blocked promotion should finish") + .expect("promotion task should not panic") + .expect_err("promotion must recheck the sticky gate after waiting"); + assert!( + err.to_string() + .contains("restart after all replicas are readable and consistent") + ); + + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("queued metadata should remain present"); + assert!(info.queued); + assert!(info.start_time.is_none()); + } + + #[tokio::test] + async fn test_decommission_bucket_done_rechecks_sticky_gate_before_mutation() { + let bucket = DecomBucketInfo { + name: "bucket-a".to_string(), + prefix: String::new(), + }; + let pool_meta = PoolMeta { + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + queued_buckets: vec![bucket.to_string()], + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![None]); + store + .pool_meta_save_gate + .lock() + .await + .observe_replicas(super::PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + + let err = store + .mark_decommission_bucket_done_and_save(0, &bucket) + .await + .expect_err("bucket completion must stop before mutating sticky pool metadata"); + assert!( + err.to_string() + .contains("restart after all replicas are readable and consistent") + ); + + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("decommission metadata should remain present"); + assert_eq!(info.queued_buckets, vec![bucket.to_string()]); + assert!(info.decommissioned_buckets.is_empty()); + } + + #[tokio::test] + async fn test_decommission_cancel_save_failure_preserves_generation_and_blocks_retry() { let generation = OffsetDateTime::UNIX_EPOCH; let canceler = DecommissionCanceler::new(CancellationToken::new()); let pool_meta = PoolMeta { @@ -14789,7 +15959,7 @@ mod pools_tests { let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); let err = store - .decommission_cancel_with_owner_and_save(0, Some(&canceler), |_| async { Err(Error::Timeout) }) + .decommission_cancel_with_owner_and_save(0, Some(&canceler), |_, _| async { Err(Error::Timeout) }) .await .expect_err("injected pool metadata timeout should fail cancel"); assert!(matches!(err, Error::Timeout)); @@ -14815,61 +15985,94 @@ mod pools_tests { assert!(!canceler.is_cancelled()); assert!(store.decommission_terminal_retryable_for_operation(0, &canceler).await); - let persisted = Arc::new(std::sync::Mutex::new(None)); - let persisted_for_save = persisted.clone(); - store - .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move { - let data = snapshot.encode_config_data()?; - *persisted_for_save - .lock() - .expect("persisted snapshot lock should not be poisoned") = Some(data); + let retry_save_called = Arc::new(AtomicBool::new(false)); + let retry_save_called_by_closure = retry_save_called.clone(); + let retry_err = store + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |_, _| async move { + retry_save_called_by_closure.store(true, Ordering::SeqCst); Ok(()) }) .await - .expect("cancel retry should commit"); + .expect_err("an ambiguous save failure must block same-process retry"); + assert!( + retry_err + .to_string() + .contains("restart after all replicas are readable and consistent") + ); + assert!(!retry_save_called.load(Ordering::SeqCst)); { let pool_meta = store.pool_meta.read().await; let info = pool_meta.pools[0] .decommission .as_ref() - .expect("committed cancel metadata should remain present"); - assert!(info.canceled); + .expect("blocked retry must retain decommission metadata"); + assert_eq!(info.start_time, Some(generation)); + assert!(!info.canceled); assert!(!info.complete); assert!(!info.failed); - assert!(info.start_time.is_none()); } - assert!(store.decommission_cancelers.read().await[0].is_none()); - assert!(!canceler.is_active()); - assert!(canceler.is_cancelled()); + assert!(store.decommission_cancelers.read().await[0].is_some()); + assert!(canceler.is_active()); + assert!(!canceler.is_cancelled()); + } - let repeated_save_called = Arc::new(AtomicBool::new(false)); - let repeated_save_called_by_closure = repeated_save_called.clone(); + #[tokio::test] + async fn test_decommission_cancel_stays_blocked_after_unreadable_pool_meta_replica() { + let generation = OffsetDateTime::UNIX_EPOCH; + let canceler = DecommissionCanceler::new(CancellationToken::new()); + let pool_meta = PoolMeta { + version: super::POOL_META_VERSION, + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + )], + ..Default::default() + }; + let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]); store - .decommission_cancel_with_owner_and_save(0, None, move |_| async move { - repeated_save_called_by_closure.store(true, Ordering::SeqCst); + .pool_meta_save_gate + .lock() + .await + .observe_replicas(super::PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + let save_called = Arc::new(AtomicBool::new(false)); + let save_called_by_closure = save_called.clone(); + + let err = store + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |_, _| async move { + save_called_by_closure.store(true, Ordering::SeqCst); Ok(()) }) .await - .expect("repeated cancel should be idempotent"); - assert!(!repeated_save_called.load(Ordering::SeqCst)); + .expect_err("cancel must remain blocked until restart after an unreadable replica"); - let persisted = persisted - .lock() - .expect("persisted snapshot lock should not be poisoned") - .take() - .expect("successful retry should capture persisted bytes"); - let mut restarted = PoolMeta::default(); - restarted - .load_from_config_data(persisted) - .expect("a restarted process should decode the committed cancel"); assert!( - restarted.pools[0] - .decommission - .as_ref() - .is_some_and(|info| info.canceled && info.start_time.is_none()) + err.to_string() + .contains("restart after all replicas are readable and consistent") + ); + assert!(!save_called.load(Ordering::SeqCst)); + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("blocked cancel must preserve decommission metadata"); + assert_eq!(info.start_time, Some(generation)); + assert!(!info.canceled); + drop(pool_meta); + assert!(canceler.is_active()); + assert!(!canceler.is_cancelled()); + let cancelers = store.decommission_cancelers.read().await; + assert!( + cancelers[0] + .as_ref() + .is_some_and(|current| current.owns_same_operation(&canceler)) ); - assert!(resumable_decommission_queue_indices(&restarted).is_empty()); } #[tokio::test] @@ -14898,7 +16101,7 @@ mod pools_tests { let save_release = save_release.clone(); async move { store - .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move { + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot, _| async move { persisted_tx .send(snapshot.encode_config_data()?) .map_err(|_| Error::other("failed to expose saved cancel snapshot"))?; @@ -14976,7 +16179,7 @@ mod pools_tests { let save_release = save_release.clone(); async move { store - .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move { + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot, _| async move { assert!( snapshot.pools[0] .decommission @@ -15080,7 +16283,7 @@ mod pools_tests { let save_entered = save_entered.clone(); async move { store - .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |_| async move { + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |_, _| async move { save_entered.store(true, Ordering::SeqCst); save_started.notify_one(); save_release.notified().await; @@ -15220,7 +16423,7 @@ mod pools_tests { let queued_replacement_for_save = queued_replacement.clone(); store - .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move { + .decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot, _| async move { let saved_cancel = snapshot.pools[0] .decommission .as_ref() diff --git a/crates/ecstore/src/data_usage/mod.rs b/crates/ecstore/src/data_usage/mod.rs index 0a613e02f..37f14c4ed 100644 --- a/crates/ecstore/src/data_usage/mod.rs +++ b/crates/ecstore/src/data_usage/mod.rs @@ -2776,7 +2776,7 @@ mod tests { rebalance_meta: RwLock::new(None), decommission_cancelers: RwLock::new(Vec::new()), start_gate: TokioMutex::new(()), - pool_meta_save_gate: TokioMutex::new(()), + pool_meta_save_gate: TokioMutex::default(), ctx, bucket_fence_registry: Arc::default(), }) diff --git a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs index 08d8f6b0c..d475a606c 100644 --- a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs +++ b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs @@ -2999,7 +2999,7 @@ fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc Result { + let movement_gate = self.ctx.data_movement_operation_gate(); + let _movement_guard = movement_gate.read().await; + let save_guard = self.pool_meta_save_gate.lock().await; + save_guard.ensure_write_safe("bucket heal cannot run while pool metadata requires recovery")?; let mut fenced_pools = BTreeSet::new(); { let pool_meta = self.pool_meta.read().await; @@ -320,9 +330,10 @@ impl ECStore { } let dispatch_fenced_pools = fenced_pools.iter().copied().collect::>(); + drop(save_guard); let mut res = self .peer_sys - .heal_bucket_with_fence(bucket, opts, &dispatch_fenced_pools) + .heal_bucket_with_fence_from_movement_guarded_coordinator(bucket, opts, &dispatch_fenced_pools) .await?; { let pool_meta = self.pool_meta.read().await; @@ -479,17 +490,113 @@ impl ECStore { mod tests { use super::*; use crate::bucket::metadata_sys; - use crate::core::pools::{PoolDecommissionInfo, PoolStatus}; + use crate::cluster::rpc::PeerS3Client; + use crate::core::pools::{PoolDecommissionInfo, PoolMetaReplicaState, PoolStatus}; + use crate::disk::error::Result as DiskResult; use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk}; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::runtime::instance::InstanceContext; use crate::services::rebalance::{RebalanceInfo, RebalanceStats}; - use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions}; + use crate::storage_api_contracts::bucket::{ + BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, + }; use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations}; use crate::store::init_format::{load_format_erasure, save_format_file}; use crate::store::init_local_disks_with_instance_ctx; use tokio_util::sync::CancellationToken; + #[derive(Debug)] + struct BlockingHealPeer { + started: Arc, + release: Arc, + } + + #[async_trait::async_trait] + impl PeerS3Client for BlockingHealPeer { + async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> DiskResult { + self.started.notify_one(); + self.release.notified().await; + Ok(HealResultItem::default()) + } + + async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> DiskResult<()> { + Ok(()) + } + + async fn list_bucket(&self, _opts: &BucketOptions) -> DiskResult> { + Ok(Vec::new()) + } + + async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> DiskResult<()> { + Ok(()) + } + + async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> DiskResult { + Ok(BucketInfo::default()) + } + + fn get_pools(&self) -> Option> { + Some(vec![0, 1]) + } + } + + #[derive(Debug)] + struct WriterQueuedLocalHealPeer { + movement_gate: Arc>, + writer_queued: Arc, + writer_acquired: Arc, + } + + #[async_trait::async_trait] + impl PeerS3Client for WriterQueuedLocalHealPeer { + async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> DiskResult { + let _movement_guard = self + .movement_gate + .try_read() + .map_err(|_| crate::error::StorageError::SlowDown)?; + Ok(HealResultItem::default()) + } + + async fn heal_bucket_with_fence_from_movement_guarded_coordinator( + &self, + _bucket: &str, + _opts: &HealOpts, + _fenced_pools: &[usize], + ) -> DiskResult { + Ok(HealResultItem::default()) + } + + async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> DiskResult<()> { + Ok(()) + } + + async fn list_bucket(&self, _opts: &BucketOptions) -> DiskResult> { + Ok(Vec::new()) + } + + async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> DiskResult<()> { + Ok(()) + } + + async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> DiskResult { + let movement_gate = self.movement_gate.clone(); + let writer_acquired = self.writer_acquired.clone(); + tokio::spawn(async move { + let _movement_guard = movement_gate.write().await; + writer_acquired.notify_one(); + }); + while self.movement_gate.try_read().is_ok() { + tokio::task::yield_now().await; + } + self.writer_queued.notify_one(); + Ok(BucketInfo::default()) + } + + fn get_pools(&self) -> Option> { + Some(vec![0, 1]) + } + } + async fn minimal_heal_pool(pool_idx: usize) -> Arc { let format = FormatV3::new(1, 1); let endpoint_url = format!("http://127.0.0.1:{}/data", 19000 + pool_idx); @@ -529,7 +636,7 @@ mod tests { rebalance_meta: RwLock::new(None), decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), - pool_meta_save_gate: Mutex::new(()), + pool_meta_save_gate: Mutex::default(), ctx: crate::runtime::instance::bootstrap_ctx(), bucket_fence_registry: std::sync::Arc::default(), } @@ -955,6 +1062,93 @@ mod tests { ); } + #[tokio::test] + async fn bucket_heal_blocks_before_dispatch_after_unreadable_pool_meta_replica() { + let store = minimal_heal_store().await; + store.pool_meta_save_gate.lock().await.observe_replicas(PoolMetaReplicaState { + needs_repair: true, + repair_write_safe: false, + }); + + let err = store + .handle_heal_bucket("bucket", &HealOpts::default()) + .await + .expect_err("bucket heal must stay blocked until restart after an unreadable replica"); + + assert!( + err.to_string() + .contains("restart after all replicas are readable and consistent") + ); + } + + #[tokio::test] + async fn bucket_heal_releases_save_gate_before_peer_dispatch_and_holds_movement_snapshot() { + let mut store = minimal_heal_store().await; + let started = Arc::new(tokio::sync::Notify::new()); + let release = Arc::new(tokio::sync::Notify::new()); + let peer: Box = Box::new(BlockingHealPeer { + started: started.clone(), + release: release.clone(), + }); + store.peer_sys.clients = vec![Arc::new(peer)]; + let store = Arc::new(store); + let movement_gate = store.ctx.data_movement_operation_gate(); + let mut heal = tokio::spawn({ + let store = store.clone(); + async move { store.handle_heal_bucket("bucket", &HealOpts::default()).await } + }); + + tokio::time::timeout(std::time::Duration::from_secs(1), started.notified()) + .await + .expect("peer dispatch should start"); + assert!( + store.pool_meta_save_gate.try_lock().is_ok(), + "coordinator must release its local save gate before waiting for peers" + ); + assert!( + movement_gate.try_write().is_err(), + "bucket heal must hold the movement snapshot through peer dispatch" + ); + + release.notify_one(); + tokio::time::timeout(std::time::Duration::from_secs(1), &mut heal) + .await + .expect("bucket heal should finish after peer release") + .expect("bucket heal task should not panic") + .expect("bucket heal should succeed"); + } + + #[tokio::test] + async fn bucket_heal_local_fanout_does_not_reenter_movement_read_behind_queued_writer() { + let mut store = minimal_heal_store().await; + let movement_gate = store.ctx.data_movement_operation_gate(); + let writer_queued = Arc::new(tokio::sync::Notify::new()); + let writer_acquired = Arc::new(tokio::sync::Notify::new()); + let peer: Box = Box::new(WriterQueuedLocalHealPeer { + movement_gate: movement_gate.clone(), + writer_queued: writer_queued.clone(), + writer_acquired: writer_acquired.clone(), + }); + store.peer_sys.clients = vec![Arc::new(peer)]; + let store = Arc::new(store); + let mut heal = tokio::spawn({ + let store = store.clone(); + async move { store.handle_heal_bucket("bucket", &HealOpts::default()).await } + }); + + tokio::time::timeout(std::time::Duration::from_secs(1), writer_queued.notified()) + .await + .expect("movement writer should queue during local peer lookup"); + tokio::time::timeout(std::time::Duration::from_secs(1), &mut heal) + .await + .expect("local fan-out must not reenter movement read behind the queued writer") + .expect("bucket heal task should not panic") + .expect("bucket heal should succeed"); + tokio::time::timeout(std::time::Duration::from_secs(1), writer_acquired.notified()) + .await + .expect("queued movement writer should proceed after bucket heal releases its read guard"); + } + #[tokio::test] #[serial_test::serial] async fn unscoped_heal_object_suspended_owner_semantics() { @@ -1282,7 +1476,7 @@ mod tests { rebalance_meta: RwLock::new(None), decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), - pool_meta_save_gate: Mutex::new(()), + pool_meta_save_gate: Mutex::default(), ctx: crate::runtime::instance::bootstrap_ctx(), bucket_fence_registry: std::sync::Arc::default(), }; diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 6cd1452ae..330e8c144 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -13,7 +13,9 @@ // limitations under the License. use super::*; -use crate::core::pools::{PoolMetaReplicaState, local_decommission_queue_prefix, pool_meta_has_active_decommission}; +use crate::core::pools::{ + PoolMetaReplicaState, PoolMetaWriteState, local_decommission_queue_prefix, pool_meta_has_active_decommission, +}; use crate::error::is_err_decommission_running; use crate::runtime::instance::InstanceContext; use crate::runtime::sources as runtime_sources; @@ -109,6 +111,14 @@ fn should_auto_start_rebalance_after_init(decommission_running: bool, rebalance_ rebalance_meta_loaded && !decommission_running } +fn should_schedule_local_decommission_resume( + pool_indices: &[usize], + pool_meta_replica_state: PoolMetaReplicaState, + pool_meta_write_safe: bool, +) -> bool { + !pool_indices.is_empty() && pool_meta_replica_state.repair_write_safe && pool_meta_write_safe +} + async fn wait_for_local_decommission_resume_delay(rx: &CancellationToken, delay: Duration) -> bool { tokio::select! { _ = rx.cancelled() => false, @@ -120,15 +130,19 @@ fn resolve_store_init_stage_result(result: Result<()>, stage: &str) -> Result<() result.map_err(|err| Error::other(format!("store init failed during {stage}: {err}"))) } -async fn load_pool_meta_for_startup(pools: Vec>) -> Result<(PoolMeta, PoolMetaReplicaState)> +async fn load_pool_meta_for_startup( + pools: Vec>, + write_state: &mut PoolMetaWriteState, +) -> Result<(PoolMeta, PoolMetaReplicaState)> where S: EcstoreObjectIO, { let mut meta = PoolMeta::default(); let replica_state = meta - .load_no_lock_from_replicas(pools) + .load_no_lock_from_replicas_observing(pools, write_state) .await .map_err(|err| Error::other(format!("store init failed during load_pool_meta: {err}")))?; + write_state.observe_replicas(replica_state); Ok((meta, replica_state)) } @@ -143,6 +157,7 @@ async fn persist_pool_meta_for_startup_if_safe( meta: &PoolMeta, pools: Vec>, replica_state: PoolMetaReplicaState, + write_state: PoolMetaWriteState, topology_update: bool, elected_writer: bool, ) -> Result<()> @@ -152,10 +167,14 @@ where if !elected_writer { return Ok(()); } + let should_write = topology_update || (replica_state.needs_repair && replica_state.repair_write_safe); if topology_update { replica_state.ensure_write_safe("store init failed during save_validated_pool_meta")?; } - if topology_update || (replica_state.needs_repair && replica_state.repair_write_safe) { + if should_write { + write_state.ensure_write_safe("store init failed during save_validated_pool_meta")?; + } + if should_write { save_validated_pool_meta_for_startup(meta, pools).await?; } Ok(()) @@ -431,7 +450,7 @@ impl ECStore { rebalance_meta: RwLock::new(None), decommission_cancelers, start_gate: Mutex::new(()), - pool_meta_save_gate: Mutex::new(()), + pool_meta_save_gate: Mutex::default(), // Adopt the caller's context (the process bootstrap one on the // legacy path) so startup writes (erasure type recorded before // this point) and later reads share one cell. @@ -475,7 +494,10 @@ impl ECStore { pub async fn init(self: &Arc, rx: CancellationToken) -> Result<()> { runtime_sources::ensure_boot_time().await; - let (meta, pool_meta_replica_state) = load_pool_meta_for_startup(self.pools.clone()).await?; + let (meta, pool_meta_replica_state) = { + let mut write_state = self.pool_meta_save_gate.lock().await; + load_pool_meta_for_startup(self.pools.clone(), &mut write_state).await? + }; let update = meta.validate(self.pools.clone())?; let endpoints = runtime_sources::endpoint_pools_or_default(); let should_persist_pool_meta = runtime_sources::first_cluster_node_is_local().await; @@ -487,14 +509,18 @@ impl ECStore { }; // Only one local node should persist validated pool metadata here; otherwise // distributed startup can race on the same lock and replay the prior init bug. - persist_pool_meta_for_startup_if_safe( - &installed_pool_meta, - self.pools.clone(), - pool_meta_replica_state, - update, - should_persist_pool_meta, - ) - .await?; + { + let write_state = self.pool_meta_save_gate.lock().await; + persist_pool_meta_for_startup_if_safe( + &installed_pool_meta, + self.pools.clone(), + pool_meta_replica_state, + *write_state, + update, + should_persist_pool_meta, + ) + .await?; + } { let mut pool_meta = self.pool_meta.write().await; @@ -538,7 +564,11 @@ impl ECStore { } let local_pool_indices = local_decommission_queue_prefix(&endpoints, &pool_indices)?; - if !local_pool_indices.is_empty() { + let pool_meta_write_safe = self + .ensure_pool_meta_side_effects_safe("decommission resume blocked while pool metadata requires recovery") + .await + .is_ok(); + if should_schedule_local_decommission_resume(&local_pool_indices, pool_meta_replica_state, pool_meta_write_safe) { let store = self.clone(); tokio::spawn(async move { @@ -547,6 +577,16 @@ impl ECStore { } resume_local_decommission_after_init(store, rx, local_pool_indices).await; }); + } else if !local_pool_indices.is_empty() { + error!( + event = EVENT_DECOMMISSION_RESUME_FAILED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_STORE_INIT, + state = "blocked", + pool_indices = ?local_pool_indices, + reason = "pool_meta_write_blocked", + "Decommission resume blocked until pool metadata replicas are readable and consistent" + ); } runtime_sources::init_bucket_monitor_for_current_endpoints(); @@ -575,11 +615,11 @@ impl ECStore { #[cfg(test)] mod tests { use super::{ - LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES, load_pool_meta_for_startup, persist_pool_meta_for_startup_if_safe, - pool_first_endpoint_is_local, pool_meta_has_active_decommission, preflight_startup_rpc_secret_with, - resolve_startup_pool_defaults_with, resolve_store_init_stage_result, save_validated_pool_meta_for_startup, - should_auto_start_rebalance_after_init, should_retry_format_load, should_retry_local_decommission_resume, - wait_for_local_decommission_resume_delay, + LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES, PoolMetaWriteState, load_pool_meta_for_startup, + persist_pool_meta_for_startup_if_safe, pool_first_endpoint_is_local, pool_meta_has_active_decommission, + preflight_startup_rpc_secret_with, resolve_startup_pool_defaults_with, resolve_store_init_stage_result, + save_validated_pool_meta_for_startup, should_auto_start_rebalance_after_init, should_retry_format_load, + should_retry_local_decommission_resume, wait_for_local_decommission_resume_delay, }; #[cfg(feature = "test-util")] use crate::disk::DiskAPI; @@ -796,8 +836,9 @@ mod tests { #[tokio::test] async fn test_store_init_pool_meta_io_bypasses_namespace_lock_surface() { let storage = Arc::new(StartupPoolMetaStorage::new(Vec::new())); + let mut write_state = PoolMetaWriteState::default(); - let (loaded, replica_state) = load_pool_meta_for_startup(vec![storage.clone()]) + let (loaded, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut write_state) .await .expect("startup pool metadata load should tolerate missing metadata without locks"); assert!(loaded.pools.is_empty()); @@ -822,8 +863,9 @@ mod tests { let corrupt = Arc::new(StartupPoolMetaStorage::new(vec![0, 1, 2])); let expected = init_test_pool_meta(None); let backup = Arc::new(StartupPoolMetaStorage::new(startup_pool_meta_payload(&expected))); + let mut write_state = PoolMetaWriteState::default(); - let (loaded, replica_state) = load_pool_meta_for_startup(vec![corrupt.clone(), backup.clone()]) + let (loaded, replica_state) = load_pool_meta_for_startup(vec![corrupt.clone(), backup.clone()], &mut write_state) .await .expect("startup should select the validated backup replica"); @@ -834,9 +876,16 @@ mod tests { assert!(corrupt.read_without_lock.load(Ordering::SeqCst)); assert!(backup.read_without_lock.load(Ordering::SeqCst)); - persist_pool_meta_for_startup_if_safe(&loaded, vec![corrupt.clone(), backup.clone()], replica_state, false, true) - .await - .expect("the elected startup writer should repair validated corrupt replicas"); + persist_pool_meta_for_startup_if_safe( + &loaded, + vec![corrupt.clone(), backup.clone()], + replica_state, + write_state, + false, + true, + ) + .await + .expect("the elected startup writer should repair validated corrupt replicas"); let corrupt_write = corrupt .written_payload @@ -858,26 +907,76 @@ mod tests { async fn test_store_init_pool_meta_does_not_repair_unreadable_replica() { let valid = Arc::new(StartupPoolMetaStorage::new(startup_pool_meta_payload(&init_test_pool_meta(None)))); let unreadable = Arc::new(StartupPoolMetaStorage::unreadable()); + let mut write_state = PoolMetaWriteState::default(); - let (loaded, replica_state) = load_pool_meta_for_startup(vec![valid.clone(), unreadable.clone()]) + let (loaded, replica_state) = load_pool_meta_for_startup(vec![valid.clone(), unreadable.clone()], &mut write_state) .await .expect("startup should use a validated replica without overwriting an unreadable copy"); assert!(replica_state.needs_repair); assert!(!replica_state.repair_write_safe); - persist_pool_meta_for_startup_if_safe(&loaded, vec![valid.clone(), unreadable.clone()], replica_state, false, true) - .await - .expect("an unreadable copy should defer repair when no topology write is needed"); + persist_pool_meta_for_startup_if_safe( + &loaded, + vec![valid.clone(), unreadable.clone()], + replica_state, + write_state, + false, + true, + ) + .await + .expect("an unreadable copy should defer repair when no topology write is needed"); assert!(!valid.wrote_without_lock.load(Ordering::SeqCst)); assert!(!unreadable.wrote_without_lock.load(Ordering::SeqCst)); - let err = - persist_pool_meta_for_startup_if_safe(&loaded, vec![valid.clone(), unreadable.clone()], replica_state, true, true) - .await - .expect_err("a topology update must not overwrite an unreadable replica"); + let err = persist_pool_meta_for_startup_if_safe( + &loaded, + vec![valid.clone(), unreadable.clone()], + replica_state, + write_state, + true, + true, + ) + .await + .expect_err("a topology update must not overwrite an unreadable replica"); assert!(err.to_string().contains("cannot overwrite an unreadable replica")); assert!(!valid.wrote_without_lock.load(Ordering::SeqCst)); assert!(!unreadable.wrote_without_lock.load(Ordering::SeqCst)); + assert!(!super::should_schedule_local_decommission_resume(&[0], replica_state, true)); + assert!(!super::should_schedule_local_decommission_resume( + &[0], + crate::core::pools::PoolMetaReplicaState { + needs_repair: false, + repair_write_safe: true, + }, + false, + )); + } + + #[tokio::test] + async fn test_store_init_pool_meta_stays_blocked_after_all_replicas_were_unreadable() { + let unreadable_a = Arc::new(StartupPoolMetaStorage::unreadable()); + let unreadable_b = Arc::new(StartupPoolMetaStorage::unreadable()); + let mut write_state = PoolMetaWriteState::default(); + + load_pool_meta_for_startup(vec![unreadable_a, unreadable_b], &mut write_state) + .await + .expect_err("startup must fail when no readable pool metadata replica exists"); + + let repaired = Arc::new(StartupPoolMetaStorage::new(startup_pool_meta_payload(&init_test_pool_meta(None)))); + let (loaded, replica_state) = load_pool_meta_for_startup(vec![repaired.clone()], &mut write_state) + .await + .expect("a later startup retry may read the repaired replica"); + assert!(replica_state.repair_write_safe); + + let err = persist_pool_meta_for_startup_if_safe(&loaded, vec![repaired.clone()], replica_state, write_state, true, true) + .await + .expect_err("the same store instance must not write after observing unreadable replicas"); + assert!( + err.to_string() + .contains("restart after all replicas are readable and consistent") + ); + assert!(!repaired.wrote_without_lock.load(Ordering::SeqCst)); + assert!(!super::should_schedule_local_decommission_resume(&[0], replica_state, false)); } #[test] diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 7853dbca7..a3d5df0d0 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -33,7 +33,7 @@ use crate::bucket::utils::check_put_object_part_args; use crate::bucket::utils::{check_valid_bucket_name, check_valid_bucket_name_strict, is_meta_bucketname}; use crate::cluster::rpc::{RemoteClient, S3PeerSys}; use crate::config::storageclass; -use crate::core::pools::{DecommissionCanceler, PoolMeta}; +use crate::core::pools::{DecommissionCanceler, PoolMeta, PoolMetaWriteState}; use crate::disk::endpoint::{Endpoint, EndpointType}; use crate::disk::{DiskAPI, DiskInfo, DiskInfoOptions}; use crate::error::{Error, Result}; @@ -185,12 +185,12 @@ pub struct ECStore { /// or `decommission_cancelers`. The guarded sections may perform bounded /// async metadata work so check/init/start cannot race across operations. pub(crate) start_gate: Mutex<()>, - /// Serializes full-document pool metadata saves. + /// Serializes full-document pool metadata saves and retains a fail-closed + /// write block after startup observes an unreadable replica. /// - /// Lock order: acquire `pool_meta_save_gate` without holding `pool_meta`. - /// The saver then clones the latest `pool_meta` under a short read lock and - /// releases it before awaiting disk writes. - pub(crate) pool_meta_save_gate: Mutex<()>, + /// Lock order: acquire `pool_meta_save_gate`, then the distributed + /// `pool.bin` fence, then clone `pool_meta` under a short read lock. + pub(crate) pool_meta_save_gate: Mutex, /// Per-instance runtime state (Phase 5, backlog#939). /// /// Carries this instance's identity/runtime out of the process globals so @@ -1114,7 +1114,7 @@ mod tests { rebalance_meta: RwLock::new(None), decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), - pool_meta_save_gate: Mutex::new(()), + pool_meta_save_gate: Mutex::default(), ctx, bucket_fence_registry: Arc::default(), }) diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index 74837e46a..fd4a4750b 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -960,7 +960,7 @@ mod tests { rebalance_meta: RwLock::new(None), decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), - pool_meta_save_gate: Mutex::new(()), + pool_meta_save_gate: Mutex::default(), ctx: crate::runtime::instance::bootstrap_ctx(), bucket_fence_registry: std::sync::Arc::default(), } diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 0373e4afd..c5b7c321c 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -5449,7 +5449,7 @@ mod tests { rebalance_meta: RwLock::new(None), decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), - pool_meta_save_gate: Mutex::new(()), + pool_meta_save_gate: Mutex::default(), ctx: crate::runtime::instance::bootstrap_ctx(), bucket_fence_registry: std::sync::Arc::default(), } @@ -5512,7 +5512,7 @@ mod tests { rebalance_meta: RwLock::new(None), decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), - pool_meta_save_gate: Mutex::new(()), + pool_meta_save_gate: Mutex::default(), ctx, bucket_fence_registry: std::sync::Arc::default(), } diff --git a/crates/ecstore/src/store/rebalance.rs b/crates/ecstore/src/store/rebalance.rs index 22d6f89ea..946b4460c 100644 --- a/crates/ecstore/src/store/rebalance.rs +++ b/crates/ecstore/src/store/rebalance.rs @@ -722,9 +722,8 @@ impl ECStore { // overwrite a newer local transition after the writer commits. let movement_gate = self.ctx.data_movement_operation_gate(); let _movement_guard = movement_gate.write().await; - let mut reloaded = PoolMeta::default(); - resolve_store_rebalance_pool_meta_reload_result( - reloaded.load(self.pools[0].clone(), self.pools.clone()).await, + let reloaded = resolve_store_rebalance_pool_meta_reload_result( + self.load_runtime_pool_meta("store rebalance pool meta reload failed").await, "reload_pool_meta", )?; @@ -933,11 +932,12 @@ mod tests { use super::*; use crate::bucket::replication::{ReplicationStatusType, VersionPurgeStatusType}; use crate::config::storageclass::{CLASS_RRS, CLASS_STANDARD, lookup_config_for_pools_without_env}; - use crate::core::pools::{POOL_META_VERSION, PoolDecommissionInfo, PoolStatus}; + use crate::core::pools::{POOL_META_NAME, POOL_META_VERSION, PoolDecommissionInfo, PoolStatus}; use crate::disk::error::DiskError; use crate::layout::endpoint::Endpoint; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::object_api::ObjectLockConfigSnapshot; + use crate::set_disk::get_lock_acquire_timeout; use crate::storage_api_contracts::bucket::MakeBucketOptions; use crate::storage_api_contracts::object::ObjectIO as _; use arc_swap::ArcSwap; @@ -2122,7 +2122,7 @@ mod tests { #[test] fn resolve_store_rebalance_pool_meta_reload_result_wraps_error_context() { - let err = resolve_store_rebalance_pool_meta_reload_result(Err(Error::SlowDown), "reload_pool_meta") + let err = resolve_store_rebalance_pool_meta_reload_result::<()>(Err(Error::SlowDown), "reload_pool_meta") .expect_err("failed pool meta reload should be wrapped"); let err_message = err.to_string(); assert!(err_message.contains("store rebalance pool meta reload failed during reload_pool_meta")); @@ -2321,6 +2321,60 @@ mod tests { .expect("pool meta snapshot should persist to every pool"); } + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_runtime_load_waits_for_multi_pool_commit_fence() { + let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-load-fence", &[2, 2]).await; + let old = PoolMeta::new(&store.pools, &PoolMeta::default()); + old.save(store.pools.clone()).await.expect("old pool metadata should persist"); + let mut newer = old.clone(); + newer.pools[0].last_update += TimeDuration::seconds(1); + + let pool_meta_lock = store.pools[0] + .new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME) + .await + .expect("pool metadata lock should be created"); + let pool_meta_guard = pool_meta_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .expect("pool metadata write fence should be acquired"); + newer + .save_for_startup(vec![store.pools[0].clone()]) + .await + .expect("first replica should enter the new generation"); + + let started = Arc::new(tokio::sync::Notify::new()); + let mut load_task = tokio::spawn({ + let store = store.clone(); + let started = started.clone(); + async move { + started.notify_one(); + store.load_runtime_pool_meta("test runtime pool metadata load").await + } + }); + started.notified().await; + assert!( + tokio::time::timeout(std::time::Duration::from_millis(100), &mut load_task) + .await + .is_err(), + "runtime load must not observe a valid-old/valid-new intermediate state" + ); + + newer + .save_for_startup(vec![store.pools[1].clone()]) + .await + .expect("second replica should enter the new generation"); + drop(pool_meta_guard); + + let loaded = tokio::time::timeout(std::time::Duration::from_secs(5), load_task) + .await + .expect("runtime load should finish after commit publication") + .expect("runtime load task should not panic") + .expect("runtime load should select the completed snapshot"); + assert_eq!(loaded.pools[0].last_update, newer.pools[0].last_update); + shutdown.cancel(); + } + #[tokio::test] #[serial_test::serial] async fn peer_pool_meta_reload_does_not_rollback_newer_local_states() { diff --git a/crates/ecstore/src/store/rebalance/support.rs b/crates/ecstore/src/store/rebalance/support.rs index 25cd5914f..e1a773a54 100644 --- a/crates/ecstore/src/store/rebalance/support.rs +++ b/crates/ecstore/src/store/rebalance/support.rs @@ -67,7 +67,7 @@ pub(super) fn pool_lookup_not_found_error(bucket: &str, object: &str, opts: &Obj } } -pub(super) fn resolve_store_rebalance_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> { +pub(super) fn resolve_store_rebalance_pool_meta_reload_result(result: Result, stage: &str) -> Result { result.map_err(|err| Error::other(format!("store rebalance pool meta reload failed during {stage}: {err}"))) }