fix(ecstore): fence pool metadata replica updates (#6466)

* fix(ecstore): fence pool metadata replica updates

* fix(ecstore): block decommission on unsafe pool metadata

* fix(ecstore): block writes after pool metadata save errors

* fix(ecstore): latch pool metadata writes before await
This commit is contained in:
Zhengchao An
2026-08-24 14:25:19 +08:00
committed by GitHub
parent eec0e0e056
commit e091a7e702
11 changed files with 2012 additions and 323 deletions
+150 -11
View File
@@ -18,7 +18,7 @@ use crate::cluster::rpc::client::{
node_service_time_out_client,
};
use crate::cluster::rpc::set_tonic_mutation_body_digest;
use crate::core::pools::PoolMeta;
use crate::core::pools::{PoolMeta, PoolMetaWriteState};
use crate::disk::error::DiskError;
use crate::disk::error::{Error, Result};
use crate::disk::error_reduce::{BUCKET_OP_IGNORED_ERRS, is_all_buckets_not_found, reduce_write_quorum_errs};
@@ -324,6 +324,14 @@ pub trait PeerS3Client: Debug + Sync + Send + 'static {
async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, _fenced_pools: &[usize]) -> Result<HealResultItem> {
self.heal_bucket(bucket, opts).await
}
async fn heal_bucket_with_fence_from_movement_guarded_coordinator(
&self,
bucket: &str,
opts: &HealOpts,
fenced_pools: &[usize],
) -> Result<HealResultItem> {
self.heal_bucket_with_fence(bucket, opts, fenced_pools).await
}
async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()>;
async fn list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>>;
async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()>;
@@ -381,6 +389,25 @@ impl S3PeerSys {
}
pub async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, fenced_pools: &[usize]) -> Result<HealResultItem> {
self.heal_bucket_with_fence_inner(bucket, opts, fenced_pools, false).await
}
pub async fn heal_bucket_with_fence_from_movement_guarded_coordinator(
&self,
bucket: &str,
opts: &HealOpts,
fenced_pools: &[usize],
) -> Result<HealResultItem> {
self.heal_bucket_with_fence_inner(bucket, opts, fenced_pools, true).await
}
async fn heal_bucket_with_fence_inner(
&self,
bucket: &str,
opts: &HealOpts,
fenced_pools: &[usize],
movement_guard_held: bool,
) -> Result<HealResultItem> {
let mut opts = *opts;
let mut futures = Vec::with_capacity(self.clients.len());
for client in self.clients.iter() {
@@ -403,7 +430,14 @@ impl S3PeerSys {
let opts_clone = opts;
let heal_bucket_results_clone = heal_bucket_results.clone();
futures.push(async move {
match client.heal_bucket_with_fence(bucket, &opts_clone, fenced_pools).await {
let result = if movement_guard_held {
client
.heal_bucket_with_fence_from_movement_guarded_coordinator(bucket, &opts_clone, fenced_pools)
.await
} else {
client.heal_bucket_with_fence(bucket, &opts_clone, fenced_pools).await
};
match result {
Ok(res) => {
heal_bucket_results_clone.write().await[idx] = res;
None
@@ -698,6 +732,63 @@ impl LocalPeerS3Client {
.filter(|disk| usize::try_from(disk.endpoint().pool_idx).is_ok_and(|pool_idx| pools.contains(&pool_idx)))
.collect()
}
async fn heal_bucket_with_fence_inner(
&self,
bucket: &str,
opts: &HealOpts,
fenced_pools: &[usize],
movement_guard_held: bool,
) -> Result<HealResultItem> {
let disks = self.local_disks_for_pools().await.into_iter().map(Some).collect();
let store = runtime_sources::object_store_handle().filter(|store| Arc::ptr_eq(&store.ctx, &self.instance_ctx));
#[cfg(not(test))]
if store.is_none() {
return Err(Error::other("bucket heal refused: pool metadata is unavailable for this instance"));
}
let movement_gate = store.as_ref().map(|store| store.ctx.data_movement_operation_gate());
let movement_guard = try_acquire_bucket_heal_movement_guard(movement_gate.as_ref(), movement_guard_held)?;
let save_guard = acquire_bucket_heal_write_guard(store.as_ref().map(|store| &store.pool_meta_save_gate)).await?;
let result = heal_bucket_local_on_disks_with_pool_meta(
bucket,
opts,
disks,
store.as_ref().map(|store| &store.pool_meta),
fenced_pools,
)
.await;
drop(save_guard);
drop(movement_guard);
result
}
}
fn try_acquire_bucket_heal_movement_guard<'a>(
gate: Option<&'a Arc<tokio::sync::RwLock<()>>>,
movement_guard_held: bool,
) -> Result<Option<tokio::sync::RwLockReadGuard<'a, ()>>> {
if movement_guard_held {
return Ok(None);
}
let Some(gate) = gate else {
return Ok(None);
};
// Do not queue a receiver behind a movement writer while its coordinator
// holds another node's read guard; failing fast breaks that cross-node cycle.
gate.try_read()
.map(Some)
.map_err(|_| crate::error::StorageError::SlowDown.into())
}
async fn acquire_bucket_heal_write_guard<'a>(
gate: Option<&'a tokio::sync::Mutex<PoolMetaWriteState>>,
) -> Result<Option<tokio::sync::MutexGuard<'a, PoolMetaWriteState>>> {
let Some(gate) = gate else {
return Ok(None);
};
let guard = gate.lock().await;
guard.ensure_write_safe("bucket heal cannot run while pool metadata requires recovery")?;
Ok(Some(guard))
}
#[async_trait]
@@ -711,14 +802,16 @@ impl PeerS3Client for LocalPeerS3Client {
}
async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, fenced_pools: &[usize]) -> Result<HealResultItem> {
let disks = self.local_disks_for_pools().await.into_iter().map(Some).collect();
let store = runtime_sources::object_store_handle().filter(|store| Arc::ptr_eq(&store.ctx, &self.instance_ctx));
#[cfg(not(test))]
if store.is_none() {
return Err(Error::other("bucket heal refused: pool metadata is unavailable for this instance"));
}
heal_bucket_local_on_disks_with_pool_meta(bucket, opts, disks, store.as_ref().map(|store| &store.pool_meta), fenced_pools)
.await
self.heal_bucket_with_fence_inner(bucket, opts, fenced_pools, false).await
}
async fn heal_bucket_with_fence_from_movement_guarded_coordinator(
&self,
bucket: &str,
opts: &HealOpts,
fenced_pools: &[usize],
) -> Result<HealResultItem> {
self.heal_bucket_with_fence_inner(bucket, opts, fenced_pools, true).await
}
async fn list_bucket(&self, _opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
@@ -1656,7 +1749,7 @@ async fn clone_drives() -> Vec<Option<DiskStore>> {
#[cfg(test)]
mod tests {
use super::*;
use crate::core::pools::{PoolDecommissionInfo, PoolStatus};
use crate::core::pools::{PoolDecommissionInfo, PoolMetaReplicaState, PoolStatus};
use crate::disk::WalkDirOptions;
use crate::disk::disk_store::LocalDiskWrapper;
use crate::disk::endpoint::Endpoint;
@@ -2241,6 +2334,52 @@ mod tests {
reset_local_disk_test_state().await;
}
#[tokio::test]
async fn local_bucket_heal_refuses_receiver_with_unsafe_pool_metadata() {
let gate = tokio::sync::Mutex::new(PoolMetaWriteState::default());
gate.lock().await.observe_replicas(PoolMetaReplicaState {
needs_repair: true,
repair_write_safe: false,
});
let err = acquire_bucket_heal_write_guard(Some(&gate))
.await
.expect_err("receiver-side bucket heal must honor the local pool metadata gate");
assert!(
err.to_string()
.contains("bucket heal cannot run while pool metadata requires recovery")
);
}
#[tokio::test]
async fn receiver_bucket_heal_fails_fast_behind_queued_movement_writer() {
let gate = Arc::new(tokio::sync::RwLock::new(()));
let coordinator_guard = gate.read().await;
let writer_gate = gate.clone();
let mut writer = tokio::spawn(async move {
let _writer_guard = writer_gate.write().await;
});
while gate.try_read().is_ok() {
tokio::task::yield_now().await;
}
let err = try_acquire_bucket_heal_movement_guard(Some(&gate), false)
.expect_err("receiver must not wait behind a queued movement writer");
assert_eq!(err, crate::error::StorageError::SlowDown.into());
assert!(
try_acquire_bucket_heal_movement_guard(Some(&gate), true)
.expect("coordinator-owned movement guard should be reused")
.is_none()
);
drop(coordinator_guard);
tokio::time::timeout(std::time::Duration::from_secs(1), &mut writer)
.await
.expect("queued movement writer should proceed after the coordinator guard is released")
.expect("movement writer task should not panic");
}
#[tokio::test]
#[serial]
async fn heal_bucket_keeps_suspended_pool_volume_on_remove() {
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -2776,7 +2776,7 @@ mod tests {
rebalance_meta: RwLock::new(None),
decommission_cancelers: RwLock::new(Vec::new()),
start_gate: TokioMutex::new(()),
pool_meta_save_gate: TokioMutex::new(()),
pool_meta_save_gate: TokioMutex::default(),
ctx,
bucket_fence_registry: Arc::default(),
})
@@ -2999,7 +2999,7 @@ fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc<crate::store::ECSt
rebalance_meta: tokio::sync::RwLock::new(Some(meta)),
decommission_cancelers: tokio::sync::RwLock::new(Vec::new()),
start_gate: tokio::sync::Mutex::new(()),
pool_meta_save_gate: tokio::sync::Mutex::new(()),
pool_meta_save_gate: tokio::sync::Mutex::default(),
ctx: crate::runtime::instance::bootstrap_ctx(),
bucket_fence_registry: std::sync::Arc::default(),
})
+201 -7
View File
@@ -97,6 +97,8 @@ impl ECStore {
.first()
.cloned()
.ok_or_else(|| Error::other("heal format requires at least one storage pool"))?;
let mut write_state = self.pool_meta_save_gate.lock().await;
write_state.ensure_write_safe("heal format fence failed")?;
// Metadata fence order is part of the decommission/rebalance protocol:
// pool.bin must always be acquired before rebalance.bin.
@@ -110,8 +112,12 @@ impl ECStore {
}
let mut pool_meta = PoolMeta::default();
let replica_state = pool_meta.load_no_lock_from_replicas(self.pools.clone()).await?;
replica_state.ensure_write_safe("heal format fence failed")?;
let replica_state = pool_meta
.load_no_lock_from_replicas_observing(self.pools.clone(), &mut write_state)
.await?;
write_state.observe_replicas(replica_state);
write_state.ensure_write_safe("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
@@ -293,6 +299,10 @@ impl ECStore {
#[instrument(skip(self))]
pub(super) async fn handle_heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
let movement_gate = self.ctx.data_movement_operation_gate();
let _movement_guard = movement_gate.read().await;
let save_guard = self.pool_meta_save_gate.lock().await;
save_guard.ensure_write_safe("bucket heal cannot run while pool metadata requires recovery")?;
let mut fenced_pools = BTreeSet::new();
{
let pool_meta = self.pool_meta.read().await;
@@ -320,9 +330,10 @@ impl ECStore {
}
let dispatch_fenced_pools = fenced_pools.iter().copied().collect::<Vec<_>>();
drop(save_guard);
let mut res = self
.peer_sys
.heal_bucket_with_fence(bucket, opts, &dispatch_fenced_pools)
.heal_bucket_with_fence_from_movement_guarded_coordinator(bucket, opts, &dispatch_fenced_pools)
.await?;
{
let pool_meta = self.pool_meta.read().await;
@@ -479,17 +490,113 @@ impl ECStore {
mod tests {
use super::*;
use crate::bucket::metadata_sys;
use crate::core::pools::{PoolDecommissionInfo, PoolStatus};
use crate::cluster::rpc::PeerS3Client;
use crate::core::pools::{PoolDecommissionInfo, PoolMetaReplicaState, PoolStatus};
use crate::disk::error::Result as DiskResult;
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::runtime::instance::InstanceContext;
use crate::services::rebalance::{RebalanceInfo, RebalanceStats};
use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions};
use crate::storage_api_contracts::bucket::{
BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions,
};
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations};
use crate::store::init_format::{load_format_erasure, save_format_file};
use crate::store::init_local_disks_with_instance_ctx;
use tokio_util::sync::CancellationToken;
#[derive(Debug)]
struct BlockingHealPeer {
started: Arc<tokio::sync::Notify>,
release: Arc<tokio::sync::Notify>,
}
#[async_trait::async_trait]
impl PeerS3Client for BlockingHealPeer {
async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> DiskResult<HealResultItem> {
self.started.notify_one();
self.release.notified().await;
Ok(HealResultItem::default())
}
async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> DiskResult<()> {
Ok(())
}
async fn list_bucket(&self, _opts: &BucketOptions) -> DiskResult<Vec<BucketInfo>> {
Ok(Vec::new())
}
async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> DiskResult<()> {
Ok(())
}
async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> DiskResult<BucketInfo> {
Ok(BucketInfo::default())
}
fn get_pools(&self) -> Option<Vec<usize>> {
Some(vec![0, 1])
}
}
#[derive(Debug)]
struct WriterQueuedLocalHealPeer {
movement_gate: Arc<tokio::sync::RwLock<()>>,
writer_queued: Arc<tokio::sync::Notify>,
writer_acquired: Arc<tokio::sync::Notify>,
}
#[async_trait::async_trait]
impl PeerS3Client for WriterQueuedLocalHealPeer {
async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> DiskResult<HealResultItem> {
let _movement_guard = self
.movement_gate
.try_read()
.map_err(|_| crate::error::StorageError::SlowDown)?;
Ok(HealResultItem::default())
}
async fn heal_bucket_with_fence_from_movement_guarded_coordinator(
&self,
_bucket: &str,
_opts: &HealOpts,
_fenced_pools: &[usize],
) -> DiskResult<HealResultItem> {
Ok(HealResultItem::default())
}
async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> DiskResult<()> {
Ok(())
}
async fn list_bucket(&self, _opts: &BucketOptions) -> DiskResult<Vec<BucketInfo>> {
Ok(Vec::new())
}
async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> DiskResult<()> {
Ok(())
}
async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> DiskResult<BucketInfo> {
let movement_gate = self.movement_gate.clone();
let writer_acquired = self.writer_acquired.clone();
tokio::spawn(async move {
let _movement_guard = movement_gate.write().await;
writer_acquired.notify_one();
});
while self.movement_gate.try_read().is_ok() {
tokio::task::yield_now().await;
}
self.writer_queued.notify_one();
Ok(BucketInfo::default())
}
fn get_pools(&self) -> Option<Vec<usize>> {
Some(vec![0, 1])
}
}
async fn minimal_heal_pool(pool_idx: usize) -> Arc<Sets> {
let format = FormatV3::new(1, 1);
let endpoint_url = format!("http://127.0.0.1:{}/data", 19000 + pool_idx);
@@ -529,7 +636,7 @@ mod tests {
rebalance_meta: RwLock::new(None),
decommission_cancelers: RwLock::new(Vec::new()),
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::default(),
ctx: crate::runtime::instance::bootstrap_ctx(),
bucket_fence_registry: std::sync::Arc::default(),
}
@@ -955,6 +1062,93 @@ mod tests {
);
}
#[tokio::test]
async fn bucket_heal_blocks_before_dispatch_after_unreadable_pool_meta_replica() {
let store = minimal_heal_store().await;
store.pool_meta_save_gate.lock().await.observe_replicas(PoolMetaReplicaState {
needs_repair: true,
repair_write_safe: false,
});
let err = store
.handle_heal_bucket("bucket", &HealOpts::default())
.await
.expect_err("bucket heal must stay blocked until restart after an unreadable replica");
assert!(
err.to_string()
.contains("restart after all replicas are readable and consistent")
);
}
#[tokio::test]
async fn bucket_heal_releases_save_gate_before_peer_dispatch_and_holds_movement_snapshot() {
let mut store = minimal_heal_store().await;
let started = Arc::new(tokio::sync::Notify::new());
let release = Arc::new(tokio::sync::Notify::new());
let peer: Box<dyn PeerS3Client> = Box::new(BlockingHealPeer {
started: started.clone(),
release: release.clone(),
});
store.peer_sys.clients = vec![Arc::new(peer)];
let store = Arc::new(store);
let movement_gate = store.ctx.data_movement_operation_gate();
let mut heal = tokio::spawn({
let store = store.clone();
async move { store.handle_heal_bucket("bucket", &HealOpts::default()).await }
});
tokio::time::timeout(std::time::Duration::from_secs(1), started.notified())
.await
.expect("peer dispatch should start");
assert!(
store.pool_meta_save_gate.try_lock().is_ok(),
"coordinator must release its local save gate before waiting for peers"
);
assert!(
movement_gate.try_write().is_err(),
"bucket heal must hold the movement snapshot through peer dispatch"
);
release.notify_one();
tokio::time::timeout(std::time::Duration::from_secs(1), &mut heal)
.await
.expect("bucket heal should finish after peer release")
.expect("bucket heal task should not panic")
.expect("bucket heal should succeed");
}
#[tokio::test]
async fn bucket_heal_local_fanout_does_not_reenter_movement_read_behind_queued_writer() {
let mut store = minimal_heal_store().await;
let movement_gate = store.ctx.data_movement_operation_gate();
let writer_queued = Arc::new(tokio::sync::Notify::new());
let writer_acquired = Arc::new(tokio::sync::Notify::new());
let peer: Box<dyn PeerS3Client> = Box::new(WriterQueuedLocalHealPeer {
movement_gate: movement_gate.clone(),
writer_queued: writer_queued.clone(),
writer_acquired: writer_acquired.clone(),
});
store.peer_sys.clients = vec![Arc::new(peer)];
let store = Arc::new(store);
let mut heal = tokio::spawn({
let store = store.clone();
async move { store.handle_heal_bucket("bucket", &HealOpts::default()).await }
});
tokio::time::timeout(std::time::Duration::from_secs(1), writer_queued.notified())
.await
.expect("movement writer should queue during local peer lookup");
tokio::time::timeout(std::time::Duration::from_secs(1), &mut heal)
.await
.expect("local fan-out must not reenter movement read behind the queued writer")
.expect("bucket heal task should not panic")
.expect("bucket heal should succeed");
tokio::time::timeout(std::time::Duration::from_secs(1), writer_acquired.notified())
.await
.expect("queued movement writer should proceed after bucket heal releases its read guard");
}
#[tokio::test]
#[serial_test::serial]
async fn unscoped_heal_object_suspended_owner_semantics() {
@@ -1282,7 +1476,7 @@ mod tests {
rebalance_meta: RwLock::new(None),
decommission_cancelers: RwLock::new(Vec::new()),
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::default(),
ctx: crate::runtime::instance::bootstrap_ctx(),
bucket_fence_registry: std::sync::Arc::default(),
};
+132 -33
View File
@@ -13,7 +13,9 @@
// limitations under the License.
use super::*;
use crate::core::pools::{PoolMetaReplicaState, local_decommission_queue_prefix, pool_meta_has_active_decommission};
use crate::core::pools::{
PoolMetaReplicaState, PoolMetaWriteState, local_decommission_queue_prefix, pool_meta_has_active_decommission,
};
use crate::error::is_err_decommission_running;
use crate::runtime::instance::InstanceContext;
use crate::runtime::sources as runtime_sources;
@@ -109,6 +111,14 @@ fn should_auto_start_rebalance_after_init(decommission_running: bool, rebalance_
rebalance_meta_loaded && !decommission_running
}
fn should_schedule_local_decommission_resume(
pool_indices: &[usize],
pool_meta_replica_state: PoolMetaReplicaState,
pool_meta_write_safe: bool,
) -> bool {
!pool_indices.is_empty() && pool_meta_replica_state.repair_write_safe && pool_meta_write_safe
}
async fn wait_for_local_decommission_resume_delay(rx: &CancellationToken, delay: Duration) -> bool {
tokio::select! {
_ = rx.cancelled() => false,
@@ -120,15 +130,19 @@ fn resolve_store_init_stage_result(result: Result<()>, stage: &str) -> Result<()
result.map_err(|err| Error::other(format!("store init failed during {stage}: {err}")))
}
async fn load_pool_meta_for_startup<S>(pools: Vec<Arc<S>>) -> Result<(PoolMeta, PoolMetaReplicaState)>
async fn load_pool_meta_for_startup<S>(
pools: Vec<Arc<S>>,
write_state: &mut PoolMetaWriteState,
) -> Result<(PoolMeta, PoolMetaReplicaState)>
where
S: EcstoreObjectIO,
{
let mut meta = PoolMeta::default();
let replica_state = meta
.load_no_lock_from_replicas(pools)
.load_no_lock_from_replicas_observing(pools, write_state)
.await
.map_err(|err| Error::other(format!("store init failed during load_pool_meta: {err}")))?;
write_state.observe_replicas(replica_state);
Ok((meta, replica_state))
}
@@ -143,6 +157,7 @@ async fn persist_pool_meta_for_startup_if_safe<S>(
meta: &PoolMeta,
pools: Vec<Arc<S>>,
replica_state: PoolMetaReplicaState,
write_state: PoolMetaWriteState,
topology_update: bool,
elected_writer: bool,
) -> Result<()>
@@ -152,10 +167,14 @@ where
if !elected_writer {
return Ok(());
}
let should_write = topology_update || (replica_state.needs_repair && replica_state.repair_write_safe);
if topology_update {
replica_state.ensure_write_safe("store init failed during save_validated_pool_meta")?;
}
if topology_update || (replica_state.needs_repair && replica_state.repair_write_safe) {
if should_write {
write_state.ensure_write_safe("store init failed during save_validated_pool_meta")?;
}
if should_write {
save_validated_pool_meta_for_startup(meta, pools).await?;
}
Ok(())
@@ -431,7 +450,7 @@ impl ECStore {
rebalance_meta: RwLock::new(None),
decommission_cancelers,
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::default(),
// Adopt the caller's context (the process bootstrap one on the
// legacy path) so startup writes (erasure type recorded before
// this point) and later reads share one cell.
@@ -475,7 +494,10 @@ impl ECStore {
pub async fn init(self: &Arc<Self>, rx: CancellationToken) -> Result<()> {
runtime_sources::ensure_boot_time().await;
let (meta, pool_meta_replica_state) = load_pool_meta_for_startup(self.pools.clone()).await?;
let (meta, pool_meta_replica_state) = {
let mut write_state = self.pool_meta_save_gate.lock().await;
load_pool_meta_for_startup(self.pools.clone(), &mut write_state).await?
};
let update = meta.validate(self.pools.clone())?;
let endpoints = runtime_sources::endpoint_pools_or_default();
let should_persist_pool_meta = runtime_sources::first_cluster_node_is_local().await;
@@ -487,14 +509,18 @@ impl ECStore {
};
// Only one local node should persist validated pool metadata here; otherwise
// distributed startup can race on the same lock and replay the prior init bug.
persist_pool_meta_for_startup_if_safe(
&installed_pool_meta,
self.pools.clone(),
pool_meta_replica_state,
update,
should_persist_pool_meta,
)
.await?;
{
let write_state = self.pool_meta_save_gate.lock().await;
persist_pool_meta_for_startup_if_safe(
&installed_pool_meta,
self.pools.clone(),
pool_meta_replica_state,
*write_state,
update,
should_persist_pool_meta,
)
.await?;
}
{
let mut pool_meta = self.pool_meta.write().await;
@@ -538,7 +564,11 @@ impl ECStore {
}
let local_pool_indices = local_decommission_queue_prefix(&endpoints, &pool_indices)?;
if !local_pool_indices.is_empty() {
let pool_meta_write_safe = self
.ensure_pool_meta_side_effects_safe("decommission resume blocked while pool metadata requires recovery")
.await
.is_ok();
if should_schedule_local_decommission_resume(&local_pool_indices, pool_meta_replica_state, pool_meta_write_safe) {
let store = self.clone();
tokio::spawn(async move {
@@ -547,6 +577,16 @@ impl ECStore {
}
resume_local_decommission_after_init(store, rx, local_pool_indices).await;
});
} else if !local_pool_indices.is_empty() {
error!(
event = EVENT_DECOMMISSION_RESUME_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_STORE_INIT,
state = "blocked",
pool_indices = ?local_pool_indices,
reason = "pool_meta_write_blocked",
"Decommission resume blocked until pool metadata replicas are readable and consistent"
);
}
runtime_sources::init_bucket_monitor_for_current_endpoints();
@@ -575,11 +615,11 @@ impl ECStore {
#[cfg(test)]
mod tests {
use super::{
LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES, load_pool_meta_for_startup, persist_pool_meta_for_startup_if_safe,
pool_first_endpoint_is_local, pool_meta_has_active_decommission, preflight_startup_rpc_secret_with,
resolve_startup_pool_defaults_with, resolve_store_init_stage_result, save_validated_pool_meta_for_startup,
should_auto_start_rebalance_after_init, should_retry_format_load, should_retry_local_decommission_resume,
wait_for_local_decommission_resume_delay,
LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES, PoolMetaWriteState, load_pool_meta_for_startup,
persist_pool_meta_for_startup_if_safe, pool_first_endpoint_is_local, pool_meta_has_active_decommission,
preflight_startup_rpc_secret_with, resolve_startup_pool_defaults_with, resolve_store_init_stage_result,
save_validated_pool_meta_for_startup, should_auto_start_rebalance_after_init, should_retry_format_load,
should_retry_local_decommission_resume, wait_for_local_decommission_resume_delay,
};
#[cfg(feature = "test-util")]
use crate::disk::DiskAPI;
@@ -796,8 +836,9 @@ mod tests {
#[tokio::test]
async fn test_store_init_pool_meta_io_bypasses_namespace_lock_surface() {
let storage = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let mut write_state = PoolMetaWriteState::default();
let (loaded, replica_state) = load_pool_meta_for_startup(vec![storage.clone()])
let (loaded, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut write_state)
.await
.expect("startup pool metadata load should tolerate missing metadata without locks");
assert!(loaded.pools.is_empty());
@@ -822,8 +863,9 @@ mod tests {
let corrupt = Arc::new(StartupPoolMetaStorage::new(vec![0, 1, 2]));
let expected = init_test_pool_meta(None);
let backup = Arc::new(StartupPoolMetaStorage::new(startup_pool_meta_payload(&expected)));
let mut write_state = PoolMetaWriteState::default();
let (loaded, replica_state) = load_pool_meta_for_startup(vec![corrupt.clone(), backup.clone()])
let (loaded, replica_state) = load_pool_meta_for_startup(vec![corrupt.clone(), backup.clone()], &mut write_state)
.await
.expect("startup should select the validated backup replica");
@@ -834,9 +876,16 @@ mod tests {
assert!(corrupt.read_without_lock.load(Ordering::SeqCst));
assert!(backup.read_without_lock.load(Ordering::SeqCst));
persist_pool_meta_for_startup_if_safe(&loaded, vec![corrupt.clone(), backup.clone()], replica_state, false, true)
.await
.expect("the elected startup writer should repair validated corrupt replicas");
persist_pool_meta_for_startup_if_safe(
&loaded,
vec![corrupt.clone(), backup.clone()],
replica_state,
write_state,
false,
true,
)
.await
.expect("the elected startup writer should repair validated corrupt replicas");
let corrupt_write = corrupt
.written_payload
@@ -858,26 +907,76 @@ mod tests {
async fn test_store_init_pool_meta_does_not_repair_unreadable_replica() {
let valid = Arc::new(StartupPoolMetaStorage::new(startup_pool_meta_payload(&init_test_pool_meta(None))));
let unreadable = Arc::new(StartupPoolMetaStorage::unreadable());
let mut write_state = PoolMetaWriteState::default();
let (loaded, replica_state) = load_pool_meta_for_startup(vec![valid.clone(), unreadable.clone()])
let (loaded, replica_state) = load_pool_meta_for_startup(vec![valid.clone(), unreadable.clone()], &mut write_state)
.await
.expect("startup should use a validated replica without overwriting an unreadable copy");
assert!(replica_state.needs_repair);
assert!(!replica_state.repair_write_safe);
persist_pool_meta_for_startup_if_safe(&loaded, vec![valid.clone(), unreadable.clone()], replica_state, false, true)
.await
.expect("an unreadable copy should defer repair when no topology write is needed");
persist_pool_meta_for_startup_if_safe(
&loaded,
vec![valid.clone(), unreadable.clone()],
replica_state,
write_state,
false,
true,
)
.await
.expect("an unreadable copy should defer repair when no topology write is needed");
assert!(!valid.wrote_without_lock.load(Ordering::SeqCst));
assert!(!unreadable.wrote_without_lock.load(Ordering::SeqCst));
let err =
persist_pool_meta_for_startup_if_safe(&loaded, vec![valid.clone(), unreadable.clone()], replica_state, true, true)
.await
.expect_err("a topology update must not overwrite an unreadable replica");
let err = persist_pool_meta_for_startup_if_safe(
&loaded,
vec![valid.clone(), unreadable.clone()],
replica_state,
write_state,
true,
true,
)
.await
.expect_err("a topology update must not overwrite an unreadable replica");
assert!(err.to_string().contains("cannot overwrite an unreadable replica"));
assert!(!valid.wrote_without_lock.load(Ordering::SeqCst));
assert!(!unreadable.wrote_without_lock.load(Ordering::SeqCst));
assert!(!super::should_schedule_local_decommission_resume(&[0], replica_state, true));
assert!(!super::should_schedule_local_decommission_resume(
&[0],
crate::core::pools::PoolMetaReplicaState {
needs_repair: false,
repair_write_safe: true,
},
false,
));
}
#[tokio::test]
async fn test_store_init_pool_meta_stays_blocked_after_all_replicas_were_unreadable() {
let unreadable_a = Arc::new(StartupPoolMetaStorage::unreadable());
let unreadable_b = Arc::new(StartupPoolMetaStorage::unreadable());
let mut write_state = PoolMetaWriteState::default();
load_pool_meta_for_startup(vec![unreadable_a, unreadable_b], &mut write_state)
.await
.expect_err("startup must fail when no readable pool metadata replica exists");
let repaired = Arc::new(StartupPoolMetaStorage::new(startup_pool_meta_payload(&init_test_pool_meta(None))));
let (loaded, replica_state) = load_pool_meta_for_startup(vec![repaired.clone()], &mut write_state)
.await
.expect("a later startup retry may read the repaired replica");
assert!(replica_state.repair_write_safe);
let err = persist_pool_meta_for_startup_if_safe(&loaded, vec![repaired.clone()], replica_state, write_state, true, true)
.await
.expect_err("the same store instance must not write after observing unreadable replicas");
assert!(
err.to_string()
.contains("restart after all replicas are readable and consistent")
);
assert!(!repaired.wrote_without_lock.load(Ordering::SeqCst));
assert!(!super::should_schedule_local_decommission_resume(&[0], replica_state, false));
}
#[test]
+7 -7
View File
@@ -33,7 +33,7 @@ use crate::bucket::utils::check_put_object_part_args;
use crate::bucket::utils::{check_valid_bucket_name, check_valid_bucket_name_strict, is_meta_bucketname};
use crate::cluster::rpc::{RemoteClient, S3PeerSys};
use crate::config::storageclass;
use crate::core::pools::{DecommissionCanceler, PoolMeta};
use crate::core::pools::{DecommissionCanceler, PoolMeta, PoolMetaWriteState};
use crate::disk::endpoint::{Endpoint, EndpointType};
use crate::disk::{DiskAPI, DiskInfo, DiskInfoOptions};
use crate::error::{Error, Result};
@@ -185,12 +185,12 @@ pub struct ECStore {
/// or `decommission_cancelers`. The guarded sections may perform bounded
/// async metadata work so check/init/start cannot race across operations.
pub(crate) start_gate: Mutex<()>,
/// Serializes full-document pool metadata saves.
/// Serializes full-document pool metadata saves and retains a fail-closed
/// write block after startup observes an unreadable replica.
///
/// Lock order: acquire `pool_meta_save_gate` without holding `pool_meta`.
/// The saver then clones the latest `pool_meta` under a short read lock and
/// releases it before awaiting disk writes.
pub(crate) pool_meta_save_gate: Mutex<()>,
/// Lock order: acquire `pool_meta_save_gate`, then the distributed
/// `pool.bin` fence, then clone `pool_meta` under a short read lock.
pub(crate) pool_meta_save_gate: Mutex<PoolMetaWriteState>,
/// Per-instance runtime state (Phase 5, backlog#939).
///
/// Carries this instance's identity/runtime out of the process globals so
@@ -1114,7 +1114,7 @@ mod tests {
rebalance_meta: RwLock::new(None),
decommission_cancelers: RwLock::new(Vec::new()),
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::default(),
ctx,
bucket_fence_registry: Arc::default(),
})
+1 -1
View File
@@ -960,7 +960,7 @@ mod tests {
rebalance_meta: RwLock::new(None),
decommission_cancelers: RwLock::new(Vec::new()),
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::default(),
ctx: crate::runtime::instance::bootstrap_ctx(),
bucket_fence_registry: std::sync::Arc::default(),
}
+2 -2
View File
@@ -5449,7 +5449,7 @@ mod tests {
rebalance_meta: RwLock::new(None),
decommission_cancelers: RwLock::new(Vec::new()),
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::default(),
ctx: crate::runtime::instance::bootstrap_ctx(),
bucket_fence_registry: std::sync::Arc::default(),
}
@@ -5512,7 +5512,7 @@ mod tests {
rebalance_meta: RwLock::new(None),
decommission_cancelers: RwLock::new(Vec::new()),
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::default(),
ctx,
bucket_fence_registry: std::sync::Arc::default(),
}
+59 -5
View File
@@ -722,9 +722,8 @@ impl ECStore {
// overwrite a newer local transition after the writer commits.
let movement_gate = self.ctx.data_movement_operation_gate();
let _movement_guard = movement_gate.write().await;
let mut reloaded = PoolMeta::default();
resolve_store_rebalance_pool_meta_reload_result(
reloaded.load(self.pools[0].clone(), self.pools.clone()).await,
let reloaded = resolve_store_rebalance_pool_meta_reload_result(
self.load_runtime_pool_meta("store rebalance pool meta reload failed").await,
"reload_pool_meta",
)?;
@@ -933,11 +932,12 @@ mod tests {
use super::*;
use crate::bucket::replication::{ReplicationStatusType, VersionPurgeStatusType};
use crate::config::storageclass::{CLASS_RRS, CLASS_STANDARD, lookup_config_for_pools_without_env};
use crate::core::pools::{POOL_META_VERSION, PoolDecommissionInfo, PoolStatus};
use crate::core::pools::{POOL_META_NAME, POOL_META_VERSION, PoolDecommissionInfo, PoolStatus};
use crate::disk::error::DiskError;
use crate::layout::endpoint::Endpoint;
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::object_api::ObjectLockConfigSnapshot;
use crate::set_disk::get_lock_acquire_timeout;
use crate::storage_api_contracts::bucket::MakeBucketOptions;
use crate::storage_api_contracts::object::ObjectIO as _;
use arc_swap::ArcSwap;
@@ -2122,7 +2122,7 @@ mod tests {
#[test]
fn resolve_store_rebalance_pool_meta_reload_result_wraps_error_context() {
let err = resolve_store_rebalance_pool_meta_reload_result(Err(Error::SlowDown), "reload_pool_meta")
let err = resolve_store_rebalance_pool_meta_reload_result::<()>(Err(Error::SlowDown), "reload_pool_meta")
.expect_err("failed pool meta reload should be wrapped");
let err_message = err.to_string();
assert!(err_message.contains("store rebalance pool meta reload failed during reload_pool_meta"));
@@ -2321,6 +2321,60 @@ mod tests {
.expect("pool meta snapshot should persist to every pool");
}
#[tokio::test]
#[serial_test::serial]
async fn pool_meta_runtime_load_waits_for_multi_pool_commit_fence() {
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-load-fence", &[2, 2]).await;
let old = PoolMeta::new(&store.pools, &PoolMeta::default());
old.save(store.pools.clone()).await.expect("old pool metadata should persist");
let mut newer = old.clone();
newer.pools[0].last_update += TimeDuration::seconds(1);
let pool_meta_lock = store.pools[0]
.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME)
.await
.expect("pool metadata lock should be created");
let pool_meta_guard = pool_meta_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.expect("pool metadata write fence should be acquired");
newer
.save_for_startup(vec![store.pools[0].clone()])
.await
.expect("first replica should enter the new generation");
let started = Arc::new(tokio::sync::Notify::new());
let mut load_task = tokio::spawn({
let store = store.clone();
let started = started.clone();
async move {
started.notify_one();
store.load_runtime_pool_meta("test runtime pool metadata load").await
}
});
started.notified().await;
assert!(
tokio::time::timeout(std::time::Duration::from_millis(100), &mut load_task)
.await
.is_err(),
"runtime load must not observe a valid-old/valid-new intermediate state"
);
newer
.save_for_startup(vec![store.pools[1].clone()])
.await
.expect("second replica should enter the new generation");
drop(pool_meta_guard);
let loaded = tokio::time::timeout(std::time::Duration::from_secs(5), load_task)
.await
.expect("runtime load should finish after commit publication")
.expect("runtime load task should not panic")
.expect("runtime load should select the completed snapshot");
assert_eq!(loaded.pools[0].last_update, newer.pools[0].last_update);
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn peer_pool_meta_reload_does_not_rollback_newer_local_states() {
@@ -67,7 +67,7 @@ pub(super) fn pool_lookup_not_found_error(bucket: &str, object: &str, opts: &Obj
}
}
pub(super) fn resolve_store_rebalance_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> {
pub(super) fn resolve_store_rebalance_pool_meta_reload_result<T>(result: Result<T>, stage: &str) -> Result<T> {
result.map_err(|err| Error::other(format!("store rebalance pool meta reload failed during {stage}: {err}")))
}