diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 7f592685a..2444d5cad 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -416,8 +416,8 @@ pub mod disk { pub mod error { pub use crate::error::{ - Error, Result, StorageError, classify_system_path_failure_reason, is_err_bucket_not_found, is_err_object_not_found, - is_err_version_not_found, + Error, PoolMetadataError, PoolMetadataFailure, Result, StorageError, classify_system_path_failure_reason, + is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found, }; } diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 3e72dcf5d..bf6ab9dd2 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -4344,7 +4344,7 @@ enum PoolMetaReplica { }, Corrupt(String), Incompatible(String), - Unreadable(String), + Unreadable(Arc), } #[derive(Debug, Clone)] @@ -4390,6 +4390,11 @@ impl PoolMetaReplicaState { pub(crate) struct PoolMetaWriteState { write_blocked: bool, aborted_transaction: Arc, + block_context: Option, + transaction_failure: Arc>>, + active_transactions: Arc, + recovery_failure: Option>, + block_started_at: Arc>>, expected_cluster_id: Option, cluster_epoch: Option, pool_meta_absent: bool, @@ -4456,6 +4461,19 @@ impl PoolMetaWriteState { pub(crate) fn independent_clone_for_test(&self) -> Self { let mut cloned = self.clone(); cloned.aborted_transaction = Arc::new(AtomicBool::new(self.aborted_transaction.load(Ordering::Acquire))); + cloned.transaction_failure = Arc::new(std::sync::Mutex::new( + self.transaction_failure + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone(), + )); + cloned.active_transactions = Arc::new(AtomicUsize::new(0)); + cloned.block_started_at = Arc::new(std::sync::Mutex::new( + *self + .block_started_at + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner), + )); cloned } @@ -4470,21 +4488,46 @@ impl PoolMetaWriteState { } pub(crate) fn observe_replicas(&mut self, replica_state: PoolMetaReplicaState) { - self.write_blocked |= !replica_state.repair_write_safe; + if !replica_state.repair_write_safe { + self.block_writes(); + } } fn block_writes(&mut self) { + self.block_with_context(pool_metadata_error( + crate::error::PoolMetadataFailure::RecoveryRequired, + "validation", + None, + )); + } + + fn block_with_context(&mut self, mut context: crate::error::PoolMetadataError) { + if !self.write_blocked { + record_pool_meta_block_once(&self.block_started_at, &mut context); + self.block_context = Some(context); + } else if let Some(previous) = &self.block_context + && previous.source.is_none() + && context.source.is_some() + { + context.since = previous.since; + self.block_context = Some(context); + } self.write_blocked = true; } pub(crate) fn block_writes_after_fence_loss(&mut self) { - self.block_writes(); + self.block_with_context(pool_metadata_error(crate::error::PoolMetadataFailure::FenceLost, "format_heal", None)); } fn arm_transaction(&self) -> PoolMetaTransactionArm { + self.active_transactions.fetch_add(1, Ordering::SeqCst); PoolMetaTransactionArm { aborted_transaction: Arc::clone(&self.aborted_transaction), - armed: true, + failure: Arc::clone(&self.transaction_failure), + active_transactions: Arc::clone(&self.active_transactions), + block_started_at: Arc::clone(&self.block_started_at), + phase: None, + source: None, } } @@ -4549,11 +4592,8 @@ impl PoolMetaWriteState { if !self.pool_meta_absent { return Ok(()); } - let result = self.validate_missing_metadata_can_initialize(); - if result.is_err() { - self.block_writes(); - } - result + self.validate_missing_metadata_can_initialize() + .map_err(|err| block_pool_meta_validation(self, err, "metadata_absence")) } fn validate_missing_metadata_can_initialize(&self) -> Result<()> { @@ -4579,29 +4619,99 @@ impl PoolMetaWriteState { if !self.write_blocked && !self.aborted_transaction.load(Ordering::SeqCst) { 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" - ))) + let mut context = self + .block_context + .clone() + .or_else(|| { + self.transaction_failure + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + }) + .unwrap_or_else(|| pool_metadata_error(crate::error::PoolMetadataFailure::FenceLost, "external_fence", None)); + context.operation = operation.to_owned(); + Err(Error::other(context)) } } +fn pool_metadata_error( + kind: crate::error::PoolMetadataFailure, + phase: &'static str, + source: Option>, +) -> crate::error::PoolMetadataError { + crate::error::PoolMetadataError { + kind, + operation: phase.to_owned(), + phase, + since: OffsetDateTime::now_utc(), + source, + } +} + +fn record_pool_meta_block_once( + started_at: &std::sync::Mutex>, + context: &mut crate::error::PoolMetadataError, +) { + let mut started_at = started_at.lock().unwrap_or_else(std::sync::PoisonError::into_inner); + if let Some(since) = *started_at { + context.since = since; + return; + } + *started_at = Some(context.since); + error!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "pool_metadata_blocked", + reason = context.kind.as_str(), + phase = context.phase, + blocked_since = %context.since, + "Pool metadata writes blocked" + ); + metrics::counter!("rustfs_pool_metadata_blocks_total", "reason" => context.kind.as_str()).increment(1); +} + #[derive(Debug)] struct PoolMetaTransactionArm { aborted_transaction: Arc, - armed: bool, + failure: Arc>>, + active_transactions: Arc, + block_started_at: Arc>>, + phase: Option<&'static str>, + source: Option>, } impl PoolMetaTransactionArm { + fn failure(&mut self, err: Error) -> Error { + let Some(phase) = self.phase else { + return err; + }; + let source = Arc::new(err); + self.source = Some(Arc::clone(&source)); + Error::other(pool_metadata_error( + crate::error::PoolMetadataFailure::TransactionUnknown, + phase, + Some(source), + )) + } + fn disarm(&mut self) { - self.armed = false; + self.phase = None; } } impl Drop for PoolMetaTransactionArm { fn drop(&mut self) { - if self.armed { - self.aborted_transaction.store(true, Ordering::SeqCst); + if let Some(phase) = self.phase { + let mut context = + pool_metadata_error(crate::error::PoolMetadataFailure::TransactionUnknown, phase, self.source.take()); + let mut failure = self.failure.lock().unwrap_or_else(std::sync::PoisonError::into_inner); + if !self.aborted_transaction.swap(true, Ordering::SeqCst) { + record_pool_meta_block_once(&self.block_started_at, &mut context); + *failure = Some(context); + } } + self.active_transactions.fetch_sub(1, Ordering::SeqCst); } } @@ -4897,7 +5007,7 @@ where cas: PoolMetaCasToken::Missing, }, Err(err) => PoolMetaReplicaRead { - replica: PoolMetaReplica::Unreadable(err.to_string()), + replica: PoolMetaReplica::Unreadable(Arc::new(err)), cas: PoolMetaCasToken::Unsafe, }, } @@ -5118,27 +5228,37 @@ where R: Into, { let replicas = replicas.into_iter().map(Into::into).collect::>(); - if replicas - .iter() - .any(|replica| matches!(&replica.replica, PoolMetaReplica::Unreadable(_))) - { - write_state.block_writes(); + if let Some(err) = pool_meta_replica_read_failure(&replicas) { + return Err(err); } match select_pool_meta_replica_reads(replicas) { Ok(selection) => { if let Err(err) = write_state.observe_selection(&selection) { - write_state.block_writes(); - return Err(err); + return Err(block_pool_meta_validation(write_state, err, "pool_selection")); } Ok(selection) } - Err(err) => { - write_state.block_writes(); - Err(err) - } + Err(err) => Err(block_pool_meta_validation(write_state, err, "pool_selection")), } } +fn block_pool_meta_validation(write_state: &mut PoolMetaWriteState, err: Error, phase: &'static str) -> Error { + let context = pool_metadata_error(crate::error::PoolMetadataFailure::RecoveryRequired, phase, Some(Arc::new(err))); + write_state.block_with_context(context.clone()); + Error::other(context) +} + +fn pool_meta_replica_read_failure(replicas: &[PoolMetaReplicaRead]) -> Option { + replicas.iter().find_map(|read| match &read.replica { + PoolMetaReplica::Unreadable(source) => Some(Error::other(pool_metadata_error( + crate::error::PoolMetadataFailure::ReadUnavailable, + "pool_read", + Some(Arc::clone(source)), + ))), + _ => None, + }) +} + fn select_pool_meta_replicas_for_read_probe( write_state: &PoolMetaWriteState, replicas: Vec, @@ -5147,7 +5267,11 @@ fn select_pool_meta_replicas_for_read_probe( where R: Into, { - let selection = select_pool_meta_replica_reads(replicas.into_iter().map(Into::into).collect())?; + let replicas = replicas.into_iter().map(Into::into).collect::>(); + if let Some(err) = pool_meta_replica_read_failure(&replicas) { + return Err(err); + } + let selection = select_pool_meta_replica_reads(replicas)?; write_state.validate_selection(&selection)?; selection.replica_state.ensure_write_safe(operation)?; if selection.absent && (write_state.expected_cluster_id.is_some() || write_state.identity_initialized.is_some()) { @@ -5185,7 +5309,17 @@ where S: EcstoreObjectIO, { let replicas = read_pool_meta_replicas(pools, no_lock).await; - select_pool_meta_replicas_for_read_probe(write_state, replicas, operation) + select_pool_meta_replicas_for_read_probe(write_state, replicas, operation).map_err(|err| { + if err.pool_metadata_failure().is_some() { + err + } else { + Error::other(pool_metadata_error( + crate::error::PoolMetadataFailure::RecoveryRequired, + "read_probe", + Some(Arc::new(err)), + )) + } + }) } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -5228,7 +5362,7 @@ enum PoolMetaIdentityReplica { Valid(PersistedPoolMetaIdentity), Corrupt(String), Incompatible(String), - Unreadable(String), + Unreadable(Arc), } #[derive(Debug)] @@ -5320,7 +5454,7 @@ where cas: PoolMetaCasToken::Missing, }, Err(err) => PoolMetaIdentityRead { - replica: PoolMetaIdentityReplica::Unreadable(err.to_string()), + replica: PoolMetaIdentityReplica::Unreadable(Arc::new(err)), cas: PoolMetaCasToken::Unsafe, }, } @@ -5396,6 +5530,16 @@ where S: EcstoreObjectIO, { let reads = join_all(pools.into_iter().map(read_pool_meta_identity_replica)).await; + if let Some(source) = reads.iter().find_map(|read| match &read.replica { + PoolMetaIdentityReplica::Unreadable(source) => Some(Arc::clone(source)), + _ => None, + }) { + return Err(Error::other(pool_metadata_error( + crate::error::PoolMetadataFailure::ReadUnavailable, + "identity_read", + Some(source), + ))); + } select_pool_meta_identity(reads, expected_cluster_id) } @@ -5410,13 +5554,17 @@ where let selection = match load_pool_meta_identity(pools, expected_cluster_id).await { Ok(selection) => selection, Err(err) => { - write_state.block_writes(); - return Err(err); + if err + .pool_metadata_failure() + .is_some_and(|context| context.kind == crate::error::PoolMetadataFailure::ReadUnavailable) + { + return Err(err); + } + return Err(block_pool_meta_validation(write_state, err, "identity_selection")); } }; if let Err(err) = write_state.observe_identity(&selection) { - write_state.block_writes(); - return Err(err); + return Err(block_pool_meta_validation(write_state, err, "identity_selection")); } Ok(selection) } @@ -5519,6 +5667,7 @@ async fn save_pool_meta_object_cas( token: &PoolMetaCasToken, fence: &PoolMetaPersistenceFence<'_>, phase: &'static str, + transaction_arm: &mut PoolMetaTransactionArm, ) -> Result where S: EcstoreObjectIO, @@ -5532,11 +5681,30 @@ where ..Default::default() }; fence.add_to_options(&mut opts); + // Cancellation can happen at the very first poll of the storage future. + // Arm before dispatch, but not during read/encode/fence preflight. + let previous_phase = transaction_arm.phase; + transaction_arm.phase = Some(phase); let result = save_config_with_opts_and_metadata(pool, object, data, &opts).await; if matches!(&result, Err(Error::PreconditionFailed)) { record_pool_meta_stale_write_rejection(phase); + transaction_arm.phase = previous_phase; } - let object_info = result?; + let object_info = match result { + Ok(info) => info, + Err(err) => { + let source = Arc::new(err); + transaction_arm.source = Some(Arc::clone(&source)); + if matches!(source.as_ref(), Error::PreconditionFailed) { + return Err(Error::PreconditionFailed); + } + return Err(Error::other(pool_metadata_error( + crate::error::PoolMetadataFailure::TransactionUnknown, + phase, + Some(source), + ))); + } + }; fence.ensure_held()?; Ok(object_info) } @@ -5546,6 +5714,7 @@ async fn persist_pool_meta_identity( write_state: &mut PoolMetaWriteState, initialized: bool, fence: &PoolMetaPersistenceFence<'_>, + transaction_arm: &mut PoolMetaTransactionArm, ) -> Result<()> where S: EcstoreObjectIO, @@ -5588,7 +5757,17 @@ where let data = encode_pool_meta_identity(identity)?; let mut conflict = false; for (pool, token) in pools.iter().cloned().zip(&selection.cas_tokens) { - match save_pool_meta_object_cas(pool, POOL_META_IDENTITY_NAME, data.clone(), token, fence, "identity_cas").await { + match save_pool_meta_object_cas( + pool, + POOL_META_IDENTITY_NAME, + data.clone(), + token, + fence, + "identity_cas", + transaction_arm, + ) + .await + { Ok(_) => {} Err(Error::PreconditionFailed) => { conflict = true; @@ -5631,7 +5810,17 @@ pub(crate) async fn persist_pool_meta_identity_for_startup( where S: EcstoreObjectIO, { - persist_pool_meta_identity(pools, write_state, initialized, &PoolMetaPersistenceFence::Distributed(None)).await + let mut transaction_arm = write_state.arm_transaction(); + persist_pool_meta_identity( + pools, + write_state, + initialized, + &PoolMetaPersistenceFence::Distributed(None), + &mut transaction_arm, + ) + .await?; + transaction_arm.disarm(); + Ok(()) } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -6125,6 +6314,135 @@ struct PoolMetaSaveOutcome { committed: PoolMeta, } +/// Reconcile only a known, interrupted pool.bin transaction. The caller keeps +/// the original block installed and holds the namespace writer through runtime +/// publication. Never feed speculative in-memory progress into this repair. +async fn load_pool_meta_for_transaction_recovery( + pools: Vec>, + write_state: &mut PoolMetaWriteState, +) -> Result +where + S: EcstoreObjectIO, +{ + let result = async { + if let Some(cluster_id) = write_state.expected_cluster_id { + let reads = join_all(pools.iter().cloned().map(read_pool_meta_identity_replica)).await; + for read in &reads { + match &read.replica { + PoolMetaIdentityReplica::Unreadable(source) => { + return Err(Error::other(pool_metadata_error( + crate::error::PoolMetadataFailure::ReadUnavailable, + "recovery_identity_read", + Some(Arc::clone(source)), + ))); + } + PoolMetaIdentityReplica::Corrupt(_) | PoolMetaIdentityReplica::Incompatible(_) => { + return Err(Error::other("pool metadata recovery requires operator repair of cluster identity")); + } + _ => {} + } + } + let identity = select_pool_meta_identity(reads, cluster_id)?; + if !identity.identity.is_some_and(|identity| identity.initialized) { + return Err(Error::other("pool metadata recovery requires an initialized durable cluster identity")); + } + write_state.observe_identity(&identity)?; + } + let reads = read_pool_meta_replicas(pools, true).await; + if let Some(err) = pool_meta_replica_read_failure(&reads) { + return Err(err); + } + if reads + .iter() + .any(|read| matches!(&read.replica, PoolMetaReplica::Corrupt(_) | PoolMetaReplica::Incompatible(_))) + { + return Err(Error::other( + "pool metadata recovery requires operator repair of corrupt or incompatible replicas", + )); + } + let selection = select_pool_meta_replica_reads(reads)?; + write_state.observe_selection(&selection)?; + if selection.absent || selection.meta.pools.is_empty() { + return Err(Error::other("pool metadata recovery cannot initialize missing metadata")); + } + selection + .replica_state + .ensure_write_safe("pool metadata transaction recovery")?; + Ok(selection) + } + .await; + result.map_err(|err: Error| { + if err.pool_metadata_failure().is_some() { + err + } else { + Error::other(pool_metadata_error( + crate::error::PoolMetadataFailure::RecoveryRequired, + "recovery_validation", + Some(Arc::new(err)), + )) + } + }) +} + +async fn repair_pool_meta_transaction( + pools: Vec>, + write_state: &mut PoolMetaWriteState, + selection: PoolMetaSelection, + fence: &PoolMetaPersistenceFence<'_>, +) -> Result +where + S: EcstoreObjectIO, +{ + let mut transaction_arm = write_state.arm_transaction(); + let mut expected_meta = selection.meta.clone(); + let mut expected_revision = selection.revision; + // A failed first V3 prepare still establishes the format floor. Commit the + // previous authoritative state as V3 instead of erasing that observation. + if selection.generation_protocol_observed && !expected_revision.is_generation_protocol() { + let (cluster_id, epoch) = selection + .generation_identity + .ok_or_else(|| Error::other("pool metadata recovery has no V3 cluster identity"))?; + expected_meta.version = POOL_META_GENERATION_VERSION; + expected_revision = PoolMetaRevision { + version: POOL_META_GENERATION_VERSION, + cluster_id: Some(cluster_id), + epoch, + generation: 1, + transaction_id: Some(uuid::Uuid::new_v4()), + }; + } + let expected_canonical = if expected_revision.is_generation_protocol() { + encode_pool_meta_v3_envelope(&expected_meta, expected_revision, true, None)? + } else { + expected_meta.encode_config_data_for_v2_gate(true)? + }; + if selection.replica_state.needs_repair { + let data = if expected_revision.is_generation_protocol() { + expected_canonical.clone() + } else { + expected_meta.encode_config_data_for_v2_gate(expected_revision.version == POOL_META_VERSION)? + }; + for (pool, token) in pools.iter().cloned().zip(&selection.cas_tokens) { + save_pool_meta_object_cas(pool, POOL_META_NAME, data.clone(), token, fence, "recovery_cas", &mut transaction_arm) + .await?; + } + } + persist_pool_meta_identity(pools.clone(), write_state, true, fence, &mut transaction_arm).await?; + let confirmed = load_pool_meta_for_transaction_recovery(pools, write_state).await?; + if confirmed.revision != expected_revision + || confirmed.canonical.as_ref() != Some(&expected_canonical) + || confirmed.replica_state.needs_repair + || write_state.identity_requires_repair() + { + return Err(Error::PreconditionFailed); + } + fence.ensure_held()?; + Ok(PoolMetaSaveOutcome { + transaction_arm, + committed: confirmed.meta, + }) +} + impl PoolMetaSaveOutcome { fn disarm(mut self) { self.transaction_arm.disarm(); @@ -6426,6 +6744,49 @@ impl PoolMeta { Ok(selection.replica_state) } + pub(crate) async fn load_for_startup_observing( + &mut self, + pools: Vec>, + write_state: &mut PoolMetaWriteState, + ) -> Result + where + S: EcstoreObjectIO, + { + // Startup may install a readable replica without authorizing writes. + // Runtime preflight retries must not inherit this bootstrap policy. + if let Some(cluster_id) = write_state.expected_cluster_id { + let reads = join_all(pools.iter().cloned().map(read_pool_meta_identity_replica)).await; + if let Some(source) = reads.iter().find_map(|read| match &read.replica { + PoolMetaIdentityReplica::Unreadable(source) => Some(Arc::clone(source)), + _ => None, + }) { + write_state.block_with_context(pool_metadata_error( + crate::error::PoolMetadataFailure::RecoveryRequired, + "startup_identity_read", + Some(source), + )); + } + let identity = select_pool_meta_identity(reads, cluster_id) + .map_err(|err| block_pool_meta_validation(write_state, err, "startup_identity"))?; + write_state.observe_identity(&identity)?; + } + let reads = read_pool_meta_replicas(pools, true).await; + if let Some(err) = pool_meta_replica_read_failure(&reads) { + write_state.block_with_context(pool_metadata_error( + crate::error::PoolMetadataFailure::RecoveryRequired, + "startup_read", + Some(Arc::new(err)), + )); + } + let selection = select_pool_meta_replica_reads(reads) + .map_err(|err| block_pool_meta_validation(write_state, err, "startup_selection"))?; + write_state + .observe_selection(&selection) + .map_err(|err| block_pool_meta_validation(write_state, err, "startup_selection"))?; + *self = selection.meta; + Ok(selection.replica_state) + } + #[cfg(test)] fn encode_config_data(&self) -> Result> { self.encode_config_data_for_v2_gate(pool_meta_v2_writer_enabled()) @@ -6552,11 +6913,15 @@ impl PoolMeta { S: EcstoreObjectIO, { write_state.ensure_write_safe("pool metadata activation save failed")?; - let transaction_arm = write_state.arm_transaction(); + let mut transaction_arm = write_state.arm_transaction(); let fence = PoolMetaPersistenceFence::Activation(activation_fence); let committed = self - .save_no_lock_transaction(pools, write_state, &fence, Some(indices)) - .await?; + .save_no_lock_transaction(pools, write_state, &fence, Some(indices), &mut transaction_arm) + .await + .map_err(|err| transaction_arm.failure(err))?; + if transaction_arm.phase.is_some() { + transaction_arm.phase = Some("publication"); + } Ok(PoolMetaSaveOutcome { transaction_arm, committed, @@ -6591,9 +6956,15 @@ impl PoolMeta { // The arm does not block its own transaction. If this future or its // returned outcome is dropped before publication, Drop latches the // sticky recovery gate. - let transaction_arm = write_state.arm_transaction(); + let mut transaction_arm = write_state.arm_transaction(); let fence = PoolMetaPersistenceFence::Distributed(lock_lost); - let committed = self.save_no_lock_transaction(pools, write_state, &fence, indices).await?; + let committed = self + .save_no_lock_transaction(pools, write_state, &fence, indices, &mut transaction_arm) + .await + .map_err(|err| transaction_arm.failure(err))?; + if transaction_arm.phase.is_some() { + transaction_arm.phase = Some("publication"); + } Ok(PoolMetaSaveOutcome { transaction_arm, committed, @@ -6606,6 +6977,7 @@ impl PoolMeta { write_state: &mut PoolMetaWriteState, fence: &PoolMetaPersistenceFence<'_>, indices: Option<&[usize]>, + transaction_arm: &mut PoolMetaTransactionArm, ) -> Result where S: EcstoreObjectIO, @@ -6618,7 +6990,7 @@ impl PoolMeta { } for attempt in 0..POOL_META_CAS_MAX_ATTEMPTS { match self - .save_no_lock_transaction_once(pools.clone(), write_state, fence, indices) + .save_no_lock_transaction_once(pools.clone(), write_state, fence, indices, transaction_arm) .await { Ok(committed) => return Ok(committed), @@ -6635,6 +7007,7 @@ impl PoolMeta { write_state: &mut PoolMetaWriteState, fence: &PoolMetaPersistenceFence<'_>, indices: Option<&[usize]>, + transaction_arm: &mut PoolMetaTransactionArm, ) -> Result where S: EcstoreObjectIO, @@ -6665,7 +7038,7 @@ impl PoolMeta { } if !selection.absent && write_state.identity_requires_repair() { let initialized = write_state.identity_initialized != Some(false) || selection.revision.is_generation_protocol(); - persist_pool_meta_identity(pools.clone(), write_state, initialized, fence).await?; + persist_pool_meta_identity(pools.clone(), write_state, initialized, fence, transaction_arm).await?; } } // Startup is the only path allowed to create an all-missing metadata @@ -6705,8 +7078,16 @@ impl PoolMeta { let data = committed.encode_config_data_for_v2_gate(target_version == POOL_META_VERSION)?; let mut canonical_saved = false; for (pool_index, (pool, token)) in pools.iter().cloned().zip(&selection.cas_tokens).enumerate() { - let result = - save_pool_meta_object_cas(pool.clone(), POOL_META_NAME, data.clone(), token, fence, "legacy_cas").await; + let result = save_pool_meta_object_cas( + pool.clone(), + POOL_META_NAME, + data.clone(), + token, + fence, + "legacy_cas", + transaction_arm, + ) + .await; if result.is_ok() && pool_index == 0 { canonical_saved = true; #[cfg(test)] @@ -6798,7 +7179,8 @@ impl PoolMeta { let mut pending_tokens = Vec::with_capacity(pools.len()); for (pool, token) in pools.iter().cloned().zip(&selection.cas_tokens) { let object_info = - save_pool_meta_object_cas(pool, POOL_META_NAME, pending.clone(), token, fence, "prepare_cas").await?; + save_pool_meta_object_cas(pool, POOL_META_NAME, pending.clone(), token, fence, "prepare_cas", transaction_arm) + .await?; let etag = object_info .etag .filter(|etag| !etag.trim().is_empty()) @@ -6811,7 +7193,17 @@ impl PoolMeta { #[cfg(test)] let mut first_pool = true; for (pool, token) in pools.iter().cloned().zip(&pending_tokens) { - match save_pool_meta_object_cas(pool.clone(), POOL_META_NAME, durable.clone(), token, fence, "commit_cas").await { + match save_pool_meta_object_cas( + pool.clone(), + POOL_META_NAME, + durable.clone(), + token, + fence, + "commit_cas", + transaction_arm, + ) + .await + { Ok(_) => { commit_succeeded = true; #[cfg(test)] @@ -6836,7 +7228,7 @@ impl PoolMeta { confirmed }; if confirmed.revision == revision && confirmed.canonical.as_ref() == Some(&durable) { - persist_pool_meta_identity(pools, write_state, true, fence).await?; + persist_pool_meta_identity(pools, write_state, true, fence, transaction_arm).await?; return Ok(confirmed.meta); } if !commit_succeeded { @@ -10275,8 +10667,174 @@ impl ECStore { /// Read-only admission probes do not change this state; startup and real /// metadata transactions still latch it on unrecoverable conditions. pub async fn pool_meta_writes_ready(&self) -> bool { - let write_state = self.pool_meta_save_gate.lock().await; - !write_state.write_blocked && !write_state.aborted_transaction.load(Ordering::SeqCst) + self.pool_meta_write_status().await.is_ok() + } + + /// A health probe must not wait behind a metadata transaction's disk I/O. + pub async fn pool_meta_write_status(&self) -> Result<()> { + let write_state = tokio::time::timeout(std::time::Duration::from_millis(100), self.pool_meta_save_gate.lock()) + .await + .map_err(|_| Error::Timeout)?; + write_state.ensure_write_safe("pool metadata readiness") + } + + pub(crate) async fn recover_pool_meta_transaction(&self) -> Result { + { + let Ok(state) = self.pool_meta_save_gate.try_lock() else { + return Ok(false); + }; + if state.ensure_write_safe("pool metadata recovery").is_ok() { + return Ok(false); + } + if state.write_blocked + || state.active_transactions.load(Ordering::SeqCst) != 0 + || state + .transaction_failure + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .is_none() + { + return Ok(false); + } + } + let _start_guard = self.start_gate.lock().await; + // Cancellation must precede the movement writer. Let the existing + // supervisor join every worker and relinquish its operation ownership; + // a new worker must never overlap an old terminal/progress callback. + { + let cancelers = self.decommission_cancelers.read().await; + let mut draining = false; + for canceler in cancelers.iter().flatten().filter(|canceler| canceler.is_active()) { + canceler.cancel(); + draining = true; + } + if draining { + return Err(Error::OperationCanceled); + } + } + let movement_gate = self.ctx.data_movement_operation_gate(); + let _movement_guard = movement_gate.write().await; + let mut state = self.pool_meta_save_gate.lock().await; + // Outcomes can outlive the save-gate guard. Do not let an old Drop + // relatch, or an older recovery clear a newly installed block. + if state.write_blocked || state.active_transactions.load(Ordering::SeqCst) != 0 { + return Ok(false); + } + let Some(blocked) = state + .transaction_failure + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + else { + return Ok(false); + }; + if self + .rebalance_meta + .read() + .await + .as_ref() + .is_some_and(|meta| meta.cancel.is_some()) + { + return Err(Error::other("pool metadata recovery requires rebalance workers to be quiescent")); + } + let pool = self + .pools + .first() + .ok_or_else(|| Error::other("pool metadata recovery has no storage pools"))?; + let lock = pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?; + let guard = lock.get_write_lock(get_lock_acquire_timeout()).await?; + let fence = PoolMetaPersistenceFence::Distributed(guard.lock_lost_signal()); + let mut candidate = PoolMetaWriteState { + aborted_transaction: Arc::new(AtomicBool::new(false)), + transaction_failure: Arc::default(), + active_transactions: Arc::default(), + block_context: None, + recovery_failure: None, + write_blocked: false, + ..state.clone() + }; + let selection = load_pool_meta_for_transaction_recovery(self.pools.clone(), &mut candidate).await?; + if selection.meta.pools.len() != self.pools.len() + || selection + .meta + .pools + .iter() + .zip(&self.pools) + .enumerate() + .any(|(index, (meta, pool))| meta.id != index || meta.cmd_line != pool.endpoints.cmd_line) + { + return Err(Error::other( + "pool metadata recovery requires operator reconciliation of the pool topology", + )); + } + let outcome = repair_pool_meta_transaction(self.pools.clone(), &mut candidate, selection, &fence).await?; + let mut pool_meta = self.pool_meta.write().await; + fence.ensure_held()?; + self.ctx.advance_data_movement_operation_epoch(); + *pool_meta = outcome.committed.clone(); + fence.ensure_held()?; + outcome.disarm(); + // No await between final fence validation, publication, and clearing. + candidate.aborted_transaction = Arc::clone(&state.aborted_transaction); + candidate.transaction_failure = Arc::clone(&state.transaction_failure); + candidate.active_transactions = Arc::clone(&state.active_transactions); + candidate.block_started_at = Arc::clone(&state.block_started_at); + *state + .block_started_at + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = None; + *state + .transaction_failure + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = None; + state.aborted_transaction.store(false, Ordering::SeqCst); + *state = candidate; + info!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "pool_metadata_recovered", + phase = blocked.phase, + blocked_since = %blocked.since, + "Pool metadata writes recovered" + ); + metrics::counter!("rustfs_pool_metadata_recoveries_total").increment(1); + Ok(true) + } + + pub(crate) fn record_pool_meta_recovery_failure(&self, error: Error) { + let Ok(mut state) = self.pool_meta_save_gate.try_lock() else { + return; + }; + if state.ensure_write_safe("pool metadata recovery").is_ok() { + return; + } + let classify = |error: &Error| match error.pool_metadata_failure() { + Some(context) => (context.kind.as_str(), context.phase), + None => match error { + Error::OperationCanceled => ("workers_draining", "recovery_quiesce"), + Error::Timeout => ("timeout", "recovery"), + Error::Lock(_) => ("fence_unavailable", "recovery"), + _ => ("unavailable", "recovery"), + }, + }; + let (reason, phase) = classify(&error); + if state + .recovery_failure + .as_ref() + .is_none_or(|previous| classify(previous) != (reason, phase)) + { + warn!( + event = EVENT_DECOMMISSION_STATE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_POOLS, + state = "pool_metadata_recovery_pending", + reason, + phase, + "Pool metadata recovery remains blocked" + ); + } + state.recovery_failure = Some(Arc::new(error)); } async fn load_runtime_pool_meta_observing(&self, write_state: &mut PoolMetaWriteState, operation: &str) -> Result { @@ -17719,6 +18277,46 @@ mod tests { assert_eq!(total, 2); } + #[test] + fn pool_meta_block_transition_is_counted_once_across_both_latches() { + let recorder = metrics_util::debugging::DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + metrics::with_local_recorder(&recorder, || { + let mut state = PoolMetaWriteState::default(); + let mut arm = state.arm_transaction(); + arm.phase = Some("commit_cas"); + drop(arm); + let first = state + .ensure_write_safe("first block") + .unwrap_err() + .pool_metadata_failure() + .unwrap() + .since; + state.block_writes_after_fence_loss(); + state.block_writes(); + assert_eq!( + state + .ensure_write_safe("later block") + .unwrap_err() + .pool_metadata_failure() + .unwrap() + .since, + first + ); + }); + let total = snapshotter + .snapshot() + .into_vec() + .into_iter() + .filter(|(key, _, _, _)| key.key().name() == "rustfs_pool_metadata_blocks_total") + .filter_map(|(_, _, _, value)| match value { + metrics_util::debugging::DebugValue::Counter(count) => Some(count), + _ => None, + }) + .sum::(); + assert_eq!(total, 1); + } + fn pool_meta_v3_test_revision(cluster_id: uuid::Uuid, generation: u64, transaction_id: uuid::Uuid) -> PoolMetaRevision { PoolMetaRevision { version: POOL_META_GENERATION_VERSION, @@ -18012,7 +18610,7 @@ mod tests { let err = select_pool_meta_replica(vec![ PoolMetaReplica::Missing, - PoolMetaReplica::Unreadable("read quorum unavailable".to_string()), + PoolMetaReplica::Unreadable(Arc::new(Error::other("read quorum unavailable"))), ]) .expect_err("an unreadable replica must not be treated as a new deployment"); assert!(err.to_string().contains("no valid committed replica is available")); @@ -18023,7 +18621,7 @@ mod tests { fn pool_meta_replica_selection_blocks_repair_for_unreadable_copy() { let selection = select_pool_meta_replica(vec![ decode_pool_meta_replica(pool_meta_replica_test_data("pool-0")), - PoolMetaReplica::Unreadable("read quorum unavailable".to_string()), + PoolMetaReplica::Unreadable(Arc::new(Error::other("read quorum unavailable"))), ]) .expect("a validated replica should remain usable while another copy is unreadable"); @@ -18057,7 +18655,7 @@ mod tests { let write_state = PoolMetaWriteState::default(); select_pool_meta_replicas_for_read_probe( &write_state, - vec![PoolMetaReplica::Unreadable("transient read failure".to_string())], + vec![PoolMetaReplica::Unreadable(Arc::new(Error::Timeout))], "capacity probe", ) .expect_err("an unreadable probe replica must fail the current admission"); @@ -18068,6 +18666,156 @@ mod tests { ); } + #[test] + fn pool_meta_preflight_drop_does_not_arm_but_started_drop_retains_first_failure() { + let state = PoolMetaWriteState::default(); + drop(state.arm_transaction()); + state.ensure_write_safe("preflight retry").unwrap(); + let mut first = state.arm_transaction(); + first.phase = Some("identity_cas"); + let err = first.failure(Error::Timeout); + drop(first); + let context = state.ensure_write_safe("retry").unwrap_err(); + let first_failure = context.pool_metadata_failure().unwrap(); + assert_eq!(first_failure.kind, crate::error::PoolMetadataFailure::TransactionUnknown); + assert_eq!(first_failure.phase, "identity_cas"); + assert!(matches!(first_failure.source.as_deref(), Some(Error::Timeout))); + assert!(matches!(err.pool_metadata_failure().unwrap().source.as_deref(), Some(Error::Timeout))); + let mut late = state.arm_transaction(); + late.phase = Some("commit_cas"); + drop(late); + let later = state.ensure_write_safe("retry").unwrap_err(); + assert_eq!(later.pool_metadata_failure().unwrap().since, first_failure.since); + assert_eq!(later.pool_metadata_failure().unwrap().phase, "identity_cas"); + assert_eq!(state.active_transactions.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_writer_retries_after_identity_read_failure_without_restart() { + let (_dirs, store, _peer) = crate::services::rebalance::test_two_pool_stores(None).await; + let mut disks = Vec::new(); + for set in &store.pools[1].disk_set { + let mut guard = set.disks.write().await; + let disk_count = guard.len(); + let original = std::mem::replace(&mut *guard, vec![None; disk_count]); + disks.push((set.clone(), original)); + } + let mut state = store.pool_meta_save_gate.lock().await; + let err = store + .acquire_pool_meta_write_guard(&mut state, "writer preflight") + .await + .unwrap_err(); + assert_eq!( + err.pool_metadata_failure().unwrap().kind, + crate::error::PoolMetadataFailure::ReadUnavailable + ); + state.ensure_write_safe("retry").unwrap(); + for (set, original) in disks { + *set.disks.write().await = original; + } + let (_guard, _) = store.acquire_pool_meta_write_guard(&mut state, "writer retry").await.unwrap(); + } + + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_runtime_recovery_publishes_durable_state_and_drains_old_workers() { + let (_dirs, store, _peer) = crate::services::rebalance::test_two_pool_stores(None).await; + let before = store.pool_meta.read().await.clone(); + let mut requested = before.clone(); + requested.pools[0].last_update += Duration::seconds(1); + requested.version = POOL_META_GENERATION_VERSION; + { + let mut state = store.pool_meta_save_gate.lock().await; + let (fence, _) = store.acquire_pool_meta_write_guard(&mut state, "test save").await.unwrap(); + let outcome = requested + .save_no_lock_armed(store.pools.clone(), &mut state, fence.lock_lost_signal(), &[0]) + .await + .unwrap(); + assert_eq!(outcome.committed.pools[0].last_update, requested.pools[0].last_update); + drop(outcome); + } + assert!(!store.pool_meta_writes_ready().await); + let worker = DecommissionCanceler::new(CancellationToken::new()); + store.decommission_cancelers.write().await[0] = Some(worker.clone()); + assert!(store.recover_pool_meta_transaction().await.is_err()); + assert!(worker.is_cancelled()); + assert!(!store.pool_meta_writes_ready().await); + assert_eq!(store.pool_meta.read().await.pools[0].last_update, before.pools[0].last_update); + worker.release(); + store.pool_meta.write().await.pools[0].last_update = requested.pools[0].last_update + Duration::seconds(10); + // An outcome may outlive the save mutex. Its Drop must precede recovery. + let outstanding = store.pool_meta_save_gate.lock().await.arm_transaction(); + assert!(!store.recover_pool_meta_transaction().await.unwrap()); + drop(outstanding); + assert!(store.recover_pool_meta_transaction().await.unwrap()); + assert!(store.pool_meta_writes_ready().await); + assert_eq!(store.pool_meta.read().await.pools[0].last_update, requested.pools[0].last_update); + assert!(!store.recover_pool_meta_transaction().await.unwrap()); + store.pool_meta_save_gate.lock().await.block_writes_after_fence_loss(); + assert!(!store.recover_pool_meta_transaction().await.unwrap()); + assert!(!store.pool_meta_writes_ready().await); + } + + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_readiness_is_bounded_while_save_gate_is_held() { + let (_dirs, store, _peer) = crate::services::rebalance::test_two_pool_stores(None).await; + let _state = store.pool_meta_save_gate.lock().await; + let result = tokio::time::timeout(std::time::Duration::from_secs(1), store.pool_meta_write_status()) + .await + .unwrap(); + assert!(matches!(result, Err(Error::Timeout))); + } + + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_cancelled_recovery_and_new_integrity_block_never_clear_original_gate() { + let (_dirs, store, _peer) = crate::services::rebalance::test_two_pool_stores(None).await; + { + let state = store.pool_meta_save_gate.lock().await; + let mut arm = state.arm_transaction(); + arm.phase = Some("publication"); + } + let publication_guard = store.pool_meta.write().await; + let recovery = tokio::spawn({ + let store = store.clone(); + async move { store.recover_pool_meta_transaction().await } + }); + tokio::time::timeout(std::time::Duration::from_secs(10), async { + loop { + if store.pool_meta_save_gate.try_lock().is_err() { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + recovery.abort(); + assert!(recovery.await.unwrap_err().is_cancelled()); + drop(publication_guard); + assert!(!store.pool_meta_writes_ready().await); + let start_guard = store.start_gate.lock().await; + let recovery = tokio::spawn({ + let store = store.clone(); + async move { store.recover_pool_meta_transaction().await } + }); + store.pool_meta_save_gate.lock().await.block_writes_after_fence_loss(); + drop(start_guard); + assert!(!recovery.await.unwrap().unwrap()); + assert_eq!( + store + .pool_meta_write_status() + .await + .unwrap_err() + .pool_metadata_failure() + .unwrap() + .kind, + crate::error::PoolMetadataFailure::FenceLost + ); + } + #[test] fn pool_meta_read_probe_rejects_missing_runtime_metadata_without_latching() { let write_state = PoolMetaWriteState { @@ -18114,22 +18862,22 @@ mod tests { } #[test] - fn pool_meta_write_state_blocks_when_selection_has_no_valid_replica() { + fn pool_meta_writer_read_failure_is_retryable_when_all_replicas_are_unreadable() { let replicas = vec![ - PoolMetaReplica::Unreadable("pool 0 read quorum unavailable".to_string()), - PoolMetaReplica::Unreadable("pool 1 read quorum unavailable".to_string()), + PoolMetaReplica::Unreadable(Arc::new(Error::other("pool 0 read quorum unavailable"))), + PoolMetaReplica::Unreadable(Arc::new(Error::other("pool 1 read quorum unavailable"))), ]; 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") + let err = select_pool_meta_replicas_observing(&mut write_state, replicas) + .expect_err("all unreadable replicas must fail selection"); + assert_eq!( + err.pool_metadata_failure().unwrap().kind, + crate::error::PoolMetadataFailure::ReadUnavailable ); + write_state + .ensure_write_safe("pool metadata save failed") + .expect("a read-only preflight has no unknown writes and must remain retryable"); } #[test] @@ -19601,7 +20349,7 @@ mod pools_tests { }) } - #[derive(Debug)] + #[derive(Debug, Default)] struct PartialPoolMetaWriteStorage { fail_write: bool, fail_after_first_write: bool, @@ -19613,6 +20361,256 @@ mod pools_tests { identity: StdMutex, String)>>, } + #[tokio::test] + async fn pool_meta_recovery_reconciles_prepare_and_partial_commit_without_format_downgrade() { + for previous_version in [POOL_META_V1_VERSION, POOL_META_VERSION, POOL_META_GENERATION_VERSION] { + for commit_first in [false, true] { + let cluster_id = uuid::Uuid::new_v4(); + let previous_meta = PoolMeta { + version: previous_version, + pools: vec![decommission_test_pool_status(0, None)], + ..Default::default() + }; + let previous_revision = if previous_version == POOL_META_GENERATION_VERSION { + super::PoolMetaRevision { + version: previous_version, + cluster_id: Some(cluster_id), + epoch: 1, + generation: 9, + transaction_id: Some(uuid::Uuid::new_v4()), + } + } else { + super::PoolMetaRevision::legacy(previous_version) + }; + let previous_data = if previous_revision.is_generation_protocol() { + super::encode_pool_meta_v3_envelope(&previous_meta, previous_revision, true, None).unwrap() + } else { + previous_meta.encode_config_data_for_v2_gate(true).unwrap() + }; + let previous = super::PoolMetaCommittedCandidate { + canonical: previous_data, + meta: previous_meta.clone(), + revision: previous_revision, + }; + let mut next = previous_meta.clone(); + next.version = POOL_META_GENERATION_VERSION; + next.pools[0].last_update += Duration::seconds(1); + let revision = super::PoolMetaRevision { + version: POOL_META_GENERATION_VERSION, + cluster_id: Some(cluster_id), + epoch: 1, + generation: previous_revision.generation + 1, + transaction_id: Some(uuid::Uuid::new_v4()), + }; + let pending = super::encode_pool_meta_v3_envelope(&next, revision, false, Some(&previous)).unwrap(); + let durable = super::encode_pool_meta_v3_envelope(&next, revision, true, None).unwrap(); + let identity = super::initialized_pool_meta_identity_for_test(cluster_id, 1).unwrap(); + let pools = [if commit_first { durable } else { pending.clone() }, pending] + .into_iter() + .map(|data| { + Arc::new(PartialPoolMetaWriteStorage { + stored: StdMutex::new(Some((data, "initial".to_owned()))), + identity: StdMutex::new(Some((identity.clone(), "identity".to_owned()))), + ..Default::default() + }) + }) + .collect::>(); + let mut state = super::PoolMetaWriteState::for_startup(cluster_id, false); + let selection = super::load_pool_meta_for_transaction_recovery(pools.clone(), &mut state) + .await + .unwrap(); + assert_eq!( + selection.meta.pools[0].last_update, + if commit_first { + next.pools[0].last_update + } else { + previous_meta.pools[0].last_update + } + ); + let outcome = super::repair_pool_meta_transaction( + pools.clone(), + &mut state, + selection, + &PoolMetaPersistenceFence::Distributed(None), + ) + .await + .unwrap(); + assert_eq!(outcome.committed.version, POOL_META_GENERATION_VERSION); + assert_eq!( + outcome.committed.pools[0].last_update, + if commit_first { + next.pools[0].last_update + } else { + previous_meta.pools[0].last_update + } + ); + outcome.disarm(); + let confirmed = super::load_pool_meta_for_transaction_recovery(pools, &mut state) + .await + .unwrap(); + assert!(!confirmed.replica_state.needs_repair); + state.ensure_write_safe("repaired").unwrap(); + } + } + } + + #[tokio::test] + async fn pool_meta_recovery_rejects_missing_corrupt_and_incompatible_replicas_without_writes() { + let cluster_id = uuid::Uuid::new_v4(); + let identity = super::initialized_pool_meta_identity_for_test(cluster_id, 1).unwrap(); + for payload in [None, Some(Vec::new()), Some(vec![1, 0, 255, 255, 0])] { + let pool = Arc::new(PartialPoolMetaWriteStorage { + stored: StdMutex::new(payload.map(|data| (data, "initial".to_owned()))), + identity: StdMutex::new(Some((identity.clone(), "identity".to_owned()))), + ..Default::default() + }); + let mut state = super::PoolMetaWriteState::for_startup(cluster_id, false); + super::load_pool_meta_for_transaction_recovery(vec![pool.clone()], &mut state) + .await + .unwrap_err(); + assert!(!pool.wrote.load(Ordering::SeqCst)); + } + } + + #[tokio::test] + async fn pool_meta_first_cas_rejection_is_retryable_but_later_rejection_keeps_prior_write_armed() { + let pool = Arc::new(PartialPoolMetaWriteStorage::default()); + let state = super::PoolMetaWriteState::default(); + let fence = PoolMetaPersistenceFence::Distributed(None); + let mut arm = state.arm_transaction(); + save_pool_meta_object_cas( + pool.clone(), + POOL_META_NAME, + vec![1], + &PoolMetaCasToken::Existing("missing".to_owned()), + &fence, + "prepare_cas", + &mut arm, + ) + .await + .unwrap_err(); + drop(arm); + state.ensure_write_safe("first rejection").unwrap(); + let mut arm = state.arm_transaction(); + save_pool_meta_object_cas( + pool.clone(), + POOL_META_NAME, + vec![1], + &PoolMetaCasToken::Missing, + &fence, + "prepare_cas", + &mut arm, + ) + .await + .unwrap(); + save_pool_meta_object_cas(pool, POOL_META_NAME, vec![2], &PoolMetaCasToken::Missing, &fence, "commit_cas", &mut arm) + .await + .unwrap_err(); + drop(arm); + let err = state.ensure_write_safe("later rejection").unwrap_err(); + assert_eq!(err.pool_metadata_failure().unwrap().phase, "prepare_cas"); + assert!(matches!( + err.pool_metadata_failure().unwrap().source.as_deref(), + Some(Error::PreconditionFailed) + )); + } + + #[tokio::test] + async fn pool_meta_recovery_cas_conflict_after_repair_write_keeps_transaction_blocked() { + let base = PoolMeta { + version: POOL_META_VERSION, + pools: vec![decommission_test_pool_status(0, None)], + ..Default::default() + }; + let mut current = base.clone(); + current.pools[0].last_update += Duration::seconds(1); + let first = Arc::new(PartialPoolMetaWriteStorage { + stored: StdMutex::new(Some((current.encode_config_data_for_v2_gate(true).unwrap(), "first".to_owned()))), + ..Default::default() + }); + let second = Arc::new(PartialPoolMetaWriteStorage { + stored: StdMutex::new(Some((base.encode_config_data_for_v2_gate(true).unwrap(), "second".to_owned()))), + ..Default::default() + }); + let mut state = super::PoolMetaWriteState::default(); + let selection = super::load_pool_meta_for_transaction_recovery(vec![first.clone(), second.clone()], &mut state) + .await + .unwrap(); + let raced = base.encode_config_data_for_v2_gate(true).unwrap(); + *second.stored.lock().unwrap() = Some((raced.clone(), "raced".to_owned())); + let result = super::repair_pool_meta_transaction( + vec![first.clone(), second.clone()], + &mut state, + selection, + &PoolMetaPersistenceFence::Distributed(None), + ) + .await; + assert!(matches!(result, Err(Error::PreconditionFailed))); + assert!(first.wrote.load(Ordering::SeqCst)); + assert!(!second.wrote.load(Ordering::SeqCst)); + assert_eq!(second.stored.lock().unwrap().as_ref().unwrap().0, raced); + let err = state.ensure_write_safe("recovery retry").unwrap_err(); + assert_eq!(err.pool_metadata_failure().unwrap().phase, "recovery_cas"); + } + + #[tokio::test] + async fn pool_meta_cas_failure_retains_the_original_source_for_caller_and_latch() { + let pool = Arc::new(PartialPoolMetaWriteStorage { + fail_write: true, + ..Default::default() + }); + let state = super::PoolMetaWriteState::default(); + let mut arm = state.arm_transaction(); + let err = save_pool_meta_object_cas( + pool, + POOL_META_NAME, + Vec::new(), + &PoolMetaCasToken::Missing, + &PoolMetaPersistenceFence::Distributed(None), + "commit_cas", + &mut arm, + ) + .await + .unwrap_err(); + let source = err.pool_metadata_failure().unwrap().source.as_ref().unwrap(); + assert!(matches!(source.as_ref(), Error::Timeout)); + assert!(Arc::ptr_eq(source, arm.source.as_ref().unwrap())); + drop(arm); + let blocked = state.ensure_write_safe("after CAS failure").unwrap_err(); + assert!(Arc::ptr_eq(source, blocked.pool_metadata_failure().unwrap().source.as_ref().unwrap())); + } + + #[tokio::test] + async fn pool_meta_cancelled_identity_dispatch_is_armed_before_pool_bin_writes() { + let pool = Arc::new(PartialPoolMetaWriteStorage { + pending_write: true, + ..Default::default() + }); + let state = Arc::new(tokio::sync::Mutex::new(super::PoolMetaWriteState::for_startup( + uuid::Uuid::new_v4(), + true, + ))); + let task = tokio::spawn({ + let pool = pool.clone(); + let state = state.clone(); + async move { + let mut state = state.lock().await; + super::persist_pool_meta_identity_for_startup(vec![pool], &mut state, false).await + } + }); + pool.write_started.notified().await; + task.abort(); + assert!(task.await.unwrap_err().is_cancelled()); + let state = state.lock().await; + let err = state.ensure_write_safe("identity retry").unwrap_err(); + assert_eq!(err.pool_metadata_failure().unwrap().phase, "identity_cas"); + assert_eq!( + err.pool_metadata_failure().unwrap().kind, + crate::error::PoolMetadataFailure::TransactionUnknown + ); + assert!(pool.stored.lock().unwrap().is_none()); + } + #[async_trait::async_trait] impl ObjectIO for PartialPoolMetaWriteStorage { type Error = Error; @@ -19745,12 +20743,22 @@ mod pools_tests { .expect("stale pool metadata should encode"); let fence = PoolMetaPersistenceFence::Distributed(None); - save_pool_meta_object_cas(storage.clone(), POOL_META_NAME, winner.clone(), &stale_token, &fence, "prepare_cas") - .await - .expect("the first writer should create pool metadata"); - let err = save_pool_meta_object_cas(storage.clone(), POOL_META_NAME, stale, &stale_token, &fence, "prepare_cas") - .await - .expect_err("the second writer must not reuse the stale missing-object revision"); + let mut arm = super::PoolMetaWriteState::default().arm_transaction(); + save_pool_meta_object_cas( + storage.clone(), + POOL_META_NAME, + winner.clone(), + &stale_token, + &fence, + "prepare_cas", + &mut arm, + ) + .await + .expect("the first writer should create pool metadata"); + let err = + save_pool_meta_object_cas(storage.clone(), POOL_META_NAME, stale, &stale_token, &fence, "prepare_cas", &mut arm) + .await + .expect_err("the second writer must not reuse the stale missing-object revision"); assert_eq!(err, Error::PreconditionFailed); let stored = storage @@ -20059,7 +21067,7 @@ mod pools_tests { 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"); + .expect("fence loss before dispatch has no unknown durable write and must remain retryable"); } #[tokio::test] diff --git a/crates/ecstore/src/error/mod.rs b/crates/ecstore/src/error/mod.rs index ec601df32..a482cea44 100644 --- a/crates/ecstore/src/error/mod.rs +++ b/crates/ecstore/src/error/mod.rs @@ -23,17 +23,59 @@ use s3s::S3ErrorCode; pub type Error = StorageError; pub type Result = core::result::Result; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum PoolMetadataFailure { + ReadUnavailable, + RecoveryRequired, + TransactionUnknown, + FenceLost, +} + +impl PoolMetadataFailure { + fn recovery_hint(self) -> &'static str { + match self { + Self::ReadUnavailable => "read unavailable; retry after the replicas are readable", + Self::TransactionUnknown => "writes remain blocked pending fenced transaction recovery", + Self::RecoveryRequired | Self::FenceLost => { + "writes remain blocked after a recovery-required replica state; restart after all replicas are readable and consistent, with compatible formats" + } + } + } + + pub fn as_str(self) -> &'static str { + match self { + Self::ReadUnavailable => "read_unavailable", + Self::RecoveryRequired => "recovery_required", + Self::TransactionUnknown => "transaction_unknown", + Self::FenceLost => "fence_lost", + } + } +} + +/// Local control-plane context. Keep the existing storage error wire codes; +/// the HTTP boundary recognizes this typed source, not an error-message prefix. +#[derive(Debug, Clone, thiserror::Error)] +#[error("{operation}: pool metadata {hint} ({reason}, {phase}): {detail}", hint = kind.recovery_hint(), reason = kind.as_str(), detail = source.as_ref().map(ToString::to_string).unwrap_or_default())] +pub struct PoolMetadataError { + pub kind: PoolMetadataFailure, + pub operation: String, + pub phase: &'static str, + pub since: time::OffsetDateTime, + #[source] + pub source: Option>, +} + /// Keeps high-cardinality diagnostic detail in the error source while making /// the rendered `io::Error` stable for quorum aggregation. #[derive(Debug)] struct StableIoContextError { - message: &'static str, + message: std::borrow::Cow<'static, str>, source: Box, } impl std::fmt::Display for StableIoContextError { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - formatter.write_str(self.message) + formatter.write_str(&self.message) } } @@ -48,7 +90,7 @@ where E: Into>, { std::io::Error::other(StableIoContextError { - message, + message: message.into(), source: source.into(), }) } @@ -300,6 +342,22 @@ impl From for StorageError { } impl StorageError { + pub fn pool_metadata_failure(&self) -> Option<&PoolMetadataError> { + let mut current: Option<&(dyn std::error::Error + 'static)> = Some(self); + while let Some(error) = current { + if let Some(context) = error.downcast_ref::() { + return Some(context); + } + // io::Error::source skips its boxed context itself. + current = if let Some(io) = error.downcast_ref::() { + io.get_ref().map(|inner| inner as &(dyn std::error::Error + 'static)) + } else { + error.source() + }; + } + None + } + pub fn other(error: E) -> Self where E: Into>, @@ -548,7 +606,19 @@ impl PartialEq for StorageError { impl Clone for StorageError { fn clone(&self) -> Self { match self { - StorageError::Io(e) => StorageError::Io(std::io::Error::new(e.kind(), e.to_string())), + StorageError::Io(e) => { + if let Some(context) = self.pool_metadata_failure() { + Self::Io(std::io::Error::new( + e.kind(), + StableIoContextError { + message: e.to_string().into(), + source: Box::new(context.clone()), + }, + )) + } else { + StorageError::Io(std::io::Error::new(e.kind(), e.to_string())) + } + } StorageError::FaultyDisk => StorageError::FaultyDisk, StorageError::DiskFull => StorageError::DiskFull, StorageError::VolumeNotFound => StorageError::VolumeNotFound, diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 6e472ce80..579a11f68 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -14,8 +14,8 @@ use super::*; use crate::core::pools::{ - PoolMetaBootstrapAuthority, PoolMetaReplicaState, PoolMetaWriteState, load_pool_meta_identity_observing, - local_decommission_queue_prefix, persist_pool_meta_identity_for_startup, pool_meta_has_active_decommission, + PoolMetaBootstrapAuthority, PoolMetaReplicaState, PoolMetaWriteState, local_decommission_queue_prefix, + persist_pool_meta_identity_for_startup, pool_meta_has_active_decommission, }; use crate::runtime::instance::InstanceContext; use crate::runtime::sources as runtime_sources; @@ -153,14 +153,11 @@ async fn load_pool_meta_for_startup( where S: EcstoreObjectIO, { - load_pool_meta_identity_observing(pools.clone(), write_state) - .await - .map_err(|err| Error::other(format!("store init failed during load_pool_meta_identity: {err}")))?; let mut meta = PoolMeta::default(); let replica_state = meta - .load_no_lock_from_replicas_observing(pools, write_state) + .load_for_startup_observing(pools, write_state) .await - .map_err(|err| Error::other(format!("store init failed during load_pool_meta: {err}")))?; + .map_err(|err| Error::other_with_context("store init failed during load_pool_meta", err))?; write_state.observe_replicas(replica_state); write_state .ensure_missing_metadata_can_initialize() @@ -769,6 +766,33 @@ impl ECStore { }); } + let recovery_store = self.clone(); + let recovery_rx = rx.clone(); + tokio::spawn(async move { + let mut delay = std::time::Duration::from_secs(5); + loop { + tokio::select! { + _ = recovery_rx.cancelled() => return, + _ = tokio::time::sleep(delay) => {} + } + let result = tokio::select! { + _ = recovery_rx.cancelled() => return, + result = tokio::time::timeout(std::time::Duration::from_secs(30), recovery_store.recover_pool_meta_transaction()) => result, + }; + delay = match result { + Ok(Ok(_)) => std::time::Duration::from_secs(5), + failure => { + let error = match failure { + Ok(Err(error)) => error, + _ => Error::Timeout, + }; + recovery_store.record_pool_meta_recovery_failure(error); + (delay * 2).min(std::time::Duration::from_secs(60)) + } + }; + } + }); + runtime_sources::init_bucket_monitor_for_current_endpoints(); crate::bucket::bucket_target_sys::BucketTargetSys::get().start_heartbeat(); @@ -2493,6 +2517,106 @@ mod tests { shutdown.cancel(); } + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn pool_metadata_preflight_recovery_preserves_single_and_multi_pool_public_mutations() { + for layout in [vec![4], vec![4, 4]] { + let temp_dir = tempfile::tempdir().unwrap(); + let (_ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "pool-meta-retry", &layout)).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await; + let bucket = format!("pool-meta-retry-{}", Uuid::new_v4()); + store.make_bucket(&bucket, &MakeBucketOptions::default()).await.unwrap(); + let mut saved_disks = Vec::new(); + for set in &store.pools[0].disk_set { + let mut disks = set.disks.write().await; + let count = disks.len(); + saved_disks.push((set.clone(), std::mem::replace(&mut *disks, vec![None; count]))); + } + let indices = (0..layout.len()).collect::>(); + let err = store.save_current_pool_meta_for_test(&indices).await.unwrap_err(); + assert_eq!( + err.pool_metadata_failure().unwrap().kind, + crate::error::PoolMetadataFailure::ReadUnavailable + ); + for (set, disks) in saved_disks { + *set.disks.write().await = disks; + } + store.save_current_pool_meta_for_test(&indices).await.unwrap(); + assert!(store.pool_meta_writes_ready().await); + + let payload = b"pool metadata recovery payload".to_vec(); + store + .put_object(&bucket, "put", &mut PutObjReader::from_vec(payload.clone()), &ObjectOptions::default()) + .await + .unwrap(); + let mut reader = store + .get_object_reader(&bucket, "put", None, HeaderMap::new(), &ObjectOptions::default()) + .await + .unwrap(); + let mut actual = Vec::new(); + reader.stream.read_to_end(&mut actual).await.unwrap(); + assert_eq!(actual, payload); + drop(reader); + store.delete_object(&bucket, "put", ObjectOptions::default()).await.unwrap(); + assert!(crate::error::is_err_object_not_found( + &store + .get_object_info(&bucket, "put", &ObjectOptions::default()) + .await + .unwrap_err() + )); + + let upload = store + .new_multipart_upload(&bucket, "multipart", &ObjectOptions::default()) + .await + .unwrap(); + let part = store + .put_object_part( + &bucket, + "multipart", + &upload.upload_id, + 1, + &mut PutObjReader::from_vec(payload.clone()), + &ObjectOptions::default(), + ) + .await + .unwrap(); + store + .clone() + .complete_multipart_upload( + &bucket, + "multipart", + &upload.upload_id, + vec![crate::storage_api_contracts::multipart::CompletePart { + part_num: part.part_num, + etag: part.etag, + ..Default::default() + }], + &ObjectOptions::default(), + ) + .await + .unwrap(); + let mut reader = store + .get_object_reader(&bucket, "multipart", None, HeaderMap::new(), &ObjectOptions::default()) + .await + .unwrap(); + actual.clear(); + reader.stream.read_to_end(&mut actual).await.unwrap(); + assert_eq!(actual, payload); + drop(reader); + let upload = store + .new_multipart_upload(&bucket, "abort", &ObjectOptions::default()) + .await + .unwrap(); + store + .abort_multipart_upload(&bucket, "abort", &upload.upload_id, &ObjectOptions::default()) + .await + .unwrap(); + assert!(store.pool_meta_writes_ready().await); + shutdown.cancel(); + } + } + #[cfg(feature = "test-util")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[serial_test::serial(storage_class_env)] diff --git a/docs/architecture/readiness-matrix.md b/docs/architecture/readiness-matrix.md index caaa33c03..45e396101 100644 --- a/docs/architecture/readiness-matrix.md +++ b/docs/architecture/readiness-matrix.md @@ -40,6 +40,9 @@ an unknown or unsupported peer-health snapshot degrades readiness with - Liveness reports process availability and must not depend on storage, IAM, lock quorum, or peer health. - Node readiness reports local dependency readiness. +- A blocked pool metadata writer degrades node and cluster-write readiness with + `pool_metadata_blocked`. Metadata save-gate inspection is bounded to 100 ms; + contention reports `pool_metadata_check_timeout` without installing a block. - Cluster write readiness requires write quorum and the runtime dependency readiness used by `FullReady`. - Cluster read readiness may use the read-quorum path and cluster-health diff --git a/docs/operations/pool-metadata-recovery.md b/docs/operations/pool-metadata-recovery.md index ad594acc2..0cb2513ff 100644 --- a/docs/operations/pool-metadata-recovery.md +++ b/docs/operations/pool-metadata-recovery.md @@ -32,6 +32,27 @@ A V3 update first conditionally writes a pending generation containing the last Do not hand-edit a pending record or select a replica only because it is in pool zero. Preserve all copies when escalating recovery. +## Runtime write recovery + +A runtime metadata or identity read failure before any write is dispatched rejects that operation but does not permanently block subsequent retries. Errors retain their typed cause; pool metadata unavailability reaches S3 as `503 ServiceUnavailable`, without exposing internal error details. Cancellation during preflight or a rejected first conditional write is also retryable. After any identity, prepare, or commit write may have started, cancellation, an uncertain write result, or abandoned runtime publication blocks further metadata-dependent mutations. + +The node checks for interrupted pool metadata transactions every five seconds. Failed recovery attempts back off to at most sixty seconds; each attempt has a thirty-second budget and stops on shutdown. Healthy nodes do not read metadata for this worker. Recovery: + +1. Cancels old decommission workers and waits for their supervisors to drain them. It does not cancel a separate rebalance operation; an attached rebalance worker prevents recovery until quiescent. +2. Holds the local start/movement gates and distributed `pool.bin` write fence, validates the initialized deployment identity and unchanged pool topology, and selects the authoritative durable transaction. +3. Repairs pending, missing, or lagging copies using conditional writes, then rereads and verifies convergence. A prepare-only first V3 migration commits the predecessor as V3, preserving the observed format floor. +4. Invalidates old movement snapshots, installs the verified durable state, rechecks the fence, and only then clears the block. Speculative in-memory progress is never used as the recovery source. The existing decommission supervisor resumes eligible work afterward. + +An unreadable replica, lost fence, or conditional-write conflict leaves the original block in place. Recovery never initializes an all-missing metadata set. Corruption, incompatible layouts, conflicting identities/epochs/transactions, and topology changes require operator reconciliation; restore readability and consistency using the procedures below. Blocks originating in startup validation or storage-format heal are not cleared by the pool transaction worker: restart only after repairing the underlying condition. There is no force-clear switch. If an attached rebalance worker cannot quiesce, collect its status and restart the affected node after verifying the durable metadata; do not manually detach its worker token. + +### Diagnostics + +- The first block emits `decommission_state` with `state=pool_metadata_blocked`, `reason`, `phase`, and `blocked_since`. A change in recovery failure classification emits `state=pool_metadata_recovery_pending`; successful recovery emits `state=pool_metadata_recovered` with the original timestamp. +- `rustfs_pool_metadata_blocks_total{reason}` and `rustfs_pool_metadata_recoveries_total` count block and recovery transitions. The original cause and phase remain attached to local typed errors; storage/RPC error numbers and on-disk formats are unchanged. +- Node and cluster-write readiness include `pool_metadata_blocked`. Waiting for the metadata save mutex is bounded to 100 ms and reports `pool_metadata_check_timeout`, not a persistent block. Cluster probes retain their existing cache and overall timeout behavior. Liveness and cluster-read quorum checks are unchanged. + +If a block persists, inspect the first block and subsequent recovery phase, restore disk/peer readability, and verify every metadata and identity copy before restarting. Do not delete metadata to make readiness green. + ## Disk replacement and metadata erasure 1. Keep a quorum of nodes online and verify the cluster is ready. @@ -39,4 +60,4 @@ Do not hand-edit a pending record or select a replica only because it is in pool 3. Restore storage formats and the `pool.bin.identity` marker from the same deployment before rejoining it. 4. Start the node and wait for it to load the verified committed generation and repair its replicas before touching another node. -An initialized identity with every `pool.bin` missing is recovery required, as are existing storage formats with neither identity nor `pool.bin`. Format creation alone is not fresh-cluster proof: only the elected first topology node may create a durable `initialized=false` bootstrap identity with a fresh-bootstrap nonce, and only after every configured disk explicitly responds that it is unformatted. An unreachable peer, a non-elected distributed node, or an existing format is not sufficient proof. All-missing `pool.bin` replicas are accepted only by the same startup that proved the fresh topology and persisted that pending identity; when every `pool.bin` is missing, a later startup must recover even if the pending identity survived. This prevents a wiped or lagging node from rebuilding empty state and overwriting the cluster. Runtime reload, rebalance activation, and rebalance worker admission all fail closed and latch the same recovery gate until the node is restarted with readable metadata. +An initialized identity with every `pool.bin` missing is recovery required, as are existing storage formats with neither identity nor `pool.bin`. Format creation alone is not fresh-cluster proof: only the elected first topology node may create a durable `initialized=false` bootstrap identity with a fresh-bootstrap nonce, and only after every configured disk explicitly responds that it is unformatted. An unreachable peer, a non-elected distributed node, or an existing format is not sufficient proof. All-missing `pool.bin` replicas are accepted only by the same startup that proved the fresh topology and persisted that pending identity; when every `pool.bin` is missing, a later startup must recover even if the pending identity survived. This prevents a wiped or lagging node from rebuilding empty state and overwriting the cluster. Runtime reload, rebalance activation, and rebalance worker admission fail closed on this missing-authority condition; a clean probe alone cannot clear it. diff --git a/rustfs/src/error.rs b/rustfs/src/error.rs index d7e7e34b5..a95e743c2 100644 --- a/rustfs/src/error.rs +++ b/rustfs/src/error.rs @@ -333,6 +333,13 @@ impl From for S3Error { impl From for ApiError { fn from(err: StorageError) -> Self { + if err.pool_metadata_failure().is_some() { + return ApiError { + code: S3ErrorCode::ServiceUnavailable, + message: ApiError::error_code_to_message(&S3ErrorCode::ServiceUnavailable), + source: Some(Box::new(err)), + }; + } if let StorageError::Io(ref io_err) = err && let Some(inner) = io_err.get_ref() { @@ -848,6 +855,35 @@ mod tests { assert_eq!(api_error.message, "The service is unavailable. Please retry."); } + #[test] + fn pool_metadata_failures_map_to_503_and_preserve_typed_private_context() { + use crate::storage_api::error::{PoolMetadataError, PoolMetadataFailure}; + for kind in [ + PoolMetadataFailure::ReadUnavailable, + PoolMetadataFailure::RecoveryRequired, + PoolMetadataFailure::TransactionUnknown, + PoolMetadataFailure::FenceLost, + ] { + let error = StorageError::other(PoolMetadataError { + kind, + operation: "pool metadata test".to_owned(), + phase: "prepare_cas", + since: time::OffsetDateTime::now_utc(), + source: Some(std::sync::Arc::new(StorageError::other("private disk failure"))), + }); + let error = StorageError::Io(std::io::Error::new(std::io::ErrorKind::TimedOut, error)); + let cloned = error.clone(); + assert_eq!(cloned, error, "cloning must preserve the outer I/O kind and message"); + assert_eq!(cloned.pool_metadata_failure().unwrap().kind, kind); + let api = ApiError::from(cloned); + assert_eq!(api.code, S3ErrorCode::ServiceUnavailable); + assert_eq!(api.message, "The service is unavailable. Please retry."); + assert!(!api.message.contains("private")); + let source = api.source.as_ref().unwrap().downcast_ref::().unwrap(); + assert_eq!(source.pool_metadata_failure().unwrap().kind, kind); + } + } + #[test] fn test_unknown_authoritative_quota_usage_maps_to_retryable_error() { let api_error = ApiError::from(QuotaError::UsageUnavailable { diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 34210287d..67cf8938c 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -479,6 +479,8 @@ pub(crate) mod ecstore_error { pub(crate) use rustfs_ecstore::api::error::{ Error, Result, StorageError, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found, }; + #[cfg(test)] + pub(crate) use rustfs_ecstore::api::error::{PoolMetadataError, PoolMetadataFailure}; } pub(crate) mod ecstore_event { diff --git a/rustfs/src/storage_api.rs b/rustfs/src/storage_api.rs index 62f883363..d120e2ef0 100644 --- a/rustfs/src/storage_api.rs +++ b/rustfs/src/storage_api.rs @@ -81,6 +81,8 @@ pub(crate) mod error { } } + #[cfg(test)] + pub(crate) use crate::storage::storage_api::ecstore_error::{PoolMetadataError, PoolMetadataFailure}; pub(crate) use crate::storage::storage_api::{QuotaError, StorageError}; }