diff --git a/crates/kms/src/backends/vault.rs b/crates/kms/src/backends/vault.rs index 109b9b12d..01cc4f042 100644 --- a/crates/kms/src/backends/vault.rs +++ b/crates/kms/src/backends/vault.rs @@ -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, /// 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, +} + +/// 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) -> Vec { + 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, + baseline_version: Option, + 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> { - 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 { + 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> { + 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, Vec)> { - // 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> { - // 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 { + 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 { + 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 { - // 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 { + 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(¬_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) -> 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" ); }