mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-08 13:06:00 +00:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 9524b39780 | |||
| b90011bab5 | |||
| dd6845279a | |||
| b4f95bb115 | |||
| 6ec3534d20 | |||
| 9c84384871 |
@@ -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]
|
||||
|
||||
@@ -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] {
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user