fix(ecstore): version pool metadata transactions (#6604)

* fix(connect): adapt offline array predicate

* test(e2e): update smoke selection baseline

* test(ecstore): make slowtail oracle deterministic

* test(get): stage relocated fixture after reader opens

* ci: bound feature test link concurrency

* test: give lifecycle transition futures a larger stack

* fix(ecstore): version pool metadata transactions
This commit is contained in:
Zhengchao An
2026-08-26 09:33:51 +08:00
committed by GitHub
parent c0c5fc22f9
commit 1bcb396752
14 changed files with 3429 additions and 359 deletions
+22
View File
@@ -181,6 +181,22 @@ pub const DEFAULT_POOL_META_V2_FLEET_CONFIRMED: bool = false;
const _: () = assert!(!DEFAULT_POOL_META_V2_WRITE);
const _: () = assert!(!DEFAULT_POOL_META_V2_FLEET_CONFIRMED);
/// Request writing pool metadata version 3 with durable generations.
///
/// Existing deployments remain on their observed version until
/// [`ENV_POOL_META_V3_FLEET_CONFIRMED`] is also enabled. Fresh deployments may
/// initialize directly at version 3 because they have no legacy readers.
pub const ENV_POOL_META_V3_WRITE: &str = "RUSTFS_POOL_META_V3_WRITE";
pub const DEFAULT_POOL_META_V3_WRITE: bool = false;
/// Operator-attested confirmation that every pool metadata reader and writer
/// understands the version 3 generation and recovery protocol.
pub const ENV_POOL_META_V3_FLEET_CONFIRMED: &str = "RUSTFS_POOL_META_V3_FLEET_CONFIRMED";
pub const DEFAULT_POOL_META_V3_FLEET_CONFIRMED: bool = false;
const _: () = assert!(!DEFAULT_POOL_META_V3_WRITE);
const _: () = assert!(!DEFAULT_POOL_META_V3_FLEET_CONFIRMED);
// =============================================================================
// Concurrent Request Fix - Timeout and Backpressure Configuration
// =============================================================================
@@ -755,4 +771,10 @@ mod remote_version_state_tests {
assert_eq!(super::ENV_POOL_META_V2_WRITE, "RUSTFS_POOL_META_V2_WRITE");
assert_eq!(super::ENV_POOL_META_V2_FLEET_CONFIRMED, "RUSTFS_POOL_META_V2_FLEET_CONFIRMED");
}
#[test]
fn pool_meta_v3_gate_uses_stable_environment_names() {
assert_eq!(super::ENV_POOL_META_V3_WRITE, "RUSTFS_POOL_META_V3_WRITE");
assert_eq!(super::ENV_POOL_META_V3_FLEET_CONFIRMED, "RUSTFS_POOL_META_V3_FLEET_CONFIRMED");
}
}
+6 -1
View File
@@ -704,7 +704,12 @@ where
save_config_with_opts_inner(api, file, data, opts, false).await.map(|_| ())
}
async fn save_config_with_opts_and_metadata<S>(api: Arc<S>, file: &str, data: Vec<u8>, opts: &ObjectOptions) -> Result<ObjectInfo>
pub(crate) async fn save_config_with_opts_and_metadata<S>(
api: Arc<S>,
file: &str,
data: Vec<u8>,
opts: &ObjectOptions,
) -> Result<ObjectInfo>
where
S: ObjectIO<
Error = Error,
File diff suppressed because it is too large Load Diff
+103 -2
View File
@@ -106,6 +106,89 @@ impl Drop for Sets {
}
}
#[cfg(test)]
struct HealFormatAfterSaveBarrierState {
pool_key: usize,
disk_index: usize,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
static HEAL_FORMAT_AFTER_SAVE_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<HealFormatAfterSaveBarrierState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
pub(crate) struct HealFormatAfterSaveBarrier {
state: Arc<HealFormatAfterSaveBarrierState>,
}
#[cfg(test)]
impl HealFormatAfterSaveBarrier {
pub(crate) fn install(pool: &Arc<Sets>, disk_index: usize) -> Self {
let state = Arc::new(HealFormatAfterSaveBarrierState {
pool_key: Arc::as_ptr(pool) as usize,
disk_index,
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut barrier = HEAL_FORMAT_AFTER_SAVE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("heal format after-save barrier should not be poisoned");
assert!(barrier.is_none(), "heal format after-save barrier must be unique");
*barrier = Some(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(std::time::Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("format heal should reach the after-save barrier");
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for HealFormatAfterSaveBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut barrier = HEAL_FORMAT_AFTER_SAVE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("heal format after-save barrier should not be poisoned");
if barrier.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*barrier = None;
}
}
}
#[cfg(test)]
async fn pause_heal_format_after_save(pool: &Sets, disk_index: usize) {
let pool_key = std::ptr::from_ref(pool) as usize;
let barrier = {
let mut barrier = HEAL_FORMAT_AFTER_SAVE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("heal format after-save barrier should not be poisoned");
if barrier
.as_ref()
.is_some_and(|state| state.pool_key == pool_key && state.disk_index == disk_index)
{
barrier.take()
} else {
None
}
};
if let Some(barrier) = barrier {
barrier.arrived.notify_one();
barrier.release.notified().await;
}
}
impl Sets {
#[tracing::instrument(level = "debug", skip(disks, endpoints, fm, pool_idx, parity_count))]
pub async fn new(
@@ -989,9 +1072,13 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for Sets {
}
impl Sets {
pub(crate) async fn heal_format_with_fence<F>(&self, dry_run: bool, fence_lost: F) -> Result<(HealResultItem, Option<Error>)>
pub(crate) async fn heal_format_with_fence<F>(
&self,
dry_run: bool,
mut fence_lost: F,
) -> Result<(HealResultItem, Option<Error>)>
where
F: Fn() -> bool + Send + Sync,
F: FnMut() -> bool + Send,
{
let (disks, init_errs) = init_storage_disks_with_errors(
&self.endpoints.endpoints,
@@ -1074,6 +1161,11 @@ impl Sets {
}
return Ok((res, Some(err.into())));
}
#[cfg(test)]
pause_heal_format_after_save(self, index).await;
if fence_lost() {
return Ok((res, Some(StorageError::SlowDown)));
}
if let Some(saved_format) = fm.as_ref() {
res.after.drives[index].uuid = saved_format.erasure.this.to_string();
res.after.drives[index].state = DriveState::Ok.to_string();
@@ -1081,8 +1173,14 @@ impl Sets {
}
}
if fence_lost() {
return Ok((res, Some(StorageError::SlowDown)));
}
for (index, fm) in tmp_new_formats.iter().enumerate() {
if let Some(fm) = fm {
if fence_lost() {
return Ok((res, Some(StorageError::SlowDown)));
}
let (m, n) = match ref_format.find_disk_index_by_disk_id(fm.erasure.this) {
Ok((m, n)) => (m, n),
Err(_) => continue,
@@ -1094,6 +1192,9 @@ impl Sets {
}
if let Some(Some(disk)) = disks.get(index) {
if fence_lost() {
return Ok((res, Some(StorageError::SlowDown)));
}
self.disk_set[m].renew_disk(&disk.endpoint()).await;
}
}
@@ -540,9 +540,12 @@ impl ECStore {
#[cfg(test)]
crate::core::pools::observe_pool_activation_start_attempt(crate::core::pools::PoolActivationStartKind::Rebalance);
let fleet_proof = acquire_pool_activation_fleet_proof(&self.ctx).await?;
let mut pool_meta_guard = self.pool_meta_save_gate.lock().await;
pool_meta_guard.ensure_write_safe(stage)?;
let activation_fence = acquire_pool_rebalance_activation_locks(pool.clone(), fleet_proof).await?;
let mut pool_meta = PoolMeta::default();
pool_meta.load_no_lock_from_replicas(self.pools.clone()).await?;
let pool_meta = self
.load_runtime_pool_meta_under_activation_fence(&mut pool_meta_guard, &activation_fence, stage)
.await?;
ensure_rebalance_activation_pool_meta_allowed(&pool_meta)?;
merge_and_save_rebalance_meta_no_lock(
@@ -568,9 +571,12 @@ impl ECStore {
S: EcstoreObjectIO + StorageNamespaceLocking<Error = Error, NamespaceLock = rustfs_lock::NamespaceLockWrapper>,
{
let fleet_proof = acquire_pool_activation_fleet_proof(&self.ctx).await?;
let mut pool_meta_guard = self.pool_meta_save_gate.lock().await;
pool_meta_guard.ensure_write_safe("rebalance worker activation")?;
let activation_fence = acquire_pool_rebalance_activation_locks(pool.clone(), fleet_proof).await?;
let mut pool_meta = PoolMeta::default();
pool_meta.load_no_lock_from_replicas(self.pools.clone()).await?;
let pool_meta = self
.load_runtime_pool_meta_under_activation_fence(&mut pool_meta_guard, &activation_fence, "rebalance worker activation")
.await?;
ensure_rebalance_activation_pool_meta_allowed(&pool_meta)?;
let mut persisted = RebalanceMeta::new();
@@ -1301,10 +1307,40 @@ impl ECStore {
#[cfg(test)]
mod tests {
use super::*;
use crate::core::pools::{PoolActivationDurableSaveBarrier, PoolActivationStartKind, PoolActivationStartProbe};
use crate::config::com::delete_config;
use crate::core::pools::{
POOL_META_NAME, PoolActivationDurableSaveBarrier, PoolActivationStartKind, PoolActivationStartProbe, PoolMetaWriteState,
persist_pool_meta_identity_for_startup,
};
use crate::object_api::NamespaceLockFence;
use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause, hermetic_set_disks_isolated};
async fn persist_initialized_identity_then_remove_pool_meta(store: &Arc<ECStore>) {
let mut write_state = PoolMetaWriteState::for_startup(store.id, false);
persist_pool_meta_identity_for_startup(store.pools.clone(), &mut write_state, true)
.await
.expect("initialized pool metadata identity should persist");
*store.pool_meta_save_gate.lock().await = write_state;
for pool in &store.pools {
delete_config(pool.clone(), POOL_META_NAME)
.await
.expect("every pool metadata replica should be removed");
}
}
async fn assert_activation_locks_released(store: &Arc<ECStore>) {
let fleet_proof = acquire_pool_activation_fleet_proof(&store.ctx)
.await
.expect("fleet proof should remain available");
tokio::time::timeout(
std::time::Duration::from_secs(5),
acquire_pool_rebalance_activation_locks(store.pools[0].clone(), fleet_proof),
)
.await
.expect("a rejected activation must release namespace fences promptly")
.expect("a rejected activation must release both namespace fences");
}
#[tokio::test]
async fn rebalance_stop_wait_probe_matches_run_id() {
let probe = RebalanceStopWaitProbe::install("rebalance-stop-current");
@@ -1356,6 +1392,90 @@ mod tests {
assert!(cancel.is_cancelled());
}
#[tokio::test]
#[serial_test::serial]
async fn rebalance_activation_rejects_initialized_cluster_with_all_pool_meta_missing() {
let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await;
persist_initialized_identity_then_remove_pool_meta(&store).await;
set_rebalance_disk_stats_override_for_test(
store.id,
vec![
DiskStat {
total_space: 100,
available_space: 0,
},
DiskStat {
total_space: 100,
available_space: 100,
},
],
);
let err = store
.init_rebalance_start(vec!["missing-pool-meta".to_string()])
.await
.expect_err("rebalance activation must fail closed when every pool.bin is missing");
assert!(err.to_string().contains("initialized cluster identity exists"));
store
.ensure_pool_meta_side_effects_safe("rebalance activation after missing pool metadata")
.await
.expect_err("the missing metadata observation must latch the shared runtime gate");
let mut persisted = RebalanceMeta::new();
assert!(
matches!(persisted.load(store.pools[0].clone()).await, Err(Error::ConfigNotFound)),
"rejected activation must not create rebalance metadata"
);
assert_activation_locks_released(&store).await;
}
#[tokio::test]
#[serial_test::serial]
async fn rebalance_worker_rejects_initialized_cluster_with_all_pool_meta_missing() {
let rebalance_id = "missing-pool-meta-worker";
let active = RebalanceMeta {
id: rebalance_id.to_string(),
percent_free_goal: 0.5,
pool_stats: vec![
RebalanceStats {
participating: true,
init_capacity: 100,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
},
RebalanceStats {
participating: true,
init_capacity: 100,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
},
],
..Default::default()
};
let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(Some(active)).await;
persist_initialized_identity_then_remove_pool_meta(&store).await;
let err = match store
.fence_rebalance_worker_activation(store.pools[0].clone(), rebalance_id)
.await
{
Ok(_) => panic!("worker activation must not return a fence when every pool.bin is missing"),
Err(err) => err,
};
assert!(err.to_string().contains("initialized cluster identity exists"));
store
.ensure_pool_meta_side_effects_safe("rebalance worker after missing pool metadata")
.await
.expect_err("worker validation must latch the shared runtime gate");
assert_activation_locks_released(&store).await;
}
#[tokio::test]
#[serial_test::serial]
async fn rebalance_activation_adopts_commit_after_post_save_fence_loss() {
+9 -3
View File
@@ -63,6 +63,12 @@ pub async fn test_store_with_persisted_rebalance_meta(
) -> (Vec<tempfile::TempDir>, std::sync::Arc<crate::store::ECStore>) {
let ctx = std::sync::Arc::new(crate::runtime::instance::InstanceContext::new());
let (temp_dirs, pool) = crate::core::sets::make_local_two_set_sets_with_ctx(ctx.clone()).await;
let pools = vec![pool.clone()];
let pool_meta = crate::core::pools::PoolMeta::new(&pools, &crate::core::pools::PoolMeta::default());
pool_meta
.save_for_startup(pools.clone())
.await
.expect("rebalance test pool metadata should be persisted");
meta.save(pool.clone())
.await
.expect("rebalance test metadata should be persisted");
@@ -70,9 +76,9 @@ pub async fn test_store_with_persisted_rebalance_meta(
let store = std::sync::Arc::new(crate::store::ECStore {
id: uuid::Uuid::new_v4(),
disk_map: std::collections::HashMap::new(),
pools: vec![pool],
pools,
peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, ctx.clone()),
pool_meta: tokio::sync::RwLock::new(crate::core::pools::PoolMeta::default()),
pool_meta: tokio::sync::RwLock::new(pool_meta),
rebalance_meta: tokio::sync::RwLock::new(Some(meta)),
decommission_cancelers: tokio::sync::RwLock::new(vec![None]),
start_gate: tokio::sync::Mutex::new(()),
@@ -139,7 +145,7 @@ async fn test_two_pool_stores_with_contexts(
}
let pool_meta = PoolMeta::new(&pools, &PoolMeta::default());
pool_meta
.save(pools.clone())
.save_for_startup(pools.clone())
.await
.expect("baseline pool metadata should be persisted");
if let Some(meta) = rebalance_meta.as_ref() {
@@ -2793,7 +2793,10 @@ async fn test_rebalance_start_save_failure_retries_persisted_completed_state() {
.start_rebalance_under_gate()
.await
.expect_err("the injected first activation save must fail through the real start path");
assert!(err.to_string().contains("injected rebalance activation save failure"));
assert!(
err.to_string().contains("injected rebalance activation save failure"),
"unexpected activation error: {err}"
);
{
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
@@ -2854,7 +2857,10 @@ async fn test_rebalance_start_save_failure_retries_persisted_stopped_state() {
.start_rebalance_under_gate()
.await
.expect_err("the injected first stopped-state save must fail through the real start path");
assert!(err.to_string().contains("injected rebalance activation save failure"));
assert!(
err.to_string().contains("injected rebalance activation save failure"),
"unexpected activation error: {err}"
);
let mut after_failure = RebalanceMeta::new();
after_failure
.load(store.pools[0].clone())
+27 -4
View File
@@ -2034,11 +2034,24 @@ impl SetDisks {
}
impl SetDisks {
#[cfg(test)]
pub(crate) async fn heal_replacement_format(
&self,
dry_run: bool,
targets: &[String],
) -> Result<(HealResultItem, Option<Error>)> {
self.heal_replacement_format_with_fence(dry_run, targets, || false).await
}
pub(crate) async fn heal_replacement_format_with_fence<F>(
&self,
dry_run: bool,
targets: &[String],
fence_lost: F,
) -> Result<(HealResultItem, Option<Error>)>
where
F: FnMut() -> bool + Send,
{
if targets.is_empty() {
return Err(Error::other("replacement format requires at least one target"));
}
@@ -2054,14 +2067,18 @@ impl SetDisks {
target_slots.push(slot);
}
self.heal_format_for_slots(dry_run, Some(&target_slots)).await
self.heal_format_for_slots(dry_run, Some(&target_slots), fence_lost).await
}
async fn heal_format_for_slots(
async fn heal_format_for_slots<F>(
&self,
dry_run: bool,
target_slots: Option<&[usize]>,
) -> Result<(HealResultItem, Option<Error>)> {
mut fence_lost: F,
) -> Result<(HealResultItem, Option<Error>)>
where
F: FnMut() -> bool + Send,
{
let disks = self.disks.read().await.clone();
let (formats, errs) = load_format_erasure_all(&disks, true).await;
if errs.iter().any(|err| {
@@ -2131,8 +2148,14 @@ impl SetDisks {
let mut new_format = ref_format.clone();
new_format.erasure.this = ref_format.erasure.sets[self.set_index][disk_idx];
if fence_lost() {
return Ok((result, Some(StorageError::SlowDown)));
}
match save_format_file(&disks[disk_idx], &Some(new_format.clone())).await {
Ok(()) => {
if fence_lost() {
return Ok((result, Some(StorageError::SlowDown)));
}
result.after.drives[disk_idx].uuid = new_format.erasure.this.to_string();
result.after.drives[disk_idx].state = DriveState::Ok.to_string();
}
@@ -2158,7 +2181,7 @@ impl crate::storage_api_contracts::heal::HealOperations for SetDisks {
#[tracing::instrument(skip(self))]
async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
self.heal_format_for_slots(dry_run, None).await
self.heal_format_for_slots(dry_run, None, || false).await
}
#[tracing::instrument(skip(self))]
+259 -11
View File
@@ -13,7 +13,7 @@
// limitations under the License.
use super::*;
use crate::core::pools::POOL_META_NAME;
use crate::core::pools::{POOL_META_NAME, load_pool_meta_identity_observing};
use crate::services::rebalance::{REBAL_META_NAME, RebalStatus};
use crate::set_disk::get_lock_acquire_timeout;
use crate::storage_api_contracts::heal::HealOperations as _;
@@ -91,7 +91,13 @@ fn heal_format_fence_lost_error() -> Error {
impl ECStore {
async fn acquire_heal_format_fence(
&self,
) -> Result<(NamespaceLockGuard, NamespaceLockGuard, PoolMeta, Option<RebalanceMeta>)> {
) -> Result<(
tokio::sync::MutexGuard<'_, PoolMetaWriteState>,
NamespaceLockGuard,
NamespaceLockGuard,
PoolMeta,
Option<RebalanceMeta>,
)> {
let metadata_pool = self
.pools
.first()
@@ -111,13 +117,14 @@ impl ECStore {
return Err(heal_format_fence_lost_error());
}
load_pool_meta_identity_observing(self.pools.clone(), &mut write_state).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_missing_metadata_can_initialize()?;
write_state.ensure_write_safe("heal format fence failed")?;
drop(write_state);
if pool_meta.pools.len() != self.pools.len()
|| pool_meta.pools.iter().enumerate().any(|(pool_idx, pool)| {
pool.id != pool_idx || pool.cmd_line.is_empty() || pool.cmd_line != self.pools[pool_idx].endpoints.cmd_line
@@ -149,11 +156,12 @@ impl ECStore {
return Err(heal_format_fence_lost_error());
}
write_state.ensure_write_safe("heal format fence failed")?;
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
return Err(heal_format_fence_lost_error());
}
Ok((pool_guard, rebalance_guard, pool_meta, rebalance_meta))
Ok((write_state, pool_guard, rebalance_guard, pool_meta, rebalance_meta))
}
fn get_pools_for_heal_object(&self, opts: &HealOpts) -> Result<Vec<Arc<Sets>>> {
@@ -234,7 +242,8 @@ impl ECStore {
let mut count_completed = 0;
let mut first_error = None;
for (pool_idx, pool) in self.pools.iter().enumerate() {
let (pool_guard, rebalance_guard, pool_meta, rebalance_meta) = self.acquire_heal_format_fence().await?;
let (mut write_state, pool_guard, rebalance_guard, pool_meta, rebalance_meta) =
self.acquire_heal_format_fence().await?;
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
first_error.get_or_insert(heal_format_fence_lost_error());
break;
@@ -249,7 +258,15 @@ impl ECStore {
continue;
}
let fence_lost = || pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost();
let fence_lost = || {
let lost = pool_guard.is_lock_lost()
|| rebalance_guard.is_lock_lost()
|| write_state.ensure_write_safe("heal format write fence failed").is_err();
if lost {
write_state.block_writes_after_fence_loss();
}
lost
};
let (mut result, err) = pool.heal_format_with_fence(dry_run, fence_lost).await?;
if let Some(err) = err {
match err {
@@ -268,7 +285,11 @@ impl ECStore {
// A lease can be lost after the final write; fail closed before
// reporting the pool as successfully healed.
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
let fence_lost = pool_guard.is_lock_lost()
|| rebalance_guard.is_lock_lost()
|| write_state.ensure_write_safe("heal format publication fence failed").is_err();
if fence_lost {
write_state.block_writes_after_fence_loss();
first_error.get_or_insert(heal_format_fence_lost_error());
break;
}
@@ -321,7 +342,33 @@ impl ECStore {
)
})?;
set.heal_replacement_format(dry_run, targets).await
let (mut write_state, pool_guard, rebalance_guard, pool_meta, rebalance_meta) = self.acquire_heal_format_fence().await?;
if let Some(skip) = classify_heal_format_pool(pool_index, &pool.endpoints.cmd_line, &pool_meta, rebalance_meta.as_ref()) {
return Ok((HealResultItem::default(), Some(heal_format_pool_skip_error(skip))));
}
let fence_lost = || {
let lost = pool_guard.is_lock_lost()
|| rebalance_guard.is_lock_lost()
|| write_state
.ensure_write_safe("replacement format write fence failed")
.is_err();
if lost {
write_state.block_writes_after_fence_loss();
}
lost
};
let result = set.heal_replacement_format_with_fence(dry_run, targets, fence_lost).await?;
let fence_lost = pool_guard.is_lock_lost()
|| rebalance_guard.is_lock_lost()
|| write_state
.ensure_write_safe("replacement format publication fence failed")
.is_err();
if fence_lost {
write_state.block_writes_after_fence_loss();
return Ok((result.0, Some(heal_format_fence_lost_error())));
}
Ok(result)
}
#[instrument(skip(self, targets), fields(pool_index, set_index, target_count = targets.len()))]
@@ -545,9 +592,13 @@ mod tests {
use super::*;
use crate::bucket::metadata_sys;
use crate::cluster::rpc::PeerS3Client;
use crate::core::pools::{PoolDecommissionInfo, PoolMetaReplicaState, PoolStatus};
use crate::config::com::{delete_config, read_config_no_lock_preserve_empty_with_metadata, save_config};
use crate::core::pools::{
POOL_META_IDENTITY_NAME, PoolDecommissionInfo, PoolMetaReplicaState, PoolStatus, initialized_pool_meta_identity_for_test,
};
use crate::core::sets::HealFormatAfterSaveBarrier;
use crate::disk::error::Result as DiskResult;
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk};
use crate::disk::{DeleteOptions, DiskOption, FORMAT_CONFIG_FILE, format::FormatV3, new_disk};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::runtime::instance::InstanceContext;
use crate::services::rebalance::{RebalanceInfo, RebalanceStats};
@@ -557,6 +608,7 @@ mod tests {
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 rustfs_common::heal_channel::DriveState;
use tokio_util::sync::CancellationToken;
#[derive(Debug)]
@@ -933,6 +985,197 @@ mod tests {
(temp_dir, store, shutdown)
}
fn heal_test_format_path(temp_dir: &tempfile::TempDir, pool_index: usize, disk_index: usize) -> std::path::PathBuf {
temp_dir
.path()
.join(format!("pool{pool_index}-disk{disk_index}"))
.join(crate::disk::RUSTFS_META_BUCKET)
.join(FORMAT_CONFIG_FILE)
}
async fn remove_heal_test_format(
temp_dir: &tempfile::TempDir,
store: &ECStore,
pool_index: usize,
disk_index: usize,
) -> String {
let target = store.pools[pool_index].endpoints.endpoints.as_ref()[disk_index].to_string();
let format_path = heal_test_format_path(temp_dir, pool_index, disk_index);
tokio::fs::remove_file(&format_path)
.await
.expect("replacement target format should be removable");
assert!(
!tokio::fs::try_exists(&format_path)
.await
.expect("replacement target format path should be inspectable")
);
target
}
async fn assert_heal_test_format_missing(temp_dir: &tempfile::TempDir, pool_index: usize, disk_index: usize) {
assert!(
!tokio::fs::try_exists(heal_test_format_path(temp_dir, pool_index, disk_index))
.await
.expect("replacement target format path should be inspectable"),
"format heal must not write the replacement target"
);
}
#[tokio::test]
#[serial_test::serial]
async fn full_format_heal_preblocked_pool_metadata_never_writes_format() {
let (temp_dir, store, shutdown) = multi_pool_heal_store().await;
remove_heal_test_format(&temp_dir, &store, 0, 3).await;
store.pool_meta_save_gate.lock().await.observe_replicas(PoolMetaReplicaState {
needs_repair: true,
repair_write_safe: false,
});
let err = store
.handle_heal_format(false)
.await
.expect_err("a preblocked pool metadata state must reject full format heal");
assert!(
err.to_string()
.contains("restart after all replicas are readable and consistent")
);
assert_heal_test_format_missing(&temp_dir, 0, 3).await;
store
.ensure_pool_meta_side_effects_safe("preblocked format heal side effect")
.await
.expect_err("the preblocked state must remain sticky");
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn full_format_heal_future_identity_never_writes_format_and_latches() {
let (temp_dir, store, shutdown) = multi_pool_heal_store().await;
remove_heal_test_format(&temp_dir, &store, 0, 3).await;
let (mut future_identity, _) =
read_config_no_lock_preserve_empty_with_metadata(store.pools[1].clone(), POOL_META_IDENTITY_NAME)
.await
.expect("current identity should be readable");
let future_version = u16::from_le_bytes([future_identity[2], future_identity[3]])
.checked_add(1)
.expect("identity version should have a future value");
future_identity[2..4].copy_from_slice(&future_version.to_le_bytes());
save_config(store.pools[1].clone(), POOL_META_IDENTITY_NAME, future_identity)
.await
.expect("future identity should be persisted");
let err = store
.handle_heal_format(false)
.await
.expect_err("a future identity must reject full format heal");
assert!(err.to_string().contains("pool metadata incompatible"));
assert_heal_test_format_missing(&temp_dir, 0, 3).await;
store
.ensure_pool_meta_side_effects_safe("future identity format heal side effect")
.await
.expect_err("the future identity rejection must latch the write gate");
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn replacement_format_heal_epoch_conflict_never_writes_format_and_latches() {
let (temp_dir, store, shutdown) = multi_pool_heal_store().await;
let target = remove_heal_test_format(&temp_dir, &store, 0, 3).await;
let conflicting_identity =
initialized_pool_meta_identity_for_test(store.id, 2).expect("conflicting identity should encode");
save_config(store.pools[1].clone(), POOL_META_IDENTITY_NAME, conflicting_identity)
.await
.expect("conflicting identity should be persisted");
let err = store
.heal_replacement_format(false, 0, 0, &[target])
.await
.expect_err("an identity epoch conflict must reject replacement format heal");
assert!(
err.to_string()
.contains("identity replicas disagree on cluster identity or epoch")
);
assert_heal_test_format_missing(&temp_dir, 0, 3).await;
store
.ensure_pool_meta_side_effects_safe("epoch conflict replacement format side effect")
.await
.expect_err("the identity epoch conflict must latch the write gate");
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn replacement_format_heal_initialized_identity_without_pool_meta_never_writes_and_latches() {
let (temp_dir, store, shutdown) = multi_pool_heal_store().await;
let target = remove_heal_test_format(&temp_dir, &store, 0, 3).await;
for pool in &store.pools {
delete_config(pool.clone(), POOL_META_NAME)
.await
.expect("pool metadata replica should be removable");
}
let err = store
.heal_replacement_format(false, 0, 0, &[target])
.await
.expect_err("initialized identity without pool metadata must reject replacement format heal");
assert!(err.to_string().contains("initialized cluster identity exists"));
assert_heal_test_format_missing(&temp_dir, 0, 3).await;
store
.ensure_pool_meta_side_effects_safe("missing pool metadata replacement format side effect")
.await
.expect_err("missing initialized pool metadata must latch the write gate");
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn full_format_heal_lost_after_last_save_does_not_publish_or_renew_and_latches() {
let (temp_dir, store, shutdown) = multi_pool_heal_store().await;
let target = remove_heal_test_format(&temp_dir, &store, 0, 3).await;
store.pools[0].disk_set[0].disks.write().await[3] = None;
let barrier = HealFormatAfterSaveBarrier::install(&store.pools[0], 3);
let recovery_latch = store.pool_meta_save_gate.lock().await.aborted_transaction_latch_for_test();
let mut heal = tokio::spawn({
let store = Arc::clone(&store);
async move { store.handle_heal_format(false).await }
});
barrier.wait_until_paused().await;
let saved = tokio::fs::read(heal_test_format_path(&temp_dir, 0, 3))
.await
.expect("the last replacement format must be durable before fence loss");
let saved = FormatV3::try_from(saved.as_slice()).expect("the durable replacement format should decode");
assert_eq!(saved.erasure.this, store.pools[0].format.erasure.sets[0][3]);
recovery_latch.store(true, std::sync::atomic::Ordering::SeqCst);
barrier.release();
let (result, err) = tokio::time::timeout(std::time::Duration::from_secs(30), &mut heal)
.await
.expect("format heal should stop after the lost fence")
.expect("format heal task should not panic")
.expect("format heal should return its fenced result");
assert!(matches!(err, Some(StorageError::SlowDown)));
assert!(
result
.after
.drives
.iter()
.any(|drive| drive.endpoint == target && drive.state == DriveState::Missing.to_string()),
"the lost fence must prevent the durable format from being published in the heal result"
);
assert!(
store.pools[0].disk_set[0].disks.read().await[3].is_none(),
"the lost fence must prevent renew_disk from attaching the replacement"
);
recovery_latch.store(false, std::sync::atomic::Ordering::SeqCst);
store
.ensure_pool_meta_side_effects_safe("post-save format fence loss side effect")
.await
.expect_err("post-save format fence loss must remain sticky");
shutdown.cancel();
}
#[tokio::test]
async fn heal_object_pool_scope_selects_only_requested_pool() {
let store = minimal_heal_store().await;
@@ -1573,13 +1816,18 @@ mod tests {
.handle_heal_format(false)
.await
.expect_err("missing pool metadata must fail closed before format writes");
assert!(matches!(err, StorageError::SlowDown));
assert!(err.to_string().contains("no durable bootstrap identity or pool.bin replica"));
store
.ensure_pool_meta_side_effects_safe("missing format-heal metadata side effect")
.await
.expect_err("missing metadata must latch the format-heal write gate");
let pool_meta = PoolMeta::new(&store.pools, &PoolMeta::default());
pool_meta
.save(store.pools.clone())
.await
.expect("pool metadata should be persisted before format heal");
*store.pool_meta_save_gate.lock().await = PoolMetaWriteState::default();
let (result, err) = store
.handle_heal_format(false)
+457 -43
View File
@@ -14,7 +14,8 @@
use super::*;
use crate::core::pools::{
PoolMetaReplicaState, PoolMetaWriteState, local_decommission_queue_prefix, pool_meta_has_active_decommission,
PoolMetaReplicaState, PoolMetaWriteState, load_pool_meta_identity_observing, local_decommission_queue_prefix,
persist_pool_meta_identity_for_startup, pool_meta_has_active_decommission,
};
use crate::error::is_err_decommission_running;
use crate::runtime::instance::InstanceContext;
@@ -137,47 +138,84 @@ async fn load_pool_meta_for_startup<S>(
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)
.await
.map_err(|err| Error::other(format!("store init failed during load_pool_meta: {err}")))?;
write_state.observe_replicas(replica_state);
write_state
.ensure_missing_metadata_can_initialize()
.map_err(|err| Error::other(format!("store init failed during classify_pool_meta_absence: {err}")))?;
Ok((meta, replica_state))
}
async fn save_validated_pool_meta_for_startup<S>(meta: &PoolMeta, pools: Vec<Arc<S>>) -> Result<()>
async fn establish_pool_meta_bootstrap_identity_if_proven<S>(
pools: Vec<Arc<S>>,
write_state: &mut PoolMetaWriteState,
elected_writer: bool,
) -> Result<()>
where
S: EcstoreObjectIO,
{
resolve_store_init_stage_result(meta.save_for_startup(pools).await, "save_validated_pool_meta")
if elected_writer && write_state.fresh_bootstrap_proven() {
persist_pool_meta_identity_for_startup(pools, write_state, false).await?;
}
Ok(())
}
async fn save_validated_pool_meta_for_startup<S>(
meta: &PoolMeta,
pools: Vec<Arc<S>>,
write_state: &mut PoolMetaWriteState,
) -> Result<PoolMeta>
where
S: EcstoreObjectIO,
{
meta.save_for_startup_observing(pools, write_state)
.await
.map_err(|err| Error::other(format!("store init failed during save_validated_pool_meta: {err}")))
}
async fn persist_pool_meta_for_startup_if_safe<S>(
meta: &PoolMeta,
pools: Vec<Arc<S>>,
replica_state: PoolMetaReplicaState,
write_state: PoolMetaWriteState,
write_state: &mut PoolMetaWriteState,
topology_update: bool,
elected_writer: bool,
) -> Result<()>
) -> Result<PoolMeta>
where
S: EcstoreObjectIO,
{
if !elected_writer {
return Ok(());
return Ok(meta.clone());
}
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 should_write {
if should_write || write_state.identity_requires_repair() {
write_state.ensure_write_safe("store init failed during save_validated_pool_meta")?;
}
let mut committed = meta.clone();
if should_write {
save_validated_pool_meta_for_startup(meta, pools).await?;
if write_state.fresh_bootstrap_proven() || write_state.identity_is_pending() {
persist_pool_meta_identity_for_startup(pools.clone(), write_state, false)
.await
.map_err(|err| Error::other(format!("store init failed during prepare_pool_meta_identity: {err}")))?;
}
committed = save_validated_pool_meta_for_startup(meta, pools.clone(), write_state).await?;
}
Ok(())
if should_write || write_state.identity_requires_repair() {
persist_pool_meta_identity_for_startup(pools, write_state, true)
.await
.map_err(|err| Error::other(format!("store init failed during commit_pool_meta_identity: {err}")))?;
}
Ok(committed)
}
async fn resume_local_decommission_after_init(store: Arc<ECStore>, rx: CancellationToken, pool_indices: Vec<usize>) {
@@ -286,6 +324,7 @@ impl ECStore {
preflight_startup_rpc_secret(&endpoint_pools)?;
let mut deployment_id = None;
let mut fresh_bootstrap_proven = true;
// let (endpoint_pools, _) = EndpointServerPools::create_server_endpoints(address.as_str(), &layouts)?;
@@ -345,7 +384,7 @@ impl ECStore {
check_disk_fatal_errs(&errs)?;
let fm = {
let loaded_format = {
let mut times = 0;
let mut interval = 1;
loop {
@@ -401,6 +440,8 @@ impl ECStore {
}
}
}?;
fresh_bootstrap_proven &= loaded_format.fresh_bootstrap_proven;
let fm = loaded_format.format;
// Format loading succeeded, enable health monitoring on all disks
for disk in disks.iter().flatten() {
@@ -436,13 +477,14 @@ impl ECStore {
runtime_sources::record_local_disks(&instance_ctx, local_disks).await;
}
let deployment_id = deployment_id.ok_or_else(|| Error::other("store init failed: deployment id is not initialized"))?;
let peer_sys = S3PeerSys::new_with_instance_ctx(&endpoint_pools, instance_ctx.clone());
let mut pool_meta = PoolMeta::new(&pools, &PoolMeta::default());
pool_meta.dont_save = true;
let decommission_cancelers = RwLock::new(vec![None; pools.len()]);
let ec = Arc::new(ECStore {
id: deployment_id.ok_or_else(|| Error::other("store init failed: deployment id is not initialized"))?,
id: deployment_id,
disk_map,
pools,
peer_sys,
@@ -450,7 +492,7 @@ impl ECStore {
rebalance_meta: RwLock::new(None),
decommission_cancelers,
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::default(),
pool_meta_save_gate: Mutex::new(PoolMetaWriteState::for_startup(deployment_id, fresh_bootstrap_proven)),
// 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.
@@ -459,10 +501,8 @@ impl ECStore {
});
// Only set it when this instance's deployment ID is not yet configured
if let Some(dep_id) = deployment_id
&& instance_ctx.deployment_id().is_none()
{
instance_ctx.set_deployment_id(dep_id);
if instance_ctx.deployment_id().is_none() {
instance_ctx.set_deployment_id(deployment_id);
}
let wait_sec = 5;
@@ -494,15 +534,21 @@ impl ECStore {
pub async fn init(self: &Arc<Self>, rx: CancellationToken) -> Result<()> {
runtime_sources::ensure_boot_time().await;
let should_persist_pool_meta = self
.pools
.first()
.is_some_and(|pool| pool_first_endpoint_is_local(&pool.endpoints));
let (meta, pool_meta_replica_state) = {
let mut write_state = self.pool_meta_save_gate.lock().await;
establish_pool_meta_bootstrap_identity_if_proven(self.pools.clone(), &mut write_state, should_persist_pool_meta)
.await
.map_err(|err| Error::other(format!("store init failed during establish_pool_meta_bootstrap_identity: {err}")))?;
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;
let installed_pool_meta = if update {
let mut installed_pool_meta = if update {
PoolMeta::new(&self.pools, &meta)
} else {
meta.clone()
@@ -510,12 +556,12 @@ 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.
{
let write_state = self.pool_meta_save_gate.lock().await;
persist_pool_meta_for_startup_if_safe(
let mut write_state = self.pool_meta_save_gate.lock().await;
installed_pool_meta = persist_pool_meta_for_startup_if_safe(
&installed_pool_meta,
self.pools.clone(),
pool_meta_replica_state,
*write_state,
&mut write_state,
update,
should_persist_pool_meta,
)
@@ -615,11 +661,11 @@ impl ECStore {
#[cfg(test)]
mod tests {
use super::{
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,
LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES, PoolMetaWriteState, establish_pool_meta_bootstrap_identity_if_proven,
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;
@@ -677,7 +723,10 @@ mod tests {
};
use crate::{
bucket::replication::{ReplicationState, ReplicationStatusType, replication_statuses_map},
core::pools::{POOL_META_VERSION, PoolDecommissionInfo, PoolMeta, PoolStatus},
core::pools::{
POOL_META_IDENTITY_NAME, POOL_META_NAME, POOL_META_VERSION, PoolDecommissionInfo, PoolMeta, PoolStatus,
pool_meta_identity_initialized_for_test, pool_meta_v3_commit_state_for_test,
},
disk::endpoint::Endpoint,
error::{Error, Result, StorageError},
io_support::rio::{WritePlan, compression_metadata_value},
@@ -718,6 +767,7 @@ mod tests {
use time::OffsetDateTime;
use tokio::io::AsyncReadExt;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
fn startup_pool_meta_payload(meta: &PoolMeta) -> Vec<u8> {
meta.encode_config_data_for_test().expect("pool metadata should encode")
@@ -731,10 +781,19 @@ mod tests {
wrote_without_lock: AtomicBool,
wrote_with_max_parity: AtomicBool,
written_payload: Mutex<Option<Vec<u8>>>,
revision: AtomicUsize,
pool_meta_write_attempts: AtomicUsize,
fail_pool_meta_write_at: AtomicUsize,
pool_meta_written_versions: Mutex<Vec<u16>>,
objects: Mutex<HashMap<String, (Vec<u8>, String)>>,
}
impl StartupPoolMetaStorage {
fn new(read_payload: Vec<u8>) -> Self {
let mut objects = HashMap::new();
if !read_payload.is_empty() {
objects.insert(POOL_META_NAME.to_string(), (read_payload.clone(), "startup-pool-meta-0".to_string()));
}
Self {
read_payload,
read_error: false,
@@ -742,6 +801,11 @@ mod tests {
wrote_without_lock: AtomicBool::new(false),
wrote_with_max_parity: AtomicBool::new(false),
written_payload: Mutex::new(None),
revision: AtomicUsize::new(0),
pool_meta_write_attempts: AtomicUsize::new(0),
fail_pool_meta_write_at: AtomicUsize::new(0),
pool_meta_written_versions: Mutex::new(Vec::new()),
objects: Mutex::new(objects),
}
}
@@ -753,15 +817,21 @@ mod tests {
wrote_without_lock: AtomicBool::new(false),
wrote_with_max_parity: AtomicBool::new(false),
written_payload: Mutex::new(None),
revision: AtomicUsize::new(0),
pool_meta_write_attempts: AtomicUsize::new(0),
fail_pool_meta_write_at: AtomicUsize::new(0),
pool_meta_written_versions: Mutex::new(Vec::new()),
objects: Mutex::new(HashMap::new()),
}
}
fn object_info(&self, bucket: &str, object: &str, size: usize) -> ObjectInfo {
fn object_info(&self, bucket: &str, object: &str, size: usize, etag: String) -> ObjectInfo {
ObjectInfo {
bucket: bucket.to_string(),
name: object.to_string(),
size: size as i64,
actual_size: size as i64,
etag: Some(etag),
..Default::default()
}
}
@@ -790,13 +860,19 @@ mod tests {
if self.read_error {
return Err(Error::other("pool metadata read quorum unavailable"));
}
if self.read_payload.is_empty() {
let Some((payload, etag)) = self
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(object)
.cloned()
else {
return Err(Error::FileNotFound);
}
};
Ok(GetObjectReader {
stream: Box::new(Cursor::new(self.read_payload.clone())),
object_info: self.object_info(bucket, object, self.read_payload.len()),
stream: Box::new(Cursor::new(payload.clone())),
object_info: self.object_info(bucket, object, payload.len(), etag),
buffered_body: None,
body_source: Default::default(),
})
@@ -812,11 +888,53 @@ mod tests {
assert!(opts.no_lock, "store init pool metadata save must not require namespace locks");
self.wrote_without_lock.store(true, Ordering::SeqCst);
self.wrote_with_max_parity.store(opts.max_parity, Ordering::SeqCst);
if object == POOL_META_NAME {
let attempt = self.pool_meta_write_attempts.fetch_add(1, Ordering::SeqCst) + 1;
if self.fail_pool_meta_write_at.load(Ordering::SeqCst) == attempt {
return Err(Error::Timeout);
}
}
let current_etag = self
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(object)
.map(|(_, etag)| etag.clone());
if opts
.http_preconditions
.as_ref()
.and_then(|preconditions| preconditions.if_none_match_value())
== Some("*")
&& current_etag.is_some()
{
return Err(Error::PreconditionFailed);
}
if let Some(expected) = opts
.http_preconditions
.as_ref()
.and_then(|preconditions| preconditions.if_match_value())
&& current_etag.as_deref() != Some(expected)
{
return Err(Error::PreconditionFailed);
}
let mut payload = Vec::new();
data.stream.read_to_end(&mut payload).await?;
let size = payload.len();
*self.written_payload.lock().unwrap_or_else(std::sync::PoisonError::into_inner) = Some(payload);
Ok(self.object_info(bucket, object, size))
if object == POOL_META_NAME {
*self.written_payload.lock().unwrap_or_else(std::sync::PoisonError::into_inner) = Some(payload.clone());
if payload.len() >= 4 {
self.pool_meta_written_versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(u16::from_le_bytes([payload[2], payload[3]]));
}
}
let etag = format!("startup-pool-meta-{}", self.revision.fetch_add(1, Ordering::SeqCst) + 1);
self.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(object.to_string(), (payload, etag.clone()));
Ok(self.object_info(bucket, object, size, etag))
}
}
@@ -834,9 +952,12 @@ mod tests {
}
#[tokio::test]
async fn test_store_init_pool_meta_io_bypasses_namespace_lock_surface() {
async fn test_fresh_bootstrap_pool_meta_save_reaches_v3_cas_without_blocking_itself() {
let storage = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let mut write_state = PoolMetaWriteState::default();
let mut write_state = PoolMetaWriteState::for_startup(Uuid::new_v4(), true);
establish_pool_meta_bootstrap_identity_if_proven(vec![storage.clone()], &mut write_state, true)
.await
.expect("the elected fresh bootstrap should persist its identity before loading pool metadata");
let (loaded, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut write_state)
.await
@@ -851,11 +972,303 @@ mod tests {
pools: Vec::new(),
dont_save: false,
};
save_validated_pool_meta_for_startup(&meta, vec![storage.clone()])
save_validated_pool_meta_for_startup(&meta, vec![storage.clone()], &mut write_state)
.await
.expect("startup pool metadata save should bypass locks");
assert!(storage.wrote_without_lock.load(Ordering::SeqCst));
assert!(storage.wrote_with_max_parity.load(Ordering::SeqCst));
assert_eq!(
storage.pool_meta_write_attempts.load(Ordering::SeqCst),
2,
"fresh bootstrap should reach both V3 prepare and commit CAS writes"
);
assert_eq!(
*storage
.pool_meta_written_versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
vec![3, 3]
);
write_state
.ensure_write_safe("fresh bootstrap publication")
.expect("successful startup publication must disarm the transaction guard");
}
#[tokio::test]
async fn test_unproven_pending_identity_cannot_authorize_all_missing_pool_meta() {
let deployment_id = Uuid::new_v4();
let storage = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let mut proven_bootstrap = PoolMetaWriteState::for_startup(deployment_id, true);
establish_pool_meta_bootstrap_identity_if_proven(vec![storage.clone()], &mut proven_bootstrap, true)
.await
.expect("fresh topology proof should persist a nonce-bound pending identity");
let identity = storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(POOL_META_IDENTITY_NAME)
.map(|(payload, _)| payload.clone())
.expect("pending identity should be durable");
assert!(!pool_meta_identity_initialized_for_test(&identity).expect("decode pending identity"));
let mut unproven_restart = PoolMetaWriteState::for_startup(deployment_id, false);
let err = load_pool_meta_for_startup(vec![storage.clone()], &mut unproven_restart)
.await
.expect_err("a pending identity alone must not authorize an all-missing restart");
assert!(err.to_string().contains("no verified fresh-bootstrap proof"));
unproven_restart
.ensure_write_safe("unproven pending identity restart")
.expect_err("the rejected restart must latch the write gate");
assert!(
!storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains_key(POOL_META_NAME),
"classification must not create pool metadata"
);
}
#[tokio::test]
async fn test_nonfresh_identity_repair_crash_never_persists_pending_bootstrap_authority() {
let deployment_id = Uuid::new_v4();
let initial = init_test_pool_meta(None);
let storage = Arc::new(StartupPoolMetaStorage::new(startup_pool_meta_payload(&initial)));
let mut write_state = PoolMetaWriteState::for_startup(deployment_id, false);
let (loaded, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut write_state)
.await
.expect("an existing pool metadata snapshot may repair its missing identity");
storage.fail_pool_meta_write_at.store(1, Ordering::SeqCst);
persist_pool_meta_for_startup_if_safe(&loaded, vec![storage.clone()], replica_state, &mut write_state, true, true)
.await
.expect_err("the injected crash boundary should stop before the pool metadata replacement");
let identity = storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(POOL_META_IDENTITY_NAME)
.map(|(payload, _)| payload.clone())
.expect("identity repair should be durable before the failed topology write");
assert!(
pool_meta_identity_initialized_for_test(&identity).expect("decode repaired identity"),
"an existing cluster must never leave pending bootstrap authority at this crash boundary"
);
storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(POOL_META_NAME);
let mut restarted = PoolMetaWriteState::for_startup(deployment_id, false);
let err = load_pool_meta_for_startup(vec![storage], &mut restarted)
.await
.expect_err("wiping pool metadata after the crash must require recovery");
assert!(err.to_string().contains("initialized cluster identity exists"));
}
#[tokio::test]
#[serial_test::serial(pool_meta_version_env)]
async fn test_existing_legacy_identity_and_replica_repair_respects_disabled_v3_gates() {
temp_env::async_with_vars(
[
(rustfs_config::ENV_POOL_META_V3_WRITE, None::<&str>),
(rustfs_config::ENV_POOL_META_V3_FLEET_CONFIRMED, None::<&str>),
],
async {
let deployment_id = Uuid::new_v4();
let corrupt = Arc::new(StartupPoolMetaStorage::new(vec![0, 1, 2]));
let legacy = init_test_pool_meta(None);
let valid = Arc::new(StartupPoolMetaStorage::new(startup_pool_meta_payload(&legacy)));
let mut write_state = PoolMetaWriteState::for_startup(deployment_id, false);
let (loaded, replica_state) = load_pool_meta_for_startup(vec![corrupt.clone(), valid.clone()], &mut write_state)
.await
.expect("existing V2 metadata should remain readable while its identity is missing");
assert!(replica_state.needs_repair);
let committed = persist_pool_meta_for_startup_if_safe(
&loaded,
vec![corrupt.clone(), valid.clone()],
replica_state,
&mut write_state,
true,
true,
)
.await
.expect("legacy replica and identity repair should succeed without crossing the V3 gate");
assert_eq!(committed.version, POOL_META_VERSION);
for storage in [corrupt, valid] {
assert_eq!(
*storage
.pool_meta_written_versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
vec![POOL_META_VERSION],
"legacy repair must write exactly one V2 snapshot"
);
let identity = storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(POOL_META_IDENTITY_NAME)
.map(|(payload, _)| payload.clone())
.expect("identity repair should persist on every pool");
assert!(pool_meta_identity_initialized_for_test(&identity).expect("decode repaired identity"));
}
},
)
.await;
}
#[tokio::test]
async fn test_store_init_distinguishes_fresh_deployment_from_wiped_lagging_node() {
let deployment_id = Uuid::new_v4();
let storage = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let mut fresh_state = PoolMetaWriteState::for_startup(deployment_id, true);
establish_pool_meta_bootstrap_identity_if_proven(vec![storage.clone()], &mut fresh_state, true)
.await
.expect("fresh topology proof should become a durable pending identity");
let (_, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut fresh_state)
.await
.expect("new formats with no identity may initialize pool metadata exactly once");
let committed = persist_pool_meta_for_startup_if_safe(
&init_test_pool_meta(None),
vec![storage.clone()],
replica_state,
&mut fresh_state,
true,
true,
)
.await
.expect("fresh deployment should durably commit identity and pool metadata");
assert_eq!(committed.pools[0].cmd_line, "pool-0");
{
let objects = storage.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
assert!(objects.contains_key(POOL_META_NAME));
assert!(objects.contains_key(POOL_META_IDENTITY_NAME));
}
storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(POOL_META_NAME);
let mut wiped_pool_state = PoolMetaWriteState::for_startup(deployment_id, false);
let err = load_pool_meta_for_startup(vec![storage.clone()], &mut wiped_pool_state)
.await
.expect_err("an initialized identity must prevent an empty pool metadata rebuild");
assert!(err.to_string().contains("initialized cluster identity exists"));
storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(POOL_META_IDENTITY_NAME);
let mut wiped_all_state = PoolMetaWriteState::for_startup(deployment_id, false);
let err = load_pool_meta_for_startup(vec![storage], &mut wiped_all_state)
.await
.expect_err("existing formats without identity or pool metadata must require recovery");
assert!(err.to_string().contains("no durable bootstrap identity"));
let distributed_node = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let mut non_elected_state = PoolMetaWriteState::for_startup(Uuid::new_v4(), true);
establish_pool_meta_bootstrap_identity_if_proven(vec![distributed_node.clone()], &mut non_elected_state, false)
.await
.expect("a non-elected distributed node must not create bootstrap authority");
let err = load_pool_meta_for_startup(vec![distributed_node], &mut non_elected_state)
.await
.expect_err("a distributed node without durable identity or pool metadata must require recovery");
assert!(err.to_string().contains("no durable bootstrap identity"));
}
#[tokio::test]
async fn test_store_init_resumes_pending_v3_bootstrap_without_legacy_overwrite() {
let deployment_id = Uuid::new_v4();
let storage = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let mut first_bootstrap = PoolMetaWriteState::for_startup(deployment_id, true);
establish_pool_meta_bootstrap_identity_if_proven(vec![storage.clone()], &mut first_bootstrap, true)
.await
.expect("fresh bootstrap should persist a pending identity");
let (_, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut first_bootstrap)
.await
.expect("the durable pending identity should authorize the initial pool metadata write");
storage.fail_pool_meta_write_at.store(2, Ordering::SeqCst);
persist_pool_meta_for_startup_if_safe(
&init_test_pool_meta(None),
vec![storage.clone()],
replica_state,
&mut first_bootstrap,
true,
true,
)
.await
.expect_err("the injected crash boundary should leave only the V3 prepare record");
let (pending_pool, pending_identity) = {
let objects = storage.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
(
objects
.get(POOL_META_NAME)
.map(|(payload, _)| payload.clone())
.expect("the V3 prepare record should be durable"),
objects
.get(POOL_META_IDENTITY_NAME)
.map(|(payload, _)| payload.clone())
.expect("the bootstrap identity should be durable"),
)
};
assert_eq!(pool_meta_v3_commit_state_for_test(pending_pool).expect("decode pending V3"), (1, false));
assert!(!pool_meta_identity_initialized_for_test(&pending_identity).expect("decode pending identity"));
storage.fail_pool_meta_write_at.store(0, Ordering::SeqCst);
let mut restarted = PoolMetaWriteState::for_startup(deployment_id, false);
let (loaded, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut restarted)
.await
.expect("restart should recover the predecessor embedded in the pending V3 record");
assert!(loaded.pools.is_empty());
assert!(replica_state.needs_repair);
let committed = persist_pool_meta_for_startup_if_safe(
&init_test_pool_meta(None),
vec![storage.clone()],
replica_state,
&mut restarted,
true,
true,
)
.await
.expect("restart should finish generation 1 before promoting the identity");
assert_eq!(committed.version, 3);
let (committed_pool, committed_identity) = {
let objects = storage.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
(
objects
.get(POOL_META_NAME)
.map(|(payload, _)| payload.clone())
.expect("the committed V3 record should be durable"),
objects
.get(POOL_META_IDENTITY_NAME)
.map(|(payload, _)| payload.clone())
.expect("the committed identity should be durable"),
)
};
assert_eq!(
pool_meta_v3_commit_state_for_test(committed_pool).expect("decode committed V3"),
(1, true)
);
assert!(pool_meta_identity_initialized_for_test(&committed_identity).expect("decode committed identity"));
assert!(
storage
.pool_meta_written_versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.iter()
.all(|version| *version == 3),
"bootstrap recovery must never overwrite a pending V3 record with legacy metadata"
);
}
#[tokio::test]
@@ -880,7 +1293,7 @@ mod tests {
&loaded,
vec![corrupt.clone(), backup.clone()],
replica_state,
write_state,
&mut write_state,
false,
true,
)
@@ -919,7 +1332,7 @@ mod tests {
&loaded,
vec![valid.clone(), unreadable.clone()],
replica_state,
write_state,
&mut write_state,
false,
true,
)
@@ -932,7 +1345,7 @@ mod tests {
&loaded,
vec![valid.clone(), unreadable.clone()],
replica_state,
write_state,
&mut write_state,
true,
true,
)
@@ -968,9 +1381,10 @@ mod tests {
.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");
let err =
persist_pool_meta_for_startup_if_safe(&loaded, vec![repaired.clone()], replica_state, &mut 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")
+44 -17
View File
@@ -77,7 +77,14 @@ pub async fn connect_load_init_formats(
deployment_id: Option<Uuid>,
) -> Result<FormatV3> {
let instance_ctx = crate::runtime::global::current_ctx();
connect_load_init_formats_with_instance_ctx(&instance_ctx, first_disk, disks, set_count, set_drive_count, deployment_id).await
connect_load_init_formats_with_instance_ctx(&instance_ctx, first_disk, disks, set_count, set_drive_count, deployment_id)
.await
.map(|loaded| loaded.format)
}
pub(crate) struct LoadedFormat {
pub(crate) format: FormatV3,
pub(crate) fresh_bootstrap_proven: bool,
}
pub(crate) async fn connect_load_init_formats_with_instance_ctx(
@@ -87,17 +94,18 @@ pub(crate) async fn connect_load_init_formats_with_instance_ctx(
set_count: usize,
set_drive_count: usize,
deployment_id: Option<Uuid>,
) -> Result<FormatV3> {
) -> Result<LoadedFormat> {
let (formats, errs) = load_format_erasure_all(disks, false).await;
check_disk_fatal_errs(&errs)?;
// Treat transient network errors (connection refused, timeout, etc.) as
// equivalent to UnformattedDisk for the bootstrap decision. During
// fresh-cluster startup a remote peer that cannot be reached is
// indistinguishable from an unformatted disk — the peer may simply not
// have started its gRPC server yet.
// Transient network errors still follow the first-node wait path so peers
// can come online, but they are never fresh-cluster authority below.
let all_unformatted = errs.iter().all(is_unformatted_or_transient_network);
// Fresh-cluster authority requires a response from every configured disk.
// A transiently unreachable peer may belong to an existing cluster and
// must never be treated as proof that the topology is new.
let fresh_bootstrap_proven = should_init_erasure_disks(&errs);
let formats_present = formats.iter().flatten().count();
let mut format_quorum = (formats_present > 0).then(|| select_format_erasure_in_quorum(&formats, 0));
if format_quorum.as_ref().is_none_or(Result::is_err)
@@ -123,7 +131,10 @@ pub(crate) async fn connect_load_init_formats_with_instance_ctx(
Ok(LegacyFormatOutcome::Migrated { format, quorum_members }) => {
info!("Migrated format from MinIO config");
retain_format_quorum_members(instance_ctx, disks, &format, &quorum_members, set_drive_count).await?;
return Ok(*format);
return Ok(LoadedFormat {
format: *format,
fresh_bootstrap_proven: false,
});
}
Ok(LegacyFormatOutcome::Incompatible) => {
error!(
@@ -138,9 +149,12 @@ pub(crate) async fn connect_load_init_formats_with_instance_ctx(
Ok(LegacyFormatOutcome::None) => {}
Err(e) => return Err(e),
}
if all_unformatted {
if fresh_bootstrap_proven {
let fm = init_format_erasure(instance_ctx, disks, set_count, set_drive_count, deployment_id).await?;
return Ok(fm);
return Ok(LoadedFormat {
format: fm,
fresh_bootstrap_proven: true,
});
}
}
@@ -166,7 +180,10 @@ pub(crate) async fn connect_load_init_formats_with_instance_ctx(
check_format_erasure_value_for_topology(&fm, formats.len(), set_drive_count)?;
retain_format_quorum_members(instance_ctx, disks, &fm, &quorum_members, set_drive_count).await?;
Ok(fm)
Ok(LoadedFormat {
format: fm,
fresh_bootstrap_proven: false,
})
}
async fn retain_format_quorum_members(
@@ -258,11 +275,8 @@ pub fn should_init_erasure_disks(errs: &[Option<DiskError>]) -> bool {
count_errs(errs, &DiskError::UnformattedDisk) == errs.len()
}
/// Returns `true` if the error represents a disk that is either unformatted
/// or unreachable due to a transient network failure. During fresh-cluster
/// bootstrap a remote peer that cannot be reached is indistinguishable from
/// an unformatted disk — the peer may simply not have started its gRPC
/// server yet.
/// Returns `true` for errors that stay on the first-node wait path. This is not
/// fresh-cluster proof; only [`should_init_erasure_disks`] grants that.
fn is_unformatted_or_transient_network(err: &Option<DiskError>) -> bool {
matches!(err, Some(DiskError::UnformattedDisk)) || err.as_ref().is_some_and(is_network_like_disk_error)
}
@@ -1294,9 +1308,14 @@ mod tests {
async fn fresh_format_load_initializes_all_disks() {
let (_temp_dir, mut disks) = local_disks(3).await;
let format = connect_load_init_formats(true, &mut disks, 1, 3, None)
let loaded = connect_load_init_formats_with_instance_ctx(&current_ctx(), true, &mut disks, 1, 3, None)
.await
.expect("fresh disks should receive a storage format");
assert!(
loaded.fresh_bootstrap_proven,
"every configured disk explicitly reporting unformatted should establish fresh topology proof"
);
let format = loaded.format;
let (formats, errors) = load_format_erasure_all(&disks, false).await;
assert!(errors.iter().all(Option::is_none), "every disk should load its fresh format: {errors:?}");
@@ -1640,6 +1659,14 @@ mod tests {
assert!(!is_unformatted_or_transient_network(&Some(DiskError::FileNotFound)));
assert!(!is_unformatted_or_transient_network(&Some(DiskError::CorruptedFormat)));
assert!(!is_unformatted_or_transient_network(&Some(DiskError::DiskFull)));
assert!(should_init_erasure_disks(&[
Some(DiskError::UnformattedDisk),
Some(DiskError::UnformattedDisk),
]));
assert!(
!should_init_erasure_disks(&[Some(DiskError::UnformattedDisk), Some(DiskError::Timeout)]),
"an unreachable peer is not fresh-topology proof"
);
}
}
+198 -5
View File
@@ -931,8 +931,12 @@ fn lifecycle_delete_all_test_failure(phase: crate::object_api::LifecycleDeleteAl
mod tests {
use super::*;
use crate::bucket::replication::{ReplicationStatusType, VersionPurgeStatusType};
use crate::config::com::{read_config_no_lock_preserve_empty_with_metadata, save_config};
use crate::config::storageclass::{CLASS_RRS, CLASS_STANDARD, lookup_config_for_pools_without_env};
use crate::core::pools::{POOL_META_NAME, POOL_META_VERSION, PoolDecommissionInfo, PoolStatus};
use crate::core::pools::{
POOL_META_IDENTITY_NAME, POOL_META_NAME, POOL_META_VERSION, PoolDecommissionInfo, PoolStatus,
initialized_pool_meta_identity_for_test, pool_meta_v3_commit_state_for_test,
};
use crate::disk::error::DiskError;
use crate::layout::endpoint::Endpoint;
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
@@ -2375,6 +2379,135 @@ mod tests {
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn runtime_save_current_pool_meta_reaches_v3_cas_and_disarms_transaction() {
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-runtime-cas", &[2]).await;
let bootstrapped = store
.load_runtime_pool_meta("verify fresh bootstrap V3 CAS save")
.await
.expect("fresh startup should durably publish pool metadata");
assert_eq!(bootstrapped.version, 3);
let (bootstrap_payload, _) = read_config_no_lock_preserve_empty_with_metadata(store.pools[0].clone(), POOL_META_NAME)
.await
.expect("fresh bootstrap V3 replica should be readable");
assert_eq!(
pool_meta_v3_commit_state_for_test(bootstrap_payload).expect("fresh bootstrap V3 envelope should decode"),
(1, true)
);
let saved_at = OffsetDateTime::now_utc() + TimeDuration::seconds(30);
{
let mut pool_meta = store.pool_meta.write().await;
pool_meta.pools[0].last_update = saved_at;
}
store
.save_current_pool_meta_for_test(&[0])
.await
.expect("runtime pool metadata save should reach V3 conditional writes");
store
.ensure_pool_meta_side_effects_safe("runtime save publication")
.await
.expect("successful runtime publication must disarm the transaction guard");
let persisted = store
.load_runtime_pool_meta("verify runtime V3 CAS save")
.await
.expect("runtime load should observe the committed save");
assert_eq!(persisted.version, 3);
assert_eq!(persisted.pools[0].last_update, saved_at);
let (runtime_payload, _) = read_config_no_lock_preserve_empty_with_metadata(store.pools[0].clone(), POOL_META_NAME)
.await
.expect("runtime V3 replica should be readable");
assert_eq!(
pool_meta_v3_commit_state_for_test(runtime_payload).expect("runtime V3 envelope should decode"),
(2, true)
);
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn stale_clear_decommission_cannot_erase_active_or_queued_replacement() {
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-stale-clear", &[2, 2]).await;
let replacement_time = OffsetDateTime::UNIX_EPOCH + TimeDuration::seconds(100);
let cases = [
(
"failed-active",
PoolDecommissionInfo {
failed: true,
..Default::default()
},
PoolDecommissionInfo {
start_time: Some(replacement_time),
..Default::default()
},
),
(
"canceled-queued",
PoolDecommissionInfo {
canceled: true,
..Default::default()
},
PoolDecommissionInfo {
queued: true,
..Default::default()
},
),
];
for (case, stale_terminal, replacement_info) in cases {
let mut replacement = store
.load_runtime_pool_meta("prepare stale clear replacement")
.await
.expect("the current durable pool metadata should load");
replacement.pools[0].last_update = replacement_time;
replacement.pools[0].decommission = Some(replacement_info.clone());
persist_reload_snapshot(&store, &replacement).await;
let mut stale_local = replacement.clone();
stale_local.pools[0].last_update = OffsetDateTime::UNIX_EPOCH;
stale_local.pools[0].decommission = Some(stale_terminal.clone());
*store.pool_meta.write().await = stale_local;
let err = store
.clear_decommission(0)
.await
.expect_err("a stale terminal clear must not erase a durable replacement");
assert!(
err.to_string()
.contains("persisted active or queued decommission cannot be cleared"),
"unexpected {case} rejection: {err}"
);
let durable = store
.load_runtime_pool_meta("verify stale clear replacement")
.await
.expect("the durable replacement should remain readable");
let durable_info = durable.pools[0]
.decommission
.as_ref()
.expect("the durable replacement must not be cleared");
assert_eq!(durable.pools[0].last_update, replacement_time, "{case} replacement revision changed");
assert_eq!(
durable_info.start_time, replacement_info.start_time,
"{case} replacement generation changed"
);
assert_eq!(durable_info.queued, replacement_info.queued, "{case} replacement queue state changed");
let local = store.pool_meta.read().await;
let local_info = local.pools[0]
.decommission
.as_ref()
.expect("the failed clear must roll back the stale local terminal state");
assert_eq!(local_info.failed, stale_terminal.failed, "{case} failed state was not rolled back");
assert_eq!(local_info.canceled, stale_terminal.canceled, "{case} canceled state was not rolled back");
}
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn peer_pool_meta_reload_does_not_rollback_newer_local_states() {
@@ -2565,7 +2698,7 @@ mod tests {
#[tokio::test]
#[serial_test::serial]
async fn peer_pool_meta_reload_fails_closed_when_persisted_metadata_is_missing() {
async fn peer_pool_meta_reload_latches_recovery_when_initialized_metadata_is_missing() {
let (temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-reload-missing", &[2]).await;
let kept_time = OffsetDateTime::now_utc();
@@ -2612,11 +2745,15 @@ mod tests {
panic!("no pool.bin found under {:?}", temp_dir.path());
}
let merged_newer = store
let err = store
.reload_pool_meta()
.await
.expect("reload with missing metadata should fail closed, not error");
assert!(!merged_newer, "missing persisted metadata must not count as merged state");
.expect_err("an initialized cluster with every pool.bin missing must require recovery");
assert!(err.to_string().contains("initialized cluster identity exists"));
store
.ensure_pool_meta_side_effects_safe("test reload side effect")
.await
.expect_err("the recovery-required reload must latch the runtime side-effect gate");
let pool_meta = store.pool_meta.read().await;
let info = pool_meta.pools[0]
@@ -2628,4 +2765,60 @@ mod tests {
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn peer_pool_meta_reload_latches_recovery_after_identity_epoch_conflict() {
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-reload-identity-conflict", &[2, 2]).await;
let conflicting_identity =
initialized_pool_meta_identity_for_test(store.id, 2).expect("conflicting initialized identity should encode");
save_config(store.pools[1].clone(), POOL_META_IDENTITY_NAME, conflicting_identity)
.await
.expect("second pool identity should be replaced for the conflict scenario");
let err = store
.reload_pool_meta()
.await
.expect_err("divergent identity epochs must reject runtime reload");
assert!(
err.to_string()
.contains("identity replicas disagree on cluster identity or epoch")
);
store
.ensure_pool_meta_side_effects_safe("identity epoch conflict side effect")
.await
.expect_err("identity epoch conflict must latch the runtime side-effect gate");
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn peer_pool_meta_reload_latches_recovery_after_future_identity_version() {
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-reload-future-identity", &[2, 2]).await;
let (mut future_identity, _) =
read_config_no_lock_preserve_empty_with_metadata(store.pools[1].clone(), POOL_META_IDENTITY_NAME)
.await
.expect("current identity should be readable");
let current_version = u16::from_le_bytes([future_identity[2], future_identity[3]]);
let future_version = current_version
.checked_add(1)
.expect("identity version should have a future value");
future_identity[2..4].copy_from_slice(&future_version.to_le_bytes());
save_config(store.pools[1].clone(), POOL_META_IDENTITY_NAME, future_identity)
.await
.expect("second pool identity should be replaced with a future version");
let err = store
.reload_pool_meta()
.await
.expect_err("a future identity version must reject runtime reload");
assert!(err.to_string().contains("pool metadata incompatible"));
store
.ensure_pool_meta_side_effects_safe("future identity version side effect")
.await
.expect_err("future identity version must latch the runtime side-effect gate");
shutdown.cancel();
}
}