fix(ecstore): resume followers after bootstrap metadata commits

(cherry picked from commit d32b46a5e0918f73b96324f8fd51d07546ef7804)
This commit is contained in:
overtrue
2026-09-08 15:50:33 +08:00
parent b4f95bb115
commit dd6845279a
2 changed files with 352 additions and 10 deletions
+145 -9
View File
@@ -4449,9 +4449,15 @@ impl PoolMetaReplicaState {
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum PoolMetaWriteBlock {
PendingBootstrapIdentity,
RecoveryRequired,
}
#[derive(Debug, Clone, Default)]
pub(crate) struct PoolMetaWriteState {
write_blocked: bool,
write_block: Option<PoolMetaWriteBlock>,
aborted_transaction: Arc<AtomicBool>,
block_context: Option<crate::error::PoolMetadataError>,
transaction_failure: Arc<std::sync::Mutex<Option<crate::error::PoolMetadataError>>>,
@@ -4565,7 +4571,8 @@ impl PoolMetaWriteState {
}
fn block_with_context(&mut self, mut context: crate::error::PoolMetadataError) {
if !self.write_blocked {
// A later hard failure must make a bootstrap wait irreversible.
if self.write_block != Some(PoolMetaWriteBlock::RecoveryRequired) {
record_pool_meta_block_once(&self.block_started_at, &mut context);
self.block_context = Some(context);
} else if let Some(previous) = &self.block_context
@@ -4575,7 +4582,7 @@ impl PoolMetaWriteState {
context.since = previous.since;
self.block_context = Some(context);
}
self.write_blocked = true;
self.write_block = Some(PoolMetaWriteBlock::RecoveryRequired);
}
pub(crate) fn block_writes_after_fence_loss(&mut self) {
@@ -4651,6 +4658,31 @@ impl PoolMetaWriteState {
Ok(())
}
pub(crate) fn ensure_startup_metadata_can_initialize(&mut self) -> Result<()> {
if !(self.pool_meta_absent
&& self.expected_cluster_id.is_some()
&& self.cluster_epoch.is_some()
&& self.identity_is_pending()
&& self.identity_fresh_bootstrap_nonce.is_some()
&& !self.bootstrap_identity_proven())
{
return self.ensure_missing_metadata_can_initialize();
}
self.validate_missing_metadata_can_initialize().map_err(|err| {
let mut context = pool_metadata_error(
crate::error::PoolMetadataFailure::RecoveryRequired,
"metadata_absence",
Some(Arc::new(err)),
);
if self.write_block.is_none() {
record_pool_meta_block_once(&self.block_started_at, &mut context);
self.block_context = Some(context.clone());
self.write_block = Some(PoolMetaWriteBlock::PendingBootstrapIdentity);
}
Error::other(context)
})
}
pub(crate) fn ensure_missing_metadata_can_initialize(&mut self) -> Result<()> {
if !self.pool_meta_absent {
return Ok(());
@@ -4679,7 +4711,7 @@ impl PoolMetaWriteState {
}
pub(crate) fn ensure_write_safe(&self, operation: &str) -> Result<()> {
if !self.write_blocked && !self.aborted_transaction.load(Ordering::SeqCst) {
if self.write_block.is_none() && !self.aborted_transaction.load(Ordering::SeqCst) {
return Ok(());
}
let mut context = self
@@ -6980,6 +7012,24 @@ impl PoolMeta {
write_state
.observe_selection(&selection)
.map_err(|err| block_pool_meta_validation(write_state, err, "startup_selection"))?;
// A pending bootstrap identity only delays an unproven startup follower.
// Any intervening hard failure promotes the block and cannot be cleared here.
// Missing or stale copies may still need repair after a safe committed selection.
if write_state.write_block == Some(PoolMetaWriteBlock::PendingBootstrapIdentity)
&& write_state.identity_initialized == Some(true)
&& !selection.absent
&& selection.revision.is_generation_protocol()
&& selection.replica_state.repair_write_safe
&& write_state.active_transactions.load(Ordering::SeqCst) == 0
&& !write_state.aborted_transaction.load(Ordering::SeqCst)
{
write_state.write_block = None;
write_state.block_context = None;
*write_state
.block_started_at
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
}
*self = selection.meta;
Ok(selection.replica_state)
}
@@ -10818,7 +10868,7 @@ impl ECStore {
};
};
let transaction_aborted = write_state.aborted_transaction.load(Ordering::SeqCst);
let writes_ready = !write_state.write_blocked && !transaction_aborted;
let writes_ready = write_state.write_block.is_none() && !transaction_aborted;
let failure = if writes_ready {
None
} else {
@@ -10828,7 +10878,7 @@ impl ECStore {
PoolMetaWriteGateStatus {
writes_ready,
check_timed_out: false,
write_blocked: write_state.write_blocked,
write_blocked: write_state.write_block.is_some(),
transaction_aborted,
pool_meta_absent: write_state.pool_meta_absent,
identity_initialized: write_state.identity_initialized,
@@ -10856,7 +10906,7 @@ impl ECStore {
if state.ensure_write_safe("pool metadata recovery").is_ok() {
return Ok(false);
}
if state.write_blocked
if state.write_block.is_some()
|| state.active_transactions.load(Ordering::SeqCst) != 0
|| state
.transaction_failure
@@ -10887,7 +10937,7 @@ impl ECStore {
let mut state = self.pool_meta_save_gate.lock().await;
// Outcomes can outlive the save-gate guard. Do not let an old Drop
// relatch, or an older recovery clear a newly installed block.
if state.write_blocked || state.active_transactions.load(Ordering::SeqCst) != 0 {
if state.write_block.is_some() || state.active_transactions.load(Ordering::SeqCst) != 0 {
return Ok(false);
}
let Some(blocked) = state
@@ -10920,7 +10970,7 @@ impl ECStore {
active_transactions: Arc::default(),
block_context: None,
recovery_failure: None,
write_blocked: false,
write_block: None,
..state.clone()
};
let selection = load_pool_meta_for_transaction_recovery(self.pools.clone(), &mut candidate).await?;
@@ -20791,6 +20841,92 @@ mod pools_tests {
identity: StdMutex<Option<(Vec<u8>, String)>>,
}
#[tokio::test]
async fn pending_bootstrap_wait_preserves_active_and_aborted_transactions() {
for abort in [false, true] {
let pool = Arc::new(PartialPoolMetaWriteStorage::default());
let cluster_id = uuid::Uuid::new_v4();
let mut writer = super::PoolMetaWriteState::for_startup(cluster_id, true);
super::persist_pool_meta_identity_for_startup(vec![pool.clone()], &mut writer, false)
.await
.expect("create a real pending identity");
let mut follower = super::PoolMetaWriteState::for_startup(cluster_id, false);
let mut loaded = PoolMeta::default();
loaded
.load_for_startup_observing(vec![pool.clone()], &mut follower)
.await
.expect("read pending identity");
follower
.ensure_startup_metadata_can_initialize()
.expect_err("pending follower must be blocked");
loaded
.load_for_startup_observing(vec![pool.clone()], &mut writer)
.await
.expect("writer reads pending metadata");
writer
.ensure_startup_metadata_can_initialize()
.expect("writer owns bootstrap authority");
let requested = PoolMeta {
version: POOL_META_VERSION,
pools: vec![decommission_test_pool_status(0, None)],
..Default::default()
};
requested
.save_for_startup_observing(vec![pool.clone()], &mut writer)
.await
.expect("prepare and commit pool metadata");
super::persist_pool_meta_identity_for_startup(vec![pool.clone()], &mut writer, true)
.await
.expect("commit the identity");
let committed = pool.stored.lock().expect("metadata lock").clone();
let identity = pool.identity.lock().expect("identity lock").clone();
let revision = pool.revision.load(Ordering::SeqCst);
let mut arm = follower.arm_transaction();
arm.phase = Some("commit_cas");
loaded
.load_for_startup_observing(vec![pool.clone()], &mut follower)
.await
.expect("committed records remain readable");
assert_eq!(follower.active_transactions.load(Ordering::SeqCst), 1);
follower
.ensure_write_safe("active transaction")
.expect_err("a startup reread cannot retire an outstanding owner");
if !abort {
arm.disarm();
}
drop(arm);
assert_eq!(follower.active_transactions.load(Ordering::SeqCst), 0);
assert_eq!(follower.aborted_transaction.load(Ordering::SeqCst), abort);
loaded
.load_for_startup_observing(vec![pool.clone()], &mut follower)
.await
.expect("reread committed records after owner drop");
if abort {
follower
.ensure_write_safe("aborted transaction")
.expect_err("startup must not erase an unknown transaction");
assert_eq!(
follower
.transaction_failure
.lock()
.expect("failure lock")
.as_ref()
.expect("retained failure")
.phase,
"commit_cas"
);
} else {
follower
.ensure_write_safe("completed transaction")
.expect("a disarmed owner permits the validated bootstrap retry");
}
assert_eq!(pool.revision.load(Ordering::SeqCst), revision);
assert_eq!(*pool.stored.lock().expect("metadata lock"), committed);
assert_eq!(*pool.identity.lock().expect("identity lock"), identity);
}
}
#[tokio::test]
async fn pool_meta_recovery_reconciles_prepare_and_partial_commit_without_format_downgrade() {
for previous_version in [POOL_META_V1_VERSION, POOL_META_VERSION, POOL_META_GENERATION_VERSION] {
+207 -1
View File
@@ -160,7 +160,7 @@ where
.map_err(|err| Error::other_with_context("store init failed during load_pool_meta", err))?;
write_state.observe_replicas(replica_state);
write_state
.ensure_missing_metadata_can_initialize()
.ensure_startup_metadata_can_initialize()
.map_err(|err| Error::other(format!("store init failed during classify_pool_meta_absence: {err}")))?;
Ok((meta, replica_state))
}
@@ -1482,6 +1482,212 @@ mod tests {
);
}
#[tokio::test]
async fn test_pending_identity_follower_accepts_safe_repairable_commit() {
let deployment_id = Uuid::new_v4();
let canonical = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let backup = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let pools = vec![canonical.clone(), backup.clone()];
let mut writer = PoolMetaWriteState::for_startup(deployment_id, true);
establish_pool_meta_bootstrap_identity_if_proven(pools.clone(), &mut writer, true)
.await
.expect("create pending identities through the elected writer");
let pending_backup = backup
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let mut follower = PoolMetaWriteState::for_startup(deployment_id, false);
for _ in 0..2 {
load_pool_meta_for_startup(pools.clone(), &mut follower)
.await
.expect_err("repeated pending reads cannot authorize follower writes");
follower
.ensure_write_safe("pending follower")
.expect_err("pending writes stay blocked");
}
assert_eq!(canonical.pool_meta_write_attempts.load(Ordering::SeqCst), 0);
assert_eq!(backup.pool_meta_write_attempts.load(Ordering::SeqCst), 0);
let (_, replica_state) = load_pool_meta_for_startup(pools.clone(), &mut writer)
.await
.expect("load fresh metadata");
let mut requested = init_test_pool_meta(None);
let mut second = requested.pools[0].clone();
second.id = 1;
second.cmd_line = "pool-1".to_string();
requested.pools.push(second);
let committed = persist_pool_meta_for_startup_if_safe(&requested, pools.clone(), replica_state, &mut writer, true, true)
.await
.expect("commit metadata and identity through the normal startup path");
assert_eq!(canonical.pool_meta_write_attempts.load(Ordering::SeqCst), 2);
assert_eq!(backup.pool_meta_write_attempts.load(Ordering::SeqCst), 2);
// Expose the backup's actual pre-commit image: pending identity and missing pool.bin.
*backup.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner) = pending_backup.clone();
let canonical_objects = canonical
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
assert!(pool_meta_identity_initialized_for_test(&canonical_objects[POOL_META_IDENTITY_NAME].0).expect("decode identity"));
assert_eq!(
pool_meta_v3_commit_state_for_test(canonical_objects[POOL_META_NAME].0.clone()).expect("decode committed record"),
(1, true)
);
canonical
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(POOL_META_IDENTITY_NAME.to_string(), pending_backup[POOL_META_IDENTITY_NAME].clone());
load_pool_meta_for_startup(pools.clone(), &mut follower)
.await
.expect("committed metadata is readable before identity promotion");
follower
.ensure_write_safe("identity still pending")
.expect_err("pool metadata alone cannot finish bootstrap");
*canonical.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner) = canonical_objects.clone();
let (loaded, replica_state) = load_pool_meta_for_startup(pools.clone(), &mut follower)
.await
.expect("a validated committed canonical copy permits safe repair of the lagging backup");
assert!(replica_state.needs_repair);
assert!(replica_state.repair_write_safe);
assert!(follower.identity_requires_repair());
assert_eq!(
serde_json::to_value(&loaded).expect("loaded metadata"),
serde_json::to_value(&committed).expect("committed metadata")
);
follower
.ensure_write_safe("safe repairable startup")
.expect("repairable copies must not permanently block a follower");
persist_pool_meta_for_startup_if_safe(&loaded, pools, replica_state, &mut follower, false, false)
.await
.expect("a non-elected follower never repairs copies itself");
assert_eq!(canonical.pool_meta_write_attempts.load(Ordering::SeqCst), 2);
assert_eq!(backup.pool_meta_write_attempts.load(Ordering::SeqCst), 2);
assert_eq!(
*canonical.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
canonical_objects
);
assert_eq!(*backup.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner), pending_backup);
}
#[tokio::test]
async fn test_pending_identity_wait_stays_blocked_after_hard_startup_failure() {
#[derive(Clone, Copy, Debug)]
enum Fault {
Unreadable,
CorruptIdentity,
CorruptMetadata,
ConflictingIdentity,
ConflictingEpoch,
InitializedMetadataMissing,
}
for fault in [
Fault::Unreadable,
Fault::CorruptIdentity,
Fault::CorruptMetadata,
Fault::ConflictingIdentity,
Fault::ConflictingEpoch,
Fault::InitializedMetadataMissing,
] {
let deployment_id = Uuid::new_v4();
let storage = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let mut writer = PoolMetaWriteState::for_startup(deployment_id, true);
establish_pool_meta_bootstrap_identity_if_proven(vec![storage.clone()], &mut writer, true)
.await
.expect("create pending identity");
let pending_objects = storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let pending = StartupPoolMetaStorage::new(Vec::new());
*pending.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner) = pending_objects;
let pending = Arc::new(pending);
let mut follower = PoolMetaWriteState::for_startup(deployment_id, false);
load_pool_meta_for_startup(vec![storage.clone()], &mut follower)
.await
.expect_err("pending identity must reject writes");
let (_, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut writer)
.await
.expect("load fresh metadata");
let committed = persist_pool_meta_for_startup_if_safe(
&init_test_pool_meta(None),
vec![storage.clone()],
replica_state,
&mut writer,
true,
true,
)
.await
.expect("commit matching identity and metadata");
let committed_objects = storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let mut faulty = StartupPoolMetaStorage::new(Vec::new());
let mut faulty_objects = committed_objects.clone();
match fault {
Fault::Unreadable => faulty.read_error = true,
Fault::CorruptIdentity => faulty_objects.get_mut(POOL_META_IDENTITY_NAME).expect("identity exists").0 = vec![0],
Fault::CorruptMetadata => faulty_objects.get_mut(POOL_META_NAME).expect("metadata exists").0 = vec![0],
Fault::ConflictingIdentity | Fault::ConflictingEpoch => {
let (cluster_id, epoch) = match fault {
Fault::ConflictingIdentity => (Uuid::new_v4(), 1),
_ => (deployment_id, 2),
};
faulty_objects.get_mut(POOL_META_IDENTITY_NAME).expect("identity exists").0 =
crate::core::pools::initialized_pool_meta_identity_for_test(cluster_id, epoch)
.expect("encode valid conflicting identity");
}
Fault::InitializedMetadataMissing => {
faulty_objects.remove(POOL_META_NAME);
}
}
*faulty.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner) = faulty_objects.clone();
let faulty = Arc::new(faulty);
load_pool_meta_for_startup(vec![faulty.clone()], &mut follower)
.await
.expect_err("hard startup failure must reject this read");
let blocked = follower
.ensure_write_safe("after hard startup failure")
.expect_err("hard failure must block writes");
let failure = blocked.pool_metadata_failure().expect("typed hard failure");
assert_eq!(faulty.pool_meta_write_attempts.load(Ordering::SeqCst), 0, "{fault:?}");
assert_eq!(
*faulty.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
faulty_objects,
"{fault:?}"
);
load_pool_meta_for_startup(vec![pending], &mut follower)
.await
.expect_err("a later pending identity cannot downgrade the hard block");
let (loaded, replicas) = load_pool_meta_for_startup(vec![storage.clone()], &mut follower)
.await
.expect("later healthy committed records remain readable");
assert!(replicas.repair_write_safe);
assert!(!replicas.needs_repair);
assert_eq!(
serde_json::to_value(&loaded).expect("loaded metadata"),
serde_json::to_value(&committed).expect("committed metadata")
);
let still_blocked = follower
.ensure_write_safe("after healthy reread")
.expect_err("a hard failure must never downgrade to a temporary bootstrap wait");
let retained = still_blocked.pool_metadata_failure().expect("retained typed failure");
assert_eq!(retained.phase, failure.phase, "{fault:?}");
assert_eq!(retained.since, failure.since, "{fault:?}");
assert_eq!(storage.pool_meta_write_attempts.load(Ordering::SeqCst), 2, "{fault:?}");
assert_eq!(
*storage.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
committed_objects,
"{fault:?}"
);
}
}
#[tokio::test]
async fn test_nonfresh_identity_repair_crash_never_persists_pending_bootstrap_authority() {
let deployment_id = Uuid::new_v4();