feat(kms): retain historical master key versions for Vault KV2 rotation (#5484)

Rotation previously had to be rejected outright because replacing the
stored material would orphan every DEK wrapped by earlier versions.
Vault KV2 now keeps each version's material in an immutable, create-only
record at {prefix}/{key_id}/versions/{N} and treats the top-level record
as the current-version pointer plus fast-path material copy:

- decrypt resolves the envelope's master_key_version to its version
  record; a missing version fails closed with KeyVersionNotFound and
  never falls back to the current material
- envelopes without a version (pre-versioning writers) resolve to the
  baseline_version frozen at the key's first rotation, so never-rotated
  keys behave exactly as before
- generate_data_key stamps the wrapping version from the same key record
  snapshot that supplied the material
- rotate_key commits in check-and-set order: freeze baseline, persist
  the next version's material, then switch the current pointer; any
  failure leaves the current pointer untouched, and concurrent rotations
  serialize on the CAS writes with monotonically unique versions
- key listings drop the versions/ directory entries and physical key
  deletion purges version records before the key record

Refs rustfs/backlog#1565
This commit is contained in:
Zhengchao An
2026-07-31 01:48:10 +08:00
committed by GitHub
parent 8368017fb2
commit 1d3ba1eb8b
+582 -53
View File
@@ -28,7 +28,7 @@ use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use tracing::{debug, info, warn};
use vaultrs::kv2;
use vaultrs::{api::kv2::requests::SetSecretRequestOptions, error::ClientError, kv2};
/// Vault KMS client implementation
pub struct VaultKmsClient {
@@ -63,6 +63,70 @@ struct VaultKeyData {
tags: HashMap<String, String>,
/// Encrypted key material (base64 encoded)
encrypted_key_material: String,
/// Version that pre-versioning envelopes (no `master_key_version`) resolve to.
///
/// Recorded once, at the key's first rotation, when the then-current material is
/// frozen as an immutable version record. `None` means the key has never been
/// rotated, so legacy envelopes keep resolving to the current version — exactly
/// the pre-versioning behavior. Optional so records written by older builds keep
/// deserializing.
#[serde(default)]
baseline_version: Option<u32>,
}
/// Immutable per-version master key material record stored under
/// `{prefix}/{key_id}/versions/{N}`.
///
/// Version records are created with a KV2 check-and-set of 0 (create-only) and are
/// never rewritten, so every master key version that ever wrapped a DEK stays
/// readable after rotation. The top-level `{prefix}/{key_id}` record keeps a copy of
/// the current material as a fast path and so binaries that predate versioned
/// storage can still read never-rotated keys.
#[derive(Debug, Clone, Serialize, Deserialize)]
struct VaultKeyVersionRecord {
/// Master key version this record holds material for
version: u32,
/// Encrypted key material (base64 encoded)
encrypted_key_material: String,
/// When this version's material was created
created_at: Zoned,
}
/// Sub-path (under each key path) reserved for immutable version records.
const KEY_VERSIONS_SUBPATH: &str = "versions";
/// Drop KV2 directory entries from a key listing.
///
/// Once a key has version records, listing the key prefix returns both the key
/// record itself ("my-key") and a directory entry for its version sub-path
/// ("my-key/"); only the former is a key.
fn filter_key_directory_entries(keys: Vec<String>) -> Vec<String> {
keys.into_iter().filter(|key| !key.ends_with('/')).collect()
}
/// Resolve which master key version wrapped an envelope.
///
/// `Some` versions are honored verbatim: if the record for that version is missing
/// the lookup must fail closed with [`KmsError::KeyVersionNotFound`], never fall
/// back to the current material. `None` (a pre-versioning envelope) resolves to the
/// key's baseline version — the deterministic version whose material was current
/// before the first rotation froze it — and, for keys that were never rotated and
/// thus have no baseline, to the current version, which matches pre-versioning
/// behavior exactly.
fn resolve_envelope_master_key_version(
envelope_version: Option<u32>,
baseline_version: Option<u32>,
current_version: u32,
) -> u32 {
envelope_version.or(baseline_version).unwrap_or(current_version)
}
/// Whether a KV2 write failed its check-and-set precondition.
fn is_cas_conflict(error: &ClientError) -> bool {
matches!(
error,
ClientError::APIError { code: 400, errors } if errors.iter().any(|message| message.contains("check-and-set"))
)
}
/// Decode and validate the stored master key material of a [`VaultKeyData`] record.
@@ -131,6 +195,16 @@ impl VaultKmsClient {
format!("{}/{}", self.key_path_prefix, key_id)
}
/// Get the path of the immutable record holding one version's material
fn key_version_path(&self, key_id: &str, version: u32) -> String {
format!("{}/{}/{}/{}", self.key_path_prefix, key_id, KEY_VERSIONS_SUBPATH, version)
}
/// Get the directory path holding a key's version records
fn key_versions_dir(&self, key_id: &str) -> String {
format!("{}/{}/{}", self.key_path_prefix, key_id, KEY_VERSIONS_SUBPATH)
}
/// Encode key material for KV2 storage.
///
/// This is plain Base64 encoding, not encryption: the KV2 backend stores master key
@@ -148,26 +222,117 @@ impl VaultKmsClient {
.map_err(|e| KmsError::cryptographic_error("decrypt", e.to_string()))
}
/// Get the actual key material for a master key
async fn get_key_material(&self, key_id: &str) -> Result<Vec<u8>> {
let key_data = self.get_key_data(key_id).await?;
decode_stored_key_material(key_id, &key_data.encrypted_key_material).inspect_err(|error| {
warn!(key_id, %error, "Vault KMS key material failed validation");
/// Read the immutable material record of one key version.
///
/// A missing record fails closed with [`KmsError::KeyVersionNotFound`]; falling
/// back to the current material would decrypt with the wrong key at best and
/// mask a tampered envelope version at worst.
async fn get_key_version_record(&self, key_id: &str, version: u32) -> Result<VaultKeyVersionRecord> {
let path = self.key_version_path(key_id, version);
let record: VaultKeyVersionRecord =
kv2::read(&self.vault().client, &self.kv_mount, &path)
.await
.map_err(|e| match e {
ClientError::ResponseWrapError => KmsError::key_version_not_found(key_id, version),
ClientError::APIError { code: 404, .. } => KmsError::key_version_not_found(key_id, version),
_ => KmsError::backend_error(format!("Failed to read key version record from Vault: {e}")),
})?;
if record.version != version {
return Err(KmsError::material_corrupt(
key_id,
format!("version record at {path} claims version {} instead of {version}", record.version),
));
}
Ok(record)
}
/// Load master key material for a specific key version.
///
/// The top-level record is the authoritative copy for the current version (a
/// never-rotated key has no version records at all); any other version must have
/// an immutable version record.
async fn get_key_material_for_version(&self, key_id: &str, key_data: &VaultKeyData, version: u32) -> Result<Vec<u8>> {
let encrypted_material = if version == key_data.version {
key_data.encrypted_key_material.clone()
} else {
self.get_key_version_record(key_id, version).await?.encrypted_key_material
};
decode_stored_key_material(key_id, &encrypted_material).inspect_err(|error| {
warn!(key_id, version, %error, "Vault KMS key material failed validation");
})
}
/// Encrypt data using a master key
async fn encrypt_with_master_key(&self, key_id: &str, plaintext: &[u8]) -> Result<(Vec<u8>, Vec<u8>)> {
// Load the actual master key material
let key_material = self.get_key_material(key_id).await?;
self.dek_crypto.encrypt(&key_material, plaintext).await
/// Read the key record together with the KV2 secret version holding it, so a
/// later write can be check-and-set against exactly this snapshot.
async fn get_key_data_versioned(&self, key_id: &str) -> Result<(u32, VaultKeyData)> {
let path = self.key_path(key_id);
let metadata = kv2::read_metadata(&self.vault().client, &self.kv_mount, &path)
.await
.map_err(|e| match e {
ClientError::ResponseWrapError => KmsError::key_not_found(key_id),
ClientError::APIError { code: 404, .. } => KmsError::key_not_found(key_id),
_ => KmsError::backend_error(format!("Failed to read key metadata from Vault: {e}")),
})?;
let cas = u32::try_from(metadata.current_version)
.map_err(|_| KmsError::backend_error(format!("KV2 secret version for key {key_id} exceeds u32")))?;
// Read the exact secret version from the metadata to keep the (cas, data)
// pair consistent even if another writer lands in between.
let key_data: VaultKeyData = kv2::read_version(&self.vault().client, &self.kv_mount, &path, metadata.current_version)
.await
.map_err(|e| match e {
ClientError::ResponseWrapError => KmsError::key_not_found(key_id),
ClientError::APIError { code: 404, .. } => KmsError::key_not_found(key_id),
_ => KmsError::backend_error(format!("Failed to read key from Vault: {e}")),
})?;
Ok((cas, key_data))
}
/// Decrypt data using a master key
async fn decrypt_with_master_key(&self, key_id: &str, ciphertext: &[u8], nonce: &[u8]) -> Result<Vec<u8>> {
// Load the actual master key material
let key_material = self.get_key_material(key_id).await?;
self.dek_crypto.decrypt(&key_material, ciphertext, nonce).await
/// Check-and-set write of the key record.
///
/// `cas` must match the KV2 secret version currently holding the record.
/// Returns the secret version created by this write so a caller can chain
/// further check-and-set writes.
async fn cas_store_key_data(&self, key_id: &str, key_data: &VaultKeyData, cas: u32) -> Result<u32> {
let path = self.key_path(key_id);
let written =
kv2::set_with_options(&self.vault().client, &self.kv_mount, &path, key_data, SetSecretRequestOptions { cas })
.await
.map_err(|e| {
if is_cas_conflict(&e) {
KmsError::invalid_operation(format!(
"Concurrent modification of key {key_id} detected, retry the rotation"
))
} else {
KmsError::backend_error(format!("Failed to store key in Vault: {e}"))
}
})?;
u32::try_from(written.version)
.map_err(|_| KmsError::backend_error(format!("KV2 secret version for key {key_id} exceeds u32")))
}
/// Create-only write of an immutable version record (KV2 check-and-set of 0).
///
/// Returns `Ok(true)` when this call created the record and `Ok(false)` when a
/// record already exists at that version; the caller decides whether the
/// existing record is acceptable. The record is never overwritten.
async fn try_create_key_version_record(&self, key_id: &str, record: &VaultKeyVersionRecord) -> Result<bool> {
let path = self.key_version_path(key_id, record.version);
match kv2::set_with_options(&self.vault().client, &self.kv_mount, &path, record, SetSecretRequestOptions { cas: 0 }).await
{
Ok(_) => Ok(true),
Err(e) if is_cas_conflict(&e) => Ok(false),
Err(e) => Err(KmsError::backend_error(format!("Failed to store key version record in Vault: {e}"))),
}
}
/// Store key data in Vault
@@ -209,6 +374,7 @@ impl VaultKmsClient {
metadata: existing_key_data.metadata.clone(),
tags: request.tags.clone(),
encrypted_key_material: existing_key_data.encrypted_key_material.clone(), // Preserve the key material
baseline_version: existing_key_data.baseline_version,
};
debug!(
@@ -240,6 +406,7 @@ impl VaultKmsClient {
// List keys under the prefix
match kv2::list(&self.vault().client, &self.kv_mount, &self.key_path_prefix).await {
Ok(keys) => {
let keys = filter_key_directory_entries(keys);
debug!("Found {} keys in Vault", keys.len());
Ok(keys)
}
@@ -260,6 +427,24 @@ impl VaultKmsClient {
async fn delete_key(&self, key_id: &str) -> Result<()> {
let path = self.key_path(key_id);
// Purge immutable version records first: if any purge fails, the top-level
// record still exists and the deletion can be retried. The reverse order
// would leave orphaned master key material in Vault after the key vanished.
let versions_dir = self.key_versions_dir(key_id);
match kv2::list(&self.vault().client, &self.kv_mount, &versions_dir).await {
Ok(versions) => {
for version in versions {
let version_path = format!("{versions_dir}/{version}");
kv2::delete_metadata(&self.vault().client, &self.kv_mount, &version_path)
.await
.map_err(|e| KmsError::backend_error(format!("Failed to delete key version record from Vault: {e}")))?;
}
}
// No version records exist (the key was never rotated).
Err(ClientError::ResponseWrapError) | Err(ClientError::APIError { code: 404, .. }) => {}
Err(e) => return Err(KmsError::backend_error(format!("Failed to list key version records in Vault: {e}"))),
}
// For this specific key path, we can safely delete the metadata
// since each key has its own unique path under the prefix
kv2::delete_metadata(&self.vault().client, &self.kv_mount, &path)
@@ -282,8 +467,16 @@ impl KmsClient for VaultKmsClient {
// Generate random data key material using the existing method
let plaintext_key = generate_key_material(&request.key_spec)?;
// Encrypt the data key with the master key
let (encrypted_key, nonce) = self.encrypt_with_master_key(&request.master_key_id, &plaintext_key).await?;
// Encrypt the data key with the current master key material. Single read of
// the key record: the material we wrap with and the version we stamp into
// the envelope must come from the same snapshot, or a concurrent rotation
// could stamp a version that never wrapped this DEK.
let key_data = self.get_key_data(&request.master_key_id).await?;
let key_material =
decode_stored_key_material(&request.master_key_id, &key_data.encrypted_key_material).inspect_err(|error| {
warn!(key_id = %request.master_key_id, %error, "Vault KMS key material failed validation");
})?;
let (encrypted_key, nonce) = self.dek_crypto.encrypt(&key_material, &plaintext_key).await?;
// Create data key envelope with master key version for rotation support
let envelope = DataKeyEnvelope {
@@ -294,7 +487,7 @@ impl KmsClient for VaultKmsClient {
nonce,
encryption_context: request.encryption_context.clone(),
created_at: Zoned::now(),
master_key_version: None,
master_key_version: Some(key_data.version),
};
// Serialize the envelope as the ciphertext
@@ -354,9 +547,16 @@ impl KmsClient for VaultKmsClient {
}
}
// Decrypt the data key
// Decrypt the data key with the master key version that wrapped it
let key_data = self.get_key_data(&envelope.master_key_id).await?;
let version =
resolve_envelope_master_key_version(envelope.master_key_version, key_data.baseline_version, key_data.version);
let key_material = self
.get_key_material_for_version(&envelope.master_key_id, &key_data, version)
.await?;
let plaintext = self
.decrypt_with_master_key(&envelope.master_key_id, &envelope.encrypted_key, &envelope.nonce)
.dek_crypto
.decrypt(&key_material, &envelope.encrypted_key, &envelope.nonce)
.await?;
debug!("Vault KMS data decrypted");
@@ -386,6 +586,7 @@ impl KmsClient for VaultKmsClient {
metadata: HashMap::new(),
tags: HashMap::new(),
encrypted_key_material: encrypted_material,
baseline_version: None,
};
// Store in Vault
@@ -514,13 +715,102 @@ impl KmsClient for VaultKmsClient {
Ok(())
}
async fn rotate_key(&self, _key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
// Rotation previously overwrote the stored master key with fresh material, which
// permanently orphaned every DEK wrapped by the prior version. Reject before any
// storage access so existing material can never be touched.
Err(KmsError::invalid_operation(
"Vault KV2 rotation is unavailable until versioned key material retention lands",
))
/// Rotate the master key while keeping every historical version decryptable.
///
/// Commit protocol (all writes check-and-set, in this order):
/// 1. First rotation only: freeze the current material as an immutable version
/// record and persist `baseline_version` so pre-versioning envelopes resolve
/// to it deterministically.
/// 2. Persist the next version's material as an immutable version record
/// (create-only) before anything references it.
/// 3. Switch the current pointer: bump `version` and mirror the new material
/// into the top-level record in a single check-and-set write.
///
/// If any step fails the current pointer is untouched, so a failed, cancelled,
/// or interrupted rotation never exposes half-committed material. Concurrent
/// rotations are serialized by the check-and-set writes: at most one caller
/// commits each version and the losers fail without side effects on current.
async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
debug!("Rotating master key: {}", key_id);
let (mut cas, mut key_data) = self.get_key_data_versioned(key_id).await?;
// The material about to be frozen must be decodable: freezing poisoned
// material would give legacy envelopes a permanently broken baseline. This
// surfaces the same typed Material* errors as the read path.
decode_stored_key_material(key_id, &key_data.encrypted_key_material)
.inspect_err(|error| warn!(key_id, %error, "Vault KMS key material failed validation"))?;
// Step 1: freeze the baseline on first rotation.
if key_data.baseline_version.is_none() {
let baseline = VaultKeyVersionRecord {
version: key_data.version,
encrypted_key_material: key_data.encrypted_key_material.clone(),
created_at: key_data.created_at.clone(),
};
if !self.try_create_key_version_record(key_id, &baseline).await? {
// Either a previous rotation attempt crashed between freezing the
// baseline and recording it in metadata, or a concurrent rotation
// got here first. Both are benign only if the existing record holds
// exactly the material being frozen; anything else means the
// version history is inconsistent and rotation must not proceed.
let existing = self.get_key_version_record(key_id, key_data.version).await?;
if existing.encrypted_key_material != key_data.encrypted_key_material {
return Err(KmsError::internal_error(format!(
"version record {} of key {key_id} does not match the current key material; refusing to rotate",
key_data.version
)));
}
}
key_data.baseline_version = Some(key_data.version);
cas = self.cas_store_key_data(key_id, &key_data, cas).await?;
}
// Step 2: durably persist the next version's material before it can become
// current.
let new_version = key_data
.version
.checked_add(1)
.ok_or_else(|| KmsError::internal_error(format!("key {key_id} exhausted the version space")))?;
let generated = generate_key_material(&key_data.algorithm)?;
let mut new_material = self.encrypt_key_material(&generated).await?;
let record = VaultKeyVersionRecord {
version: new_version,
encrypted_key_material: new_material.clone(),
created_at: Zoned::now(),
};
if !self.try_create_key_version_record(key_id, &record).await? {
// A record for the next version already exists: an interrupted rotation
// persisted it and stopped before switching the current pointer, or a
// concurrent rotation just created it. Adopt the persisted material —
// it is immutable, fully durable, and has never been current — instead
// of failing the create-only write forever. The check-and-set switch
// below still lets at most one caller commit this version.
let existing = self.get_key_version_record(key_id, new_version).await?;
decode_stored_key_material(key_id, &existing.encrypted_key_material)?;
new_material = existing.encrypted_key_material;
}
// Step 3: switch the current pointer. The top-level copy of the material is
// the fast path for new encryptions and must always match `version`.
key_data.version = new_version;
key_data.encrypted_key_material = new_material;
self.cas_store_key_data(key_id, &key_data, cas).await?;
info!(key_id, version = new_version, "Vault KMS master key rotated");
Ok(MasterKeyInfo {
key_id: key_id.to_string(),
version: new_version,
algorithm: key_data.algorithm.clone(),
usage: key_data.usage.clone(),
status: key_data.status,
description: key_data.description.clone(),
metadata: key_data.metadata.clone(),
created_at: key_data.created_at.clone(),
rotated_at: Some(Zoned::now()),
created_by: None,
})
}
async fn health_check(&self) -> Result<()> {
@@ -919,19 +1209,88 @@ mod tests {
}
#[tokio::test]
async fn test_vault_kv2_rotate_key_rejected_without_touching_storage() {
// No Vault instance needed: rotation must be rejected before any storage access,
// so the call cannot read or overwrite key material.
async fn test_key_version_paths_stay_under_the_key() {
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
.await
.expect("client");
let err = client
.rotate_key("any-key", None)
.await
.expect_err("Vault KV2 rotation must be rejected");
assert!(matches!(err, KmsError::InvalidOperation { .. }), "expected InvalidOperation, got {err:?}");
assert!(err.to_string().contains("rotation is unavailable"));
assert_eq!(client.key_path("my-key"), "rustfs/kms/keys/my-key");
assert_eq!(client.key_versions_dir("my-key"), "rustfs/kms/keys/my-key/versions");
assert_eq!(client.key_version_path("my-key", 3), "rustfs/kms/keys/my-key/versions/3");
}
#[test]
fn test_filter_key_directory_entries_drops_version_dirs() {
// Listing the key prefix returns "my-key/" as a directory entry once
// my-key has version records; only real key records may be listed.
let listed = vec!["alpha".to_string(), "alpha/".to_string(), "beta".to_string()];
assert_eq!(filter_key_directory_entries(listed), vec!["alpha".to_string(), "beta".to_string()]);
}
#[test]
fn test_resolve_envelope_master_key_version_rules() {
// An explicit envelope version is honored verbatim, even when it differs
// from both the baseline and the current version: whether material exists
// for it is decided by the versioned lookup, never by falling back.
assert_eq!(resolve_envelope_master_key_version(Some(2), Some(1), 5), 2);
assert_eq!(resolve_envelope_master_key_version(Some(9), Some(1), 5), 9);
// A pre-versioning envelope resolves to the frozen baseline, not to
// whatever version happens to be current.
assert_eq!(resolve_envelope_master_key_version(None, Some(1), 5), 1);
// Never-rotated keys have no baseline; the current version is the only
// material that ever existed, matching pre-versioning behavior.
assert_eq!(resolve_envelope_master_key_version(None, None, 1), 1);
}
#[test]
fn test_vault_key_data_without_baseline_version_deserializes() {
// Key records written before versioned storage have no baseline_version
// field and must keep deserializing with None.
let key_data = VaultKeyData {
algorithm: "AES_256".to_string(),
usage: KeyUsage::EncryptDecrypt,
created_at: Zoned::now(),
status: KeyStatus::Active,
version: 1,
description: None,
metadata: HashMap::new(),
tags: HashMap::new(),
encrypted_key_material: general_purpose::STANDARD.encode([0x42u8; 32]),
baseline_version: Some(1),
};
let mut value = serde_json::to_value(&key_data).expect("serialize key data");
value
.as_object_mut()
.expect("key data serializes to an object")
.remove("baseline_version");
let legacy: VaultKeyData = serde_json::from_value(value).expect("legacy record must deserialize");
assert_eq!(legacy.baseline_version, None);
assert_eq!(legacy.version, 1);
}
#[test]
fn test_is_cas_conflict_only_matches_cas_failures() {
let cas = ClientError::APIError {
code: 400,
errors: vec!["check-and-set parameter did not match the current version".to_string()],
};
assert!(is_cas_conflict(&cas));
let other_400 = ClientError::APIError {
code: 400,
errors: vec!["invalid request".to_string()],
};
assert!(!is_cas_conflict(&other_400));
let not_found = ClientError::APIError {
code: 404,
errors: Vec::new(),
};
assert!(!is_cas_conflict(&not_found));
}
#[tokio::test]
@@ -947,29 +1306,197 @@ mod tests {
assert!(!format!("{info:?}").contains("Transit"));
}
fn integration_generate_request(key_id: &str) -> GenerateKeyRequest {
GenerateKeyRequest {
master_key_id: key_id.to_string(),
key_spec: "AES_256".to_string(),
key_length: Some(32),
encryption_context: Default::default(),
grant_tokens: Vec::new(),
}
}
fn integration_decrypt_request(ciphertext: Vec<u8>) -> DecryptRequest {
DecryptRequest {
ciphertext,
encryption_context: Default::default(),
grant_tokens: Vec::new(),
}
}
#[tokio::test]
#[ignore] // Requires a running Vault instance (dev mode)
async fn test_vault_kv2_rotate_rejected_and_material_untouched() {
async fn test_vault_kv2_decrypt_after_rotate() {
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
.await
.expect("client");
let key_id = format!("rotate-{}", uuid::Uuid::new_v4());
let key_id = format!("rotate-retain-{}", uuid::Uuid::new_v4());
client.create_key(&key_id, "AES_256", None).await.expect("create");
let before = client.get_key_data(&key_id).await.expect("read");
let request = integration_generate_request(&key_id);
let err = client
.rotate_key(&key_id, None)
let dk_v1 = client.generate_data_key(&request, None).await.expect("generate under v1");
let env_v1: DataKeyEnvelope = serde_json::from_slice(&dk_v1.ciphertext).expect("parse v1 envelope");
assert_eq!(env_v1.master_key_version, Some(1));
let rotated = client.rotate_key(&key_id, None).await.expect("rotate to v2");
assert_eq!(rotated.version, 2);
let dk_v2 = client.generate_data_key(&request, None).await.expect("generate under v2");
let env_v2: DataKeyEnvelope = serde_json::from_slice(&dk_v2.ciphertext).expect("parse v2 envelope");
assert_eq!(env_v2.master_key_version, Some(2), "new envelopes must carry the latest version");
let rotated = client.rotate_key(&key_id, None).await.expect("rotate to v3");
assert_eq!(rotated.version, 3);
let dk_v3 = client.generate_data_key(&request, None).await.expect("generate under v3");
let env_v3: DataKeyEnvelope = serde_json::from_slice(&dk_v3.ciphertext).expect("parse v3 envelope");
assert_eq!(env_v3.master_key_version, Some(3));
// A mixed batch of envelopes from every historical version must decrypt.
for (data_key, label) in [(&dk_v1, "v1"), (&dk_v3, "v3"), (&dk_v2, "v2"), (&dk_v1, "v1 again")] {
let plaintext = client
.decrypt(&integration_decrypt_request(data_key.ciphertext.clone()), None)
.await
.unwrap_or_else(|error| panic!("envelope wrapped under {label} must stay decryptable: {error}"));
assert_eq!(Some(plaintext), data_key.plaintext, "{label} plaintext must round-trip");
}
}
#[tokio::test]
#[ignore] // Requires a running Vault instance (dev mode)
async fn test_vault_kv2_rotate_does_not_orphan_legacy_envelopes() {
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
.await
.expect_err("Vault KV2 rotation must be rejected");
assert!(matches!(err, KmsError::InvalidOperation { .. }));
.expect("client");
let after = client.get_key_data(&key_id).await.expect("reread");
assert_eq!(
after.encrypted_key_material, before.encrypted_key_material,
"rejected rotation must leave stored key material untouched"
let key_id = format!("rotate-legacy-{}", uuid::Uuid::new_v4());
client.create_key(&key_id, "AES_256", None).await.expect("create");
// Simulate an envelope written by a pre-versioning build: same wrapped DEK,
// but without the master_key_version field.
let data_key = client
.generate_data_key(&integration_generate_request(&key_id), None)
.await
.expect("generate");
let mut envelope: serde_json::Value = serde_json::from_slice(&data_key.ciphertext).expect("parse envelope");
envelope
.as_object_mut()
.expect("envelope is an object")
.remove("master_key_version");
let legacy_ciphertext = serde_json::to_vec(&envelope).expect("serialize legacy envelope");
client.rotate_key(&key_id, None).await.expect("rotate to v2");
client.rotate_key(&key_id, None).await.expect("rotate to v3");
// The baseline rule must route the legacy envelope to the frozen version 1
// material even though the current version has moved on.
let plaintext = client
.decrypt(&integration_decrypt_request(legacy_ciphertext), None)
.await
.expect("legacy envelope must stay decryptable after rotation");
assert_eq!(Some(plaintext), data_key.plaintext);
let key_data = client.get_key_data(&key_id).await.expect("read");
assert_eq!(key_data.baseline_version, Some(1), "first rotation must pin the baseline");
assert_eq!(key_data.version, 3);
}
#[tokio::test]
#[ignore] // Requires a running Vault instance (dev mode)
async fn test_vault_kv2_envelope_version_tampering_fails_closed() {
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
.await
.expect("client");
let key_id = format!("rotate-tamper-{}", uuid::Uuid::new_v4());
client.create_key(&key_id, "AES_256", None).await.expect("create");
let data_key = client
.generate_data_key(&integration_generate_request(&key_id), None)
.await
.expect("generate");
client.rotate_key(&key_id, None).await.expect("rotate");
// Point the envelope at a version that has no material record.
let mut envelope: serde_json::Value = serde_json::from_slice(&data_key.ciphertext).expect("parse envelope");
envelope
.as_object_mut()
.expect("envelope is an object")
.insert("master_key_version".to_string(), serde_json::json!(999));
let tampered = serde_json::to_vec(&envelope).expect("serialize tampered envelope");
let error = client
.decrypt(&integration_decrypt_request(tampered), None)
.await
.expect_err("nonexistent version must fail closed, not fall back to current");
assert!(
matches!(error, KmsError::KeyVersionNotFound { version: 999, key_id: ref error_key_id } if *error_key_id == key_id),
"expected KeyVersionNotFound for version 999, got {error:?}"
);
assert_eq!(after.version, before.version, "rejected rotation must not bump the key version");
// The untampered envelope still decrypts through its recorded version.
let plaintext = client
.decrypt(&integration_decrypt_request(data_key.ciphertext.clone()), None)
.await
.expect("untampered envelope must still decrypt");
assert_eq!(Some(plaintext), data_key.plaintext);
}
#[tokio::test]
#[ignore] // Requires a running Vault instance (dev mode)
async fn test_vault_kv2_concurrent_rotate_versions_monotonic() {
use std::sync::Arc;
let client = Arc::new(
VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
.await
.expect("client"),
);
let key_id = format!("rotate-concurrent-{}", uuid::Uuid::new_v4());
client.create_key(&key_id, "AES_256", None).await.expect("create");
let attempts = 4;
let tasks: Vec<_> = (0..attempts)
.map(|_| {
let client = Arc::clone(&client);
let key_id = key_id.clone();
tokio::spawn(async move { client.rotate_key(&key_id, None).await })
})
.collect();
let mut successes = 0u32;
for task in tasks {
// Losing a check-and-set race is an expected error; committing is not
// required, but every commit must account for exactly one version bump.
if task.await.expect("join rotate task").is_ok() {
successes += 1;
}
}
assert!(successes >= 1, "at least one rotation must commit");
let key_data = client.get_key_data(&key_id).await.expect("read");
assert_eq!(
key_data.version,
1 + successes,
"each successful rotation must commit exactly one new version"
);
assert_eq!(key_data.baseline_version, Some(1));
// Every version has an immutable record with unique material, and the
// top-level fast-path copy matches the current version's record.
let mut materials = std::collections::HashSet::new();
for version in 1..=key_data.version {
let record = client
.get_key_version_record(&key_id, version)
.await
.unwrap_or_else(|error| panic!("version {version} must have a record: {error}"));
assert_eq!(record.version, version);
assert!(materials.insert(record.encrypted_key_material), "version materials must be unique");
}
let current_record = client
.get_key_version_record(&key_id, key_data.version)
.await
.expect("current version record");
assert_eq!(current_record.encrypted_key_material, key_data.encrypted_key_material);
}
#[tokio::test]
@@ -991,8 +1518,9 @@ mod tests {
client.store_key_data(&key_id, &key_data).await.expect("store corrupt");
// Reading the material must now ERROR, not silently regenerate + overwrite.
let poisoned = client.get_key_data(&key_id).await.expect("read poisoned");
let error = client
.get_key_material(&key_id)
.get_key_material_for_version(&key_id, &poisoned, poisoned.version)
.await
.expect_err("corrupted key material must yield an error, not a fresh key");
assert!(
@@ -1004,7 +1532,7 @@ mod tests {
let after = client.get_key_data(&key_id).await.expect("reread");
assert_eq!(
after.encrypted_key_material, "!!!not-base64!!!",
"get_key_material must not overwrite stored master key material on failure"
"the material read path must not overwrite stored master key material on failure"
);
}
@@ -1026,8 +1554,9 @@ mod tests {
key_data.encrypted_key_material = String::new();
client.store_key_data(&key_id, &key_data).await.expect("store empty");
let poisoned = client.get_key_data(&key_id).await.expect("read poisoned");
let error = client
.get_key_material(&key_id)
.get_key_material_for_version(&key_id, &poisoned, poisoned.version)
.await
.expect_err("empty key material must yield an error, not a fresh key");
assert!(
@@ -1039,7 +1568,7 @@ mod tests {
let after = client.get_key_data(&key_id).await.expect("reread");
assert!(
after.encrypted_key_material.is_empty(),
"get_key_material must not backfill missing master key material"
"the material read path must not backfill missing master key material"
);
}