Compare commits

...

6 Commits

Author SHA1 Message Date
overtrue 9524b39780 fix: continue manual transition past in-flight objects
(cherry picked from commit 9a24389646)
2026-09-08 16:35:44 +08:00
overtrue b90011bab5 test(heal): drain PUT renames before physical census
(cherry picked from commit 618d38372cab961f293789a3bf08b8a8925e9771)
2026-09-08 16:25:46 +08:00
overtrue dd6845279a fix(ecstore): resume followers after bootstrap metadata commits
(cherry picked from commit d32b46a5e0918f73b96324f8fd51d07546ef7804)
2026-09-08 16:25:15 +08:00
overtrue b4f95bb115 test(ecstore): expose pending identity readiness latch
(cherry picked from commit 832cc60f87c5b1d734a96c5b284cb2422251d45f)
2026-09-08 16:24:45 +08:00
overtrue 6ec3534d20 fix(kms): preserve directory availability errors on main
(cherry picked from commit c2a8e476f7)
2026-09-08 14:45:04 +08:00
overtrue 9c84384871 test(kms): cover directory outages and missing keys on main
(cherry picked from commit 8fc1f0c41d)
2026-09-08 14:45:04 +08:00
6 changed files with 583 additions and 11 deletions
@@ -1047,6 +1047,8 @@ mod tests {
let mut cluster = RustFSTestClusterEnvironment::new(4).await?;
cluster.set_env("RUSTFS_UNSAFE_BYPASS_DISK_CHECK", "true");
cluster.set_env("RUSTFS_HEAL_ENABLED", "true");
// Capture physical baselines after the PUT rename fanout has drained.
cluster.set_env("RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE", "false");
// Heal control uses the first lexicographically sorted grid host.
// Keep that coordinator distinct from the remote target at index 1.
cluster.nodes.sort_by(|left, right| left.url.cmp(&right.url));
@@ -59,6 +59,36 @@ async fn test_kms_key_directory_unavailable() -> Result<(), Box<dyn std::error::
assert_eq!(put_response.server_side_encryption(), Some(&ServerSideEncryption::Aes256));
// A missing key in the healthy store is a client error, unlike a store outage.
let missing_key_object = "test-missing-kms-key";
let missing_key_error = s3_client
.put_object()
.bucket(TEST_BUCKET)
.key(missing_key_object)
.body(aws_sdk_s3::primitives::ByteStream::from_static(b"must not be published"))
.server_side_encryption(ServerSideEncryption::AwsKms)
.ssekms_key_id("rustfs-e2e-test-missing-key")
.send()
.await
.expect_err("an unknown key in a healthy Local KMS store must reject the write");
assert_eq!(missing_key_error.raw_response().map(|response| response.status().as_u16()), Some(400));
assert_eq!(
missing_key_error.as_service_error().and_then(ProvideErrorMetadata::code),
Some("KMS.NotFoundException")
);
let missing_key_absence = s3_client
.get_object()
.bucket(TEST_BUCKET)
.key(missing_key_object)
.send()
.await
.expect_err("a write rejected by a missing KMS key must not publish an object");
assert_eq!(missing_key_absence.raw_response().map(|response| response.status().as_u16()), Some(404));
assert_eq!(
missing_key_absence.as_service_error().and_then(ProvideErrorMetadata::code),
Some("NoSuchKey")
);
// Temporarily rename the key directory to simulate unavailability
info!("🔧 Simulating key directory unavailability");
let backup_dir = format!("{}.backup", kms_env.kms_keys_dir);
@@ -4036,6 +4036,10 @@ impl ManualTransitionRunReport {
|| self.skipped_queue_timeout > 0
}
fn has_enqueue_backpressure(&self) -> bool {
self.skipped_queue_full > 0 || self.skipped_queue_closed > 0 || self.skipped_queue_timeout > 0
}
pub fn was_truncated(&self) -> bool {
self.truncated_by_limit || self.truncated_by_duration || self.cancelled
}
@@ -4251,7 +4255,7 @@ pub async fn enqueue_transition_for_existing_objects_scoped(
}
report.scanned = report.scanned.saturating_add(1);
enqueue_transition_with_lifecycle_report(Some(api.clone()), object, &lc, &src, &options, &mut report).await;
if report.has_partial_enqueue() {
if report.has_enqueue_backpressure() {
report.next_marker.clone_from(&previous_marker);
report.next_version_idmarker.clone_from(&previous_version_marker);
report.continuation_token =
@@ -9950,6 +9954,18 @@ mod tests {
assert_eq!(report.skipped_queue_closed, 0);
assert_eq!(report.skipped_queue_timeout, 0);
assert!(report.has_partial_enqueue());
assert!(report.has_enqueue_backpressure());
}
#[test]
fn manual_transition_in_flight_skip_does_not_stop_the_scan() {
let options = ManualTransitionRunOptions::default();
let mut report = ManualTransitionRunReport::new("bucket", &options);
report.record_enqueue_outcome(TransitionEnqueueOutcome::AlreadyInFlight);
assert!(report.has_partial_enqueue());
assert!(!report.has_enqueue_backpressure());
}
#[test]
+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] {
+316 -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))
}
@@ -1373,6 +1373,321 @@ mod tests {
);
}
#[tokio::test]
async fn test_pending_identity_follower_becomes_write_ready_after_elected_commit() {
let deployment_id = Uuid::new_v4();
let storage = Arc::new(StartupPoolMetaStorage::new(Vec::new()));
let mut elected_writer = PoolMetaWriteState::for_startup(deployment_id, true);
establish_pool_meta_bootstrap_identity_if_proven(vec![storage.clone()], &mut elected_writer, true)
.await
.expect("the elected writer should persist its pending bootstrap identity");
let pending_objects = storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
assert!(
!pool_meta_identity_initialized_for_test(&pending_objects[POOL_META_IDENTITY_NAME].0)
.expect("decode the pending identity")
);
assert!(!pending_objects.contains_key(POOL_META_NAME));
let mut follower = PoolMetaWriteState::for_startup(deployment_id, false);
let pending_error = load_pool_meta_for_startup(vec![storage.clone()], &mut follower)
.await
.expect_err("the follower must not initialize metadata using another node's pending identity");
assert!(pending_error.to_string().contains("no verified fresh-bootstrap proof"));
let blocked = follower
.ensure_write_safe("pending identity follower")
.expect_err("the pending identity must not authorize follower writes");
assert_eq!(blocked.pool_metadata_failure().expect("typed failure").phase, "metadata_absence");
assert_eq!(storage.pool_meta_write_attempts.load(Ordering::SeqCst), 0);
assert_eq!(
*storage.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
pending_objects,
"the rejected follower must leave both metadata objects and CAS tokens untouched"
);
let (initial, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut elected_writer)
.await
.expect("the elected writer's bootstrap proof should authorize missing metadata");
assert!(initial.pools.is_empty());
assert!(!replica_state.needs_repair);
assert!(replica_state.repair_write_safe);
let committed = persist_pool_meta_for_startup_if_safe(
&init_test_pool_meta(None),
vec![storage.clone()],
replica_state,
&mut elected_writer,
true,
true,
)
.await
.expect("the elected writer should finish the normal startup metadata transaction");
elected_writer
.ensure_write_safe("elected writer after commit")
.expect("both metadata CAS phases and the identity commit must complete");
assert_eq!(storage.pool_meta_write_attempts.load(Ordering::SeqCst), 2);
assert_eq!(
*storage
.pool_meta_written_versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
vec![3, 3]
);
let committed_objects = storage
.objects
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
assert!(
pool_meta_identity_initialized_for_test(&committed_objects[POOL_META_IDENTITY_NAME].0)
.expect("decode the committed identity")
);
assert_eq!(
pool_meta_v3_commit_state_for_test(committed_objects[POOL_META_NAME].0.clone())
.expect("decode the complete V3 pool metadata"),
(1, true)
);
// Retry the same state that observed the pending identity, after normal commit completion.
let (reloaded, replica_state) = load_pool_meta_for_startup(vec![storage.clone()], &mut follower)
.await
.expect("the original follower should validate matching committed identity and metadata");
assert!(!replica_state.needs_repair);
assert!(replica_state.repair_write_safe);
assert_eq!(
serde_json::to_value(&reloaded).expect("serialize the reloaded metadata"),
serde_json::to_value(&committed).expect("serialize the committed metadata")
);
persist_pool_meta_for_startup_if_safe(&reloaded, vec![storage.clone()], replica_state, &mut follower, false, false)
.await
.expect("the non-elected follower must adopt the committed metadata without writing");
let mut late_follower = PoolMetaWriteState::for_startup(deployment_id, false);
load_pool_meta_for_startup(vec![storage.clone()], &mut late_follower)
.await
.expect("a follower first observing these same committed records must accept them");
late_follower
.ensure_write_safe("late follower after commit")
.expect("the committed records themselves must permit readiness");
assert_eq!(storage.pool_meta_write_attempts.load(Ordering::SeqCst), 2);
assert_eq!(
*storage.objects.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
committed_objects,
"neither follower may rewrite the elected writer's committed bytes or CAS tokens"
);
follower.ensure_write_safe("original follower after elected commit").expect(
"a follower retry must become write-ready after the elected writer commits matching identity and pool metadata",
);
}
#[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();
+73
View File
@@ -1233,6 +1233,9 @@ impl LocalKmsClient {
async fn decode_stored_key(&self, key_id: &str) -> Result<(StoredMasterKey, Vec<u8>)> {
let key_path = self.master_key_path(key_id)?;
if !fs::try_exists(&key_path).await? {
// Only an accessible key store can establish that a single key is
// missing; a directory outage must retain its filesystem error.
let _ = fs::read_dir(&self.config.key_dir).await?;
return Err(KmsError::key_not_found(key_id));
}
@@ -2387,6 +2390,76 @@ mod tests {
(client, temp_dir)
}
#[tokio::test]
async fn local_key_directory_outage_is_io_error_and_recovers_original_key() {
let root = TempDir::new().expect("create isolated key store");
let key_dir = root.path().join("keys");
let unavailable_dir = root.path().join("keys-unavailable");
let config = KmsConfig::local(key_dir.clone()).with_insecure_development_defaults();
let backend = LocalKmsBackend::new(config).await.expect("start Local KMS");
let key_id = "directory-outage-key";
backend
.create_key(CreateKeyRequest {
key_name: Some(key_id.to_string()),
..Default::default()
})
.await
.expect("create the original key");
let request = |key_id: &str| GenerateDataKeyRequest {
key_id: key_id.to_string(),
key_spec: KeySpec::Aes256,
encryption_context: HashMap::new(),
};
let before = backend
.generate_data_key(request(key_id))
.await
.expect("generate a data key before the outage");
let missing_key = backend.generate_data_key(request("no-such-key")).await;
let key_path = key_dir.join(format!("{key_id}.key"));
let original_record = fs::read(&key_path).await.expect("read the original key record");
fs::rename(&key_dir, &unavailable_dir)
.await
.expect("make the key directory unavailable");
let unavailable = backend.generate_data_key(request(key_id)).await;
// Restore before checking the error so the failing regression leaves no
// orphaned key store; both paths also belong to the same temporary root.
fs::rename(&unavailable_dir, &key_dir)
.await
.expect("restore the original key directory");
let after = backend
.generate_data_key(request(key_id))
.await
.expect("generate a data key after directory restoration");
for data_key in [&before, &after] {
let decrypted = backend
.decrypt(DecryptRequest {
ciphertext: data_key.ciphertext_blob.clone(),
encryption_context: HashMap::new(),
grant_tokens: Vec::new(),
})
.await
.expect("the original master key must decrypt both data keys");
assert!(
decrypted.plaintext == data_key.plaintext_key,
"directory restoration must preserve the original key material"
);
}
assert!(
fs::read(&key_path).await.expect("read the restored key record") == original_record,
"reads and recovery must not rewrite the key record"
);
assert!(
matches!(missing_key, Err(KmsError::KeyNotFound { key_id }) if key_id == "no-such-key"),
"a missing key in a readable directory must remain KeyNotFound"
);
assert!(
matches!(unavailable, Err(KmsError::IoError { .. })),
"an unavailable key directory must remain an I/O error, not KeyNotFound"
);
}
/// With the AAD write switch on, the Local backend seals the stored
/// encryption context into the wrap exactly like KV2: the bound envelope
/// round-trips, a rewritten stored context fails authentication even with