From dd6845279a9b38f66fb21d87c582794893bfbbab Mon Sep 17 00:00:00 2001 From: overtrue Date: Tue, 8 Sep 2026 15:50:33 +0800 Subject: [PATCH] fix(ecstore): resume followers after bootstrap metadata commits (cherry picked from commit d32b46a5e0918f73b96324f8fd51d07546ef7804) --- crates/ecstore/src/core/pools.rs | 154 +++++++++++++++++++++-- crates/ecstore/src/store/init.rs | 208 ++++++++++++++++++++++++++++++- 2 files changed, 352 insertions(+), 10 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index f70a65e99..726eaa07b 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -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, aborted_transaction: Arc, block_context: Option, transaction_failure: Arc>>, @@ -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, 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] { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 3d63a1c82..007277092 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -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();