Files
rustfs/crates/kms/src/backends/local.rs
T
Zhengchao An e3a8234bc9 fix: 12 P1 reliability/security defects from the full-repo audit (backlog#806) (#4256)
* fix(rio): reject corrupted short compressed/encrypted blocks instead of panicking

DecompressReader::poll_read and DecryptReader::poll_read sliced the block
body with a fixed `[0..16]` index to read the length varint. The body length
comes from an untrusted 24-bit header field, so a corrupted/truncated block
shorter than 16 bytes made the slice panic and crash the request task — a
read-path DoS on GET of tiered/corrupted data.

Pass the whole (arbitrary-length-safe) slice to uvarint and reject a
non-positive or out-of-range length prefix with InvalidData. Adds a repro
test for each reader; all existing round-trip tests still pass.

Refs rustfs/backlog#812

* fix(utils): close SSRF bypass via IPv4-mapped IPv6 addresses

validate_outbound_ip branched on the IpAddr variant, and the V6 branch's
is_loopback/is_unicast_link_local/is_unique_local checks never inspect the
embedded IPv4 of an IPv4-mapped address (::ffff:a.b.c.d). The metadata guard
also only matched the plain V4 169.254.169.254. So ::ffff:127.0.0.1,
::ffff:10.0.0.5 and ::ffff:169.254.169.254 all passed the outbound guard,
letting an attacker reach loopback/private/metadata endpoints.

Normalize IPv4-mapped IPv6 to its embedded IPv4 (via to_ipv4_mapped, which
matches only the true mapped form) before classification. Adds reject tests
for mapped loopback/private/metadata and an allow test for public IPv6.

Refs rustfs/backlog#813

* fix(ecstore): streaming last-part loss, GCS tier Range/remove, stat_all_dirs alignment

Four confirmed data-reliability defects:

- put_object_multipart_stream: the CompleteMultipartUpload part-collection loop
  used exclusive `1..total_parts_count`, dropping the final part (and collecting
  zero parts for a single-part object) — silently truncating the completed object.
  Extracted collect_complete_parts (1..=total_parts_count) with unit tests.
- GCS warm backend get() ignored the requested byte range, returning the whole
  object for a Range GET; now applies ReadRange::segment like the other backends.
- GCS warm backend remove() was an empty stub, so deleting a tiered object left
  it on GCS forever; now deletes via StorageControl (added a control-plane client),
  and in_use() actually lists (prefix-scoped) instead of always returning false.
- stat_all_dirs skipped None disk slots and dropped JoinErrors, returning a
  compressed, misaligned error vector; heal_object_dir then zipped it against the
  full disks array and could make_volume on the WRONG disk. Now returns one
  index-aligned entry per slot (None -> DiskNotFound), and heal no longer
  pre-fills the drive report (which would double it). Added an alignment test.

Refs rustfs/backlog#807

* fix(kms): stop Vault backend from destroying/reviving keys on failure

Two confirmed key-safety defects in the Vault KV2 backend:

- get_key_material() 'self-healed' a decrypt or wrong-length failure by minting a
  fresh random master key and overwriting the stored value. That destroys the
  original key material, making every DEK ever wrapped by it permanently
  undecryptable. Decryption must never mutate the stored key: both branches now
  return a cryptographic_error instead. (The empty-material bootstrap path, which
  only fills a never-initialized key, is intentionally left intact.)
- cancel_key_deletion() reset key_state to Enabled only in the returned response
  and never persisted it, so the key stayed PendingDeletion in storage and would
  still be reaped. It now writes the state back via update_key_metadata_in_storage
  and fails the request if the write fails.

Adds ignored (Vault-requiring) integration tests documenting both behaviours.

The third item (VaultTransit key state only in memory -> revived as Enabled after
restart) is deferred: a fail-closed guard would break restart availability for all
transit keys; the correct fix needs a persistent metadata store + Vault integration
testing. Tracked in rustfs/backlog#808.

Refs rustfs/backlog#808

* fix(admin): clamp STS AssumeRole duration; persist ImportBucketMetadata to disk

Two confirmed admin-API defects:

- Standard AssumeRole used the raw client-supplied DurationSeconds with no upper
  bound, so a caller could mint near-permanent temporary credentials. Clamp it to
  the AWS/MinIO STS window [900, 43200] (with 0 -> default 3600) via a shared
  clamp_assume_role_duration helper, and build the exp claim with saturating_add.
  This matches the existing AssumeRoleWithWebIdentity path.
- ImportBucketMetadata only mutated an in-memory map and returned 200, silently
  dropping every imported config. It now persists each non-empty config via
  metadata_sys::update (which merges onto existing on-disk metadata) and returns
  InternalError if a write fails. Mapping extracted to imported_configs_to_persist
  with unit tests.

Refs rustfs/backlog#809

* fix(heal): enqueue displacing request in release builds

push_displacing_lower_priority folded the real enqueue call into
debug_assert_eq!(self.push(request), Accepted). In release builds
(debug_assertions off) the whole macro — including its argument — is compiled
out, so after evicting a lower-priority queued item the new high-priority
request was silently dropped and never healed. Hoist self.push(request) out of
the assertion so the side effect runs in all builds. Adds a --release regression
test.

Refs rustfs/backlog#811

* fix(iam): propagate real delete_policy backend errors instead of swallowing them

delete_policy's is_from_notify path had its error handling inverted: a real
backend failure (disk IO / insufficient quorum) evicted the cache and returned
Ok(()), reporting a phantom success while policy.json survived on disk (to be
reloaded on the next full IAM reload); NoSuchPolicy — which should be idempotent
success — returned Err. Propagate real errors and let NoSuchPolicy fall through
to the idempotent cache-evict + Ok, matching delete_user / the notification
handler in the same file. Adds a backend-error-injection regression test.

Refs rustfs/backlog#810

* fix(utils): also normalize IPv4-compatible IPv6 in the SSRF guard

The initial fix only unwrapped IPv4-mapped (::ffff:a.b.c.d) addresses; the
deprecated IPv4-compatible form (::a.b.c.d, e.g. ::127.0.0.1 / ::169.254.169.254)
still bypassed the guard. Reject pure-IPv6 specials (::, ::1, fe80::, fc00::)
first, then normalize BOTH embedded-IPv4 forms before the IPv4 rules. Adds tests
for compatible-form loopback/metadata and confirms ::1 / :: stay rejected.

Found by adversarial review of the initial fix. Refs rustfs/backlog#813

* fix(ecstore): fix the same last-part loss in the parallel streaming path

put_object_multipart_stream_parallel had the identical off-by-one
(1..total_parts_count) that truncated the last part / produced zero parts for a
single-part upload — reachable when concurrent stream parts are enabled. Reuse
collect_complete_parts, which now returns an error instead of panicking on a gap
in the parts map. Adds a missing-part error test.

Found by adversarial review of the initial fix. Refs rustfs/backlog#807

* fix(kms): local backend must preserve key material on status change

LocalKmsClient (the default KMS backend) regenerated the master key material on
enable_key/disable_key/schedule_key_deletion/cancel_key_deletion — a pure status
change. A single disable+enable cycle therefore destroyed the original key,
making every DEK ever wrapped by it permanently undecryptable (silent data loss,
no network needed). Preserve the existing material via get_key_material and
re-save with only the status changed. Adds a hermetic regression test that wraps
a DEK, cycles all four status methods, and asserts the DEK still decrypts.

Found by adversarial review of the Vault fix. Refs rustfs/backlog#808

* test(rio): cover the length-prefix guard; correct its comment

Add a DecompressReader test that feeds an unterminated length varint so uvarint
returns 0 and the new guard (not the downstream codec) produces the InvalidData
error, and reword the guard comment which overclaimed that the > len bound
prevents a reachable panic (it is belt-and-suspenders). No behavior change.

Found by adversarial review. Refs rustfs/backlog#812

* test(rio): build test block headers via vec! to satisfy clippy

The new corrupted-block tests built the header with Vec::new() + repeated push,
tripping clippy::vec_init_then_push (-D warnings in CI). Construct the fixed
header bytes with vec![] instead. No behavior change.

---------

Co-authored-by: houseme <housemecn@gmail.com>
2026-07-04 14:24:02 +08:00

1226 lines
47 KiB
Rust

// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Local file-based KMS backend implementation
use crate::backends::{BackendInfo, KmsBackend, KmsClient};
use crate::config::KmsConfig;
use crate::config::LocalConfig;
use crate::encryption::{AesDekCrypto, DataKeyEnvelope, DekCrypto, generate_key_material};
use crate::error::{KmsError, Result};
use crate::types::*;
use aes_gcm::{
Aes256Gcm, Key, Nonce,
aead::{Aead, KeyInit},
};
use argon2::{Algorithm, Argon2, Params, Version};
use async_trait::async_trait;
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64};
use jiff::Zoned;
use rand::RngExt;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::path::PathBuf;
use std::time::Duration;
use tokio::fs;
use tokio::sync::RwLock;
use tracing::{debug, warn};
const LOCAL_KMS_MASTER_KEY_SALT_FILE: &str = ".master-key.salt";
const LOCAL_KMS_MASTER_KEY_SALT_LEN: usize = 16;
const LOCAL_KMS_MASTER_KEY_LEN: usize = 32;
const LOCAL_KMS_ARGON2_M_COST_KIB: u32 = 19 * 1024;
const LOCAL_KMS_ARGON2_T_COST: u32 = 2;
const LOCAL_KMS_ARGON2_P_COST: u32 = 1;
/// Local KMS client that stores keys in local files
pub struct LocalKmsClient {
config: LocalConfig,
/// In-memory cache of loaded keys for performance
key_cache: RwLock<HashMap<String, MasterKeyInfo>>,
/// Master encryption key for encrypting stored keys
master_cipher: Option<Aes256Gcm>,
/// DEK encryption implementation
dek_crypto: AesDekCrypto,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
enum StoredKeyProtection {
#[default]
LegacyUnspecified,
EncryptedMasterKey,
PlaintextDevOnly,
}
/// Serializable representation of a master key stored on disk
#[derive(Debug, Clone, Serialize, Deserialize)]
struct StoredMasterKey {
key_id: String,
version: u32,
algorithm: String,
usage: KeyUsage,
status: KeyStatus,
description: Option<String>,
metadata: HashMap<String, String>,
#[serde(with = "crate::time_serde::zoned")]
created_at: Zoned,
#[serde(with = "crate::time_serde::option_zoned")]
rotated_at: Option<Zoned>,
created_by: Option<String>,
/// Encrypted key material (32 bytes encoded in base64 for AES-256)
encrypted_key_material: String,
/// Nonce used for encryption
nonce: Vec<u8>,
#[serde(default)]
at_rest_protection: StoredKeyProtection,
}
impl LocalKmsClient {
/// Create a new local KMS client
pub async fn new(config: LocalConfig) -> Result<Self> {
// Create key directory if it doesn't exist
if !fs::try_exists(&config.key_dir).await? {
fs::create_dir_all(&config.key_dir).await?;
debug!(path = ?config.key_dir, "KMS key directory created");
}
// Initialize master cipher if master key is provided
let master_cipher = if let Some(ref master_key) = config.master_key {
let salt = Self::load_or_create_master_key_salt(&config).await?;
let key = Self::derive_master_key(master_key, &salt)?;
Some(Aes256Gcm::new(&key))
} else {
warn!("No master key provided - local KMS key material will use explicit plaintext-dev-only storage");
None
};
Ok(Self {
config,
key_cache: RwLock::new(HashMap::new()),
master_cipher,
dek_crypto: AesDekCrypto::new(),
})
}
/// Derive a 256-bit key from the master key string using a persistent Argon2id salt.
fn derive_master_key(master_key: &str, salt: &[u8]) -> Result<Key<Aes256Gcm>> {
let params = Params::new(
LOCAL_KMS_ARGON2_M_COST_KIB,
LOCAL_KMS_ARGON2_T_COST,
LOCAL_KMS_ARGON2_P_COST,
Some(LOCAL_KMS_MASTER_KEY_LEN),
)
.map_err(|err| KmsError::configuration_error(format!("invalid local KMS Argon2 params: {err}")))?;
let argon2 = Argon2::new(Algorithm::Argon2id, Version::V0x13, params);
let mut derived = [0u8; LOCAL_KMS_MASTER_KEY_LEN];
argon2
.hash_password_into(master_key.as_bytes(), salt, &mut derived)
.map_err(|err| KmsError::cryptographic_error("argon2id_kdf", err.to_string()))?;
let key = Key::<Aes256Gcm>::from(derived);
Ok(key)
}
fn master_key_salt_path(config: &LocalConfig) -> PathBuf {
config.key_dir.join(LOCAL_KMS_MASTER_KEY_SALT_FILE)
}
async fn load_or_create_master_key_salt(config: &LocalConfig) -> Result<[u8; LOCAL_KMS_MASTER_KEY_SALT_LEN]> {
let salt_path = Self::master_key_salt_path(config);
if fs::try_exists(&salt_path).await? {
let bytes = fs::read(&salt_path).await?;
return bytes.try_into().map_err(|_| {
KmsError::configuration_error(format!(
"Local KMS master key salt at {} must be exactly {} bytes",
salt_path.display(),
LOCAL_KMS_MASTER_KEY_SALT_LEN
))
});
}
let mut salt = [0u8; LOCAL_KMS_MASTER_KEY_SALT_LEN];
rand::rng().fill(&mut salt[..]);
fs::write(&salt_path, salt).await?;
Self::set_file_permissions(&salt_path, config.file_permissions).await?;
debug!(path = ?salt_path, "Local KMS master key salt created");
Ok(salt)
}
async fn set_file_permissions(path: &std::path::Path, permissions: Option<u32>) -> Result<()> {
#[cfg(unix)]
if let Some(mode) = permissions {
use std::os::unix::fs::PermissionsExt;
let perms = std::fs::Permissions::from_mode(mode);
fs::set_permissions(path, perms).await?;
}
let _ = permissions;
Ok(())
}
/// Get the file path for a master key
fn master_key_path(&self, key_id: &str) -> PathBuf {
self.config.key_dir.join(format!("{key_id}.key"))
}
/// Decode and decrypt a stored key file, returning both the metadata and decrypted key material
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? {
return Err(KmsError::key_not_found(key_id));
}
let content = fs::read(&key_path).await?;
let stored_key: StoredMasterKey = serde_json::from_slice(&content)?;
let encrypted_bytes = BASE64
.decode(&stored_key.encrypted_key_material)
.map_err(|e| KmsError::cryptographic_error("base64_decode", e.to_string()))?;
let effective_protection = if stored_key.at_rest_protection == StoredKeyProtection::LegacyUnspecified {
if stored_key.nonce.is_empty() {
StoredKeyProtection::PlaintextDevOnly
} else {
StoredKeyProtection::EncryptedMasterKey
}
} else {
stored_key.at_rest_protection
};
// Decrypt key material if master cipher is available.
let key_material = match effective_protection {
StoredKeyProtection::EncryptedMasterKey => {
let cipher = self.master_cipher.as_ref().ok_or_else(|| {
KmsError::configuration_error(format!(
"Local KMS key {key_id} is encrypted at rest and requires a configured master key"
))
})?;
if stored_key.nonce.len() != 12 {
return Err(KmsError::cryptographic_error("nonce", "Invalid nonce length"));
}
let mut nonce_array = [0u8; 12];
nonce_array.copy_from_slice(&stored_key.nonce);
let nonce = Nonce::from(nonce_array);
cipher
.decrypt(&nonce, encrypted_bytes.as_ref())
.map_err(|e| KmsError::cryptographic_error("decrypt", e.to_string()))?
}
StoredKeyProtection::PlaintextDevOnly | StoredKeyProtection::LegacyUnspecified => {
if self.master_cipher.is_some() && stored_key.at_rest_protection == StoredKeyProtection::PlaintextDevOnly {
warn!(
key_id,
"Local KMS loaded plaintext-dev-only key material while a master key is configured"
);
}
encrypted_bytes
}
};
Ok((stored_key, key_material))
}
/// Load a master key from disk
async fn load_master_key(&self, key_id: &str) -> Result<MasterKeyInfo> {
let (stored_key, _key_material) = self.decode_stored_key(key_id).await?;
Ok(MasterKeyInfo {
key_id: stored_key.key_id,
version: stored_key.version,
algorithm: stored_key.algorithm,
usage: stored_key.usage,
status: stored_key.status,
description: stored_key.description,
metadata: stored_key.metadata,
created_at: stored_key.created_at,
rotated_at: stored_key.rotated_at,
created_by: stored_key.created_by,
})
}
/// Save a master key to disk
async fn save_master_key(&self, master_key: &MasterKeyInfo, key_material: &[u8]) -> Result<()> {
let key_path = self.master_key_path(&master_key.key_id);
// Encrypt key material if master cipher is available
let (encrypted_key_material, nonce, at_rest_protection) = if let Some(ref cipher) = self.master_cipher {
let mut nonce_bytes = [0u8; 12];
rand::rng().fill(&mut nonce_bytes[..]);
let nonce = Nonce::from(nonce_bytes);
let encrypted = cipher
.encrypt(&nonce, key_material)
.map_err(|e| KmsError::cryptographic_error("encrypt", e.to_string()))?;
// Encode encrypted bytes to base64 string
(BASE64.encode(&encrypted), nonce.to_vec(), StoredKeyProtection::EncryptedMasterKey)
} else {
warn!(
key_id = %master_key.key_id,
"Local KMS is storing key material as plaintext-dev-only because no master key is configured"
);
(BASE64.encode(key_material), Vec::new(), StoredKeyProtection::PlaintextDevOnly)
};
let stored_key = StoredMasterKey {
key_id: master_key.key_id.clone(),
version: master_key.version,
algorithm: master_key.algorithm.clone(),
usage: master_key.usage.clone(),
status: master_key.status.clone(),
description: master_key.description.clone(),
metadata: master_key.metadata.clone(),
created_at: master_key.created_at.clone(),
rotated_at: master_key.rotated_at.clone(),
created_by: master_key.created_by.clone(),
encrypted_key_material,
nonce,
at_rest_protection,
};
let content = serde_json::to_vec_pretty(&stored_key)?;
// Write to temporary file first, then rename for atomicity
let temp_path = key_path.with_extension("tmp");
fs::write(&temp_path, &content).await?;
Self::set_file_permissions(&temp_path, self.config.file_permissions).await?;
fs::rename(&temp_path, &key_path).await?;
debug!(key_id = %master_key.key_id, path = ?key_path, "Local KMS master key saved");
Ok(())
}
/// Get the actual key material for a master key
async fn get_key_material(&self, key_id: &str) -> Result<Vec<u8>> {
let (_stored_key, key_material) = self.decode_stored_key(key_id).await?;
Ok(key_material)
}
/// 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
}
/// 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
}
}
#[async_trait]
impl KmsClient for LocalKmsClient {
async fn generate_data_key(&self, request: &GenerateKeyRequest, _context: Option<&OperationContext>) -> Result<DataKeyInfo> {
debug!("Generating data key for master key: {}", request.master_key_id);
// Generate random data key material
let key_length = match request.key_spec.as_str() {
"AES_256" => 32,
"AES_128" => 16,
_ => return Err(KmsError::unsupported_algorithm(&request.key_spec)),
};
let mut plaintext_key = vec![0u8; key_length];
rand::rng().fill(&mut plaintext_key[..]);
// Encrypt the data key with the master key
let (encrypted_key, nonce) = self.encrypt_with_master_key(&request.master_key_id, &plaintext_key).await?;
// Create data key envelope with master key version for rotation support
let envelope = DataKeyEnvelope {
key_id: uuid::Uuid::new_v4().to_string(),
master_key_id: request.master_key_id.clone(),
key_spec: request.key_spec.clone(),
encrypted_key,
nonce,
encryption_context: request.encryption_context.clone(),
created_at: Zoned::now(),
};
// Serialize the envelope as the ciphertext
let ciphertext = serde_json::to_vec(&envelope)?;
let data_key = DataKeyInfo::new(envelope.key_id, 1, Some(plaintext_key), ciphertext, request.key_spec.clone());
debug!(key_id = %request.master_key_id, "Local KMS data key generated");
Ok(data_key)
}
async fn encrypt(&self, request: &EncryptRequest, context: Option<&OperationContext>) -> Result<EncryptResponse> {
debug!("Encrypting data with key: {}", request.key_id);
// Verify key exists and is active
let key_info = self.describe_key(&request.key_id, context).await?;
if key_info.status != KeyStatus::Active {
return Err(KmsError::invalid_operation(format!(
"Key {} is not active (status: {:?})",
request.key_id, key_info.status
)));
}
let (ciphertext, _nonce) = self.encrypt_with_master_key(&request.key_id, &request.plaintext).await?;
Ok(EncryptResponse {
ciphertext,
key_id: request.key_id.clone(),
key_version: key_info.version,
algorithm: key_info.algorithm,
})
}
async fn decrypt(&self, request: &DecryptRequest, _context: Option<&OperationContext>) -> Result<Vec<u8>> {
debug!("Decrypting data");
// Parse the data key envelope from ciphertext
let envelope: DataKeyEnvelope = serde_json::from_slice(&request.ciphertext)?;
// Verify encryption context matches
// Check that all keys in envelope.encryption_context are present in request.encryption_context
// and their values match. This ensures the context used for decryption matches what was used for encryption.
for (key, expected_value) in &envelope.encryption_context {
if let Some(actual_value) = request.encryption_context.get(key) {
if actual_value != expected_value {
return Err(KmsError::context_mismatch(format!(
"Context mismatch for key '{key}': expected '{expected_value}', got '{actual_value}'"
)));
}
} else {
// If request.encryption_context is empty, allow decryption (backward compatibility)
// Otherwise, require all envelope context keys to be present
if !request.encryption_context.is_empty() {
return Err(KmsError::context_mismatch(format!("Missing context key '{key}'")));
}
}
}
// Decrypt the data key
let plaintext = self
.decrypt_with_master_key(&envelope.master_key_id, &envelope.encrypted_key, &envelope.nonce)
.await?;
debug!("Local KMS data decrypted");
Ok(plaintext)
}
async fn create_key(&self, key_id: &str, algorithm: &str, context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
debug!("Creating master key: {}", key_id);
// Check if key already exists
if self.master_key_path(key_id).exists() {
return Err(KmsError::key_already_exists(key_id));
}
// Validate algorithm
if algorithm != "AES_256" {
return Err(KmsError::unsupported_algorithm(algorithm));
}
// Generate key material
let key_material = generate_key_material(algorithm)?;
let created_by = context
.map(|ctx| ctx.principal.clone())
.unwrap_or_else(|| "local-kms".to_string());
let master_key = MasterKeyInfo::new_with_description(key_id.to_string(), algorithm.to_string(), Some(created_by), None);
// Save to disk
self.save_master_key(&master_key, &key_material).await?;
// Cache the key
let mut cache = self.key_cache.write().await;
cache.insert(key_id.to_string(), master_key.clone());
debug!(key_id, "Local KMS master key created");
Ok(master_key)
}
async fn describe_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<KeyInfo> {
debug!("Describing key: {}", key_id);
// Check cache first
{
let cache = self.key_cache.read().await;
if let Some(master_key) = cache.get(key_id) {
return Ok(master_key.clone().into());
}
}
// Load from disk
let master_key = self.load_master_key(key_id).await?;
// Update cache
{
let mut cache = self.key_cache.write().await;
cache.insert(key_id.to_string(), master_key.clone());
}
Ok(master_key.into())
}
async fn list_keys(&self, request: &ListKeysRequest, _context: Option<&OperationContext>) -> Result<ListKeysResponse> {
debug!("Listing keys");
let mut keys = Vec::new();
let limit = request.limit.unwrap_or(100) as usize;
let mut count = 0;
let mut entries = fs::read_dir(&self.config.key_dir).await?;
while let Some(entry) = entries.next_entry().await? {
if count >= limit {
break;
}
let path = entry.path();
if path.extension().is_some_and(|ext| ext == "key")
&& let Some(stem) = path.file_stem()
&& let Some(key_id) = stem.to_str()
&& let Ok(key_info) = self.describe_key(key_id, None).await
{
// Apply filters
if let Some(ref status_filter) = request.status_filter
&& &key_info.status != status_filter
{
continue;
}
if let Some(ref usage_filter) = request.usage_filter
&& &key_info.usage != usage_filter
{
continue;
}
keys.push(key_info);
count += 1;
}
}
Ok(ListKeysResponse {
keys,
next_marker: None, // Simple implementation without pagination
truncated: false,
})
}
async fn enable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Enabling key: {}", key_id);
let mut master_key = self.load_master_key(key_id).await?;
master_key.status = KeyStatus::Active;
// Preserve the existing key material. Regenerating it on a pure status change would
// destroy the original master key and make every DEK ever wrapped by it permanently
// undecryptable (silent data loss).
let key_material = self.get_key_material(key_id).await?;
self.save_master_key(&master_key, &key_material).await?;
// Update cache
let mut cache = self.key_cache.write().await;
cache.insert(key_id.to_string(), master_key);
debug!(key_id, "Local KMS key enabled");
Ok(())
}
async fn disable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Disabling key: {}", key_id);
let mut master_key = self.load_master_key(key_id).await?;
master_key.status = KeyStatus::Disabled;
// Preserve the existing key material (see enable_key): a status change must never
// regenerate the master key, or every DEK wrapped by it becomes undecryptable.
let key_material = self.get_key_material(key_id).await?;
self.save_master_key(&master_key, &key_material).await?;
// Update cache
let mut cache = self.key_cache.write().await;
cache.insert(key_id.to_string(), master_key);
debug!(key_id, "Local KMS key disabled");
Ok(())
}
async fn schedule_key_deletion(
&self,
key_id: &str,
_pending_window_days: u32,
_context: Option<&OperationContext>,
) -> Result<()> {
debug!("Scheduling deletion for key: {}", key_id);
let mut master_key = self.load_master_key(key_id).await?;
master_key.status = KeyStatus::PendingDeletion;
// Preserve the existing key material (see enable_key): scheduling deletion must not
// regenerate the master key, or cancelling the deletion later would recover a key that
// can no longer decrypt existing data.
let key_material = self.get_key_material(key_id).await?;
self.save_master_key(&master_key, &key_material).await?;
// Update cache
let mut cache = self.key_cache.write().await;
cache.insert(key_id.to_string(), master_key);
debug!(key_id, "Local KMS key deletion scheduled");
Ok(())
}
async fn cancel_key_deletion(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Canceling deletion for key: {}", key_id);
let mut master_key = self.load_master_key(key_id).await?;
master_key.status = KeyStatus::Active;
// Preserve the existing key material (see enable_key): cancelling deletion must recover
// the ORIGINAL key, not mint a new one that cannot decrypt existing data.
let key_material = self.get_key_material(key_id).await?;
self.save_master_key(&master_key, &key_material).await?;
// Update cache
let mut cache = self.key_cache.write().await;
cache.insert(key_id.to_string(), master_key);
debug!(key_id, "Local KMS key deletion canceled");
Ok(())
}
async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
debug!("Rotating key: {}", key_id);
let mut master_key = self.load_master_key(key_id).await?;
master_key.version += 1;
master_key.rotated_at = Some(Zoned::now());
// Generate new key material
let key_material = generate_key_material(&master_key.algorithm)?;
self.save_master_key(&master_key, &key_material).await?;
// Update cache
let mut cache = self.key_cache.write().await;
cache.insert(key_id.to_string(), master_key.clone());
debug!(key_id, "Local KMS key rotated");
Ok(master_key)
}
async fn health_check(&self) -> Result<()> {
// Check if key directory is accessible
if !self.config.key_dir.exists() {
return Err(KmsError::backend_error("Key directory does not exist"));
}
// Try to read the directory
let _ = fs::read_dir(&self.config.key_dir).await?;
Ok(())
}
fn backend_info(&self) -> BackendInfo {
BackendInfo::new(
"local".to_string(),
env!("CARGO_PKG_VERSION").to_string(),
self.config.key_dir.to_string_lossy().to_string(),
true, // We'll assume healthy for now
)
.with_metadata("key_dir".to_string(), self.config.key_dir.to_string_lossy().to_string())
.with_metadata("encrypted_at_rest".to_string(), self.master_cipher.is_some().to_string())
}
}
/// LocalKmsBackend wraps LocalKmsClient and implements the KmsBackend trait
pub struct LocalKmsBackend {
client: LocalKmsClient,
}
impl LocalKmsBackend {
/// Create a new LocalKmsBackend
pub async fn new(config: KmsConfig) -> Result<Self> {
config.validate()?;
let local_config = match &config.backend_config {
crate::config::BackendConfig::Local(local_config) => local_config.clone(),
crate::config::BackendConfig::VaultKv2(_) | crate::config::BackendConfig::VaultTransit(_) => {
return Err(KmsError::configuration_error("Expected Local backend configuration"));
}
};
let client = LocalKmsClient::new(local_config).await?;
Ok(Self { client })
}
}
#[async_trait]
impl KmsBackend for LocalKmsBackend {
async fn create_key(&self, request: CreateKeyRequest) -> Result<CreateKeyResponse> {
let key_id = request.key_name.unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
// Create master key with description directly
let _master_key = {
let algorithm = "AES_256";
// Generate key material
let key_material = generate_key_material(algorithm)?;
let master_key = MasterKeyInfo::new_with_description(
key_id.clone(),
algorithm.to_string(),
Some("local-kms".to_string()),
request.description.clone(),
);
// Save to disk and cache
self.client.save_master_key(&master_key, &key_material).await?;
let mut cache = self.client.key_cache.write().await;
cache.insert(key_id.clone(), master_key.clone());
master_key
};
let metadata = KeyMetadata {
key_id: key_id.clone(),
key_state: KeyState::Enabled,
key_usage: request.key_usage,
description: request.description,
creation_date: Zoned::now(),
deletion_date: None,
origin: "KMS".to_string(),
key_manager: "CUSTOMER".to_string(),
tags: request.tags,
};
Ok(CreateKeyResponse {
key_id,
key_metadata: metadata,
})
}
async fn encrypt(&self, request: EncryptRequest) -> Result<EncryptResponse> {
let encrypt_request = EncryptRequest {
key_id: request.key_id.clone(),
plaintext: request.plaintext,
encryption_context: request.encryption_context,
grant_tokens: request.grant_tokens,
};
let response = self.client.encrypt(&encrypt_request, None).await?;
Ok(EncryptResponse {
ciphertext: response.ciphertext,
key_id: response.key_id,
key_version: response.key_version,
algorithm: response.algorithm,
})
}
async fn decrypt(&self, request: DecryptRequest) -> Result<DecryptResponse> {
let plaintext = self.client.decrypt(&request, None).await?;
// For simplicity, return basic response - in real implementation would extract more info from ciphertext
Ok(DecryptResponse {
plaintext,
key_id: "unknown".to_string(), // Would be extracted from ciphertext metadata
encryption_algorithm: Some("AES-256-GCM".to_string()),
})
}
async fn generate_data_key(&self, request: GenerateDataKeyRequest) -> Result<GenerateDataKeyResponse> {
let generate_request = GenerateKeyRequest {
master_key_id: request.key_id.clone(),
key_spec: request.key_spec.as_str().to_string(),
key_length: Some(request.key_spec.key_size() as u32),
encryption_context: request.encryption_context,
grant_tokens: Vec::new(),
};
let data_key = self.client.generate_data_key(&generate_request, None).await?;
Ok(GenerateDataKeyResponse {
key_id: request.key_id,
plaintext_key: data_key.plaintext.clone().unwrap_or_default(),
ciphertext_blob: data_key.ciphertext.clone(),
})
}
async fn describe_key(&self, request: DescribeKeyRequest) -> Result<DescribeKeyResponse> {
let key_info = self.client.describe_key(&request.key_id, None).await?;
let metadata = KeyMetadata {
key_id: key_info.key_id,
key_state: match key_info.status {
KeyStatus::Active => KeyState::Enabled,
KeyStatus::Disabled => KeyState::Disabled,
KeyStatus::PendingDeletion => KeyState::PendingDeletion,
KeyStatus::Deleted => KeyState::Unavailable,
},
key_usage: key_info.usage,
description: key_info.description,
creation_date: key_info.created_at,
deletion_date: None,
origin: "KMS".to_string(),
key_manager: "CUSTOMER".to_string(),
tags: key_info.tags,
};
Ok(DescribeKeyResponse { key_metadata: metadata })
}
async fn list_keys(&self, request: ListKeysRequest) -> Result<ListKeysResponse> {
let response = self.client.list_keys(&request, None).await?;
Ok(response)
}
async fn delete_key(&self, request: DeleteKeyRequest) -> Result<DeleteKeyResponse> {
// For local backend, we'll implement immediate deletion by default
// unless a pending window is specified
let key_id = &request.key_id;
// First, load the key from disk to get the master key
let mut master_key = self
.client
.load_master_key(key_id)
.await
.map_err(|_| KmsError::key_not_found(format!("Key {key_id} not found")))?;
let (deletion_date_str, deletion_date_dt) = if request.force_immediate.unwrap_or(false) {
// For immediate deletion, actually delete the key from filesystem
let key_path = self.client.master_key_path(key_id);
tokio::fs::remove_file(&key_path)
.await
.map_err(|e| KmsError::internal_error(format!("Failed to delete key file: {e}")))?;
// Remove from cache
let mut cache = self.client.key_cache.write().await;
cache.remove(key_id);
debug!(key_id, "Local KMS key deleted immediately");
// Return success response for immediate deletion
let key_metadata = KeyMetadata {
key_id: master_key.key_id.clone(),
description: master_key.description.clone(),
key_usage: master_key.usage,
key_state: KeyState::PendingDeletion, // AWS KMS compatibility
creation_date: master_key.created_at,
deletion_date: Some(Zoned::now()),
key_manager: "CUSTOMER".to_string(),
origin: "AWS_KMS".to_string(),
tags: master_key.metadata,
};
return Ok(DeleteKeyResponse {
key_id: key_id.clone(),
deletion_date: None, // No deletion date for immediate deletion
key_metadata,
});
} else {
// Schedule for deletion (default 30 days)
let days = request.pending_window_in_days.unwrap_or(30);
if !(7..=30).contains(&days) {
return Err(KmsError::invalid_parameter("pending_window_in_days must be between 7 and 30".to_string()));
}
let deletion_date = Zoned::now() + Duration::from_secs(days as u64 * 86400);
master_key.status = KeyStatus::PendingDeletion;
(Some(deletion_date.to_string()), Some(deletion_date))
};
// Save the updated key to disk - preserve existing key material!
// Load and decode the stored key to get the existing key material
let (_stored_key, existing_key_material) = self
.client
.decode_stored_key(key_id)
.await
.map_err(|e| KmsError::internal_error(format!("Failed to decode key: {e}")))?;
self.client.save_master_key(&master_key, &existing_key_material).await?;
// Update cache
let mut cache = self.client.key_cache.write().await;
cache.insert(key_id.to_string(), master_key.clone());
// Convert master_key to KeyMetadata for response
let key_metadata = KeyMetadata {
key_id: master_key.key_id.clone(),
description: master_key.description.clone(),
key_usage: master_key.usage,
key_state: KeyState::PendingDeletion,
creation_date: master_key.created_at,
deletion_date: deletion_date_dt,
key_manager: "CUSTOMER".to_string(),
origin: "AWS_KMS".to_string(),
tags: master_key.metadata,
};
Ok(DeleteKeyResponse {
key_id: key_id.clone(),
deletion_date: deletion_date_str,
key_metadata,
})
}
async fn cancel_key_deletion(&self, request: CancelKeyDeletionRequest) -> Result<CancelKeyDeletionResponse> {
let key_id = &request.key_id;
// Load the key from disk to get the master key
let mut master_key = self
.client
.load_master_key(key_id)
.await
.map_err(|_| KmsError::key_not_found(format!("Key {key_id} not found")))?;
if master_key.status != KeyStatus::PendingDeletion {
return Err(KmsError::invalid_key_state(format!("Key {key_id} is not pending deletion")));
}
// Cancel the deletion by resetting the state
master_key.status = KeyStatus::Active;
// Save the updated key to disk - this is the missing critical step!
// Preserve existing key material instead of generating new one
let (_stored_key, existing_key_material) = self
.client
.decode_stored_key(key_id)
.await
.map_err(|e| KmsError::internal_error(format!("Failed to decode key: {e}")))?;
self.client.save_master_key(&master_key, &existing_key_material).await?;
// Update cache
let mut cache = self.client.key_cache.write().await;
cache.insert(key_id.to_string(), master_key.clone());
// Convert master_key to KeyMetadata for response
let key_metadata = KeyMetadata {
key_id: master_key.key_id.clone(),
description: master_key.description.clone(),
key_usage: master_key.usage,
key_state: KeyState::Enabled,
creation_date: master_key.created_at,
deletion_date: None,
key_manager: "CUSTOMER".to_string(),
origin: "AWS_KMS".to_string(),
tags: master_key.metadata,
};
Ok(CancelKeyDeletionResponse {
key_id: key_id.clone(),
key_metadata,
})
}
async fn health_check(&self) -> Result<bool> {
self.client.health_check().await.map(|_| true)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
use tempfile::TempDir;
async fn create_test_client() -> (LocalKmsClient, TempDir) {
let temp_dir = TempDir::new().expect("Failed to create temp dir");
let config = LocalConfig {
key_dir: temp_dir.path().to_path_buf(),
master_key: Some("test-master-key".to_string()),
file_permissions: Some(0o600),
};
let client = LocalKmsClient::new(config).await.expect("Failed to create client");
(client, temp_dir)
}
async fn create_dev_mode_client() -> (LocalKmsClient, TempDir) {
let temp_dir = TempDir::new().expect("Failed to create temp dir");
let config = LocalConfig {
key_dir: temp_dir.path().to_path_buf(),
master_key: None,
file_permissions: Some(0o600),
};
let client = LocalKmsClient::new(config).await.expect("Failed to create dev-mode client");
(client, temp_dir)
}
#[tokio::test]
async fn test_key_lifecycle() {
let (client, _temp_dir) = create_test_client().await;
let key_id = "test-key";
let algorithm = "AES_256";
// Create key
let master_key = client
.create_key(key_id, algorithm, None)
.await
.expect("Failed to create key");
assert_eq!(master_key.key_id, key_id);
assert_eq!(master_key.algorithm, algorithm);
assert_eq!(master_key.status, KeyStatus::Active);
// Describe key
let key_info = client.describe_key(key_id, None).await.expect("Failed to describe key");
assert_eq!(key_info.key_id, key_id);
assert_eq!(key_info.status, KeyStatus::Active);
// List keys
let list_response = client
.list_keys(&ListKeysRequest::default(), None)
.await
.expect("Failed to list keys");
assert_eq!(list_response.keys.len(), 1);
assert_eq!(list_response.keys[0].key_id, key_id);
// Disable key
client.disable_key(key_id, None).await.expect("Failed to disable key");
let key_info = client.describe_key(key_id, None).await.expect("Failed to describe key");
assert_eq!(key_info.status, KeyStatus::Disabled);
// Enable key
client.enable_key(key_id, None).await.expect("Failed to enable key");
let key_info = client.describe_key(key_id, None).await.expect("Failed to describe key");
assert_eq!(key_info.status, KeyStatus::Active);
}
#[tokio::test]
async fn test_data_key_operations() {
let (client, _temp_dir) = create_test_client().await;
let key_id = "test-key";
client
.create_key(key_id, "AES_256", None)
.await
.expect("Failed to create key");
// Generate data key
let request = GenerateKeyRequest::new(key_id.to_string(), "AES_256".to_string())
.with_context("bucket".to_string(), "test-bucket".to_string());
let data_key = client
.generate_data_key(&request, None)
.await
.expect("Failed to generate data key");
assert!(data_key.plaintext.is_some());
assert!(!data_key.ciphertext.is_empty());
// Decrypt data key
let decrypt_request =
DecryptRequest::new(data_key.ciphertext.clone()).with_context("bucket".to_string(), "test-bucket".to_string());
let decrypted = client.decrypt(&decrypt_request, None).await.expect("Failed to decrypt");
assert_eq!(decrypted, data_key.plaintext.clone().expect("No plaintext"));
}
#[tokio::test]
async fn key_state_transitions_preserve_master_key_material() {
// Regression: enable/disable/schedule_deletion/cancel_deletion previously regenerated the
// master key material on a pure status change, permanently destroying the ability to
// decrypt any DEK wrapped by that key. A status cycle must preserve the material.
let (client, _temp_dir) = create_test_client().await;
let key_id = "state-cycle-key";
client.create_key(key_id, "AES_256", None).await.expect("create");
let request = GenerateKeyRequest::new(key_id.to_string(), "AES_256".to_string())
.with_context("bucket".to_string(), "b".to_string());
let data_key = client.generate_data_key(&request, None).await.expect("generate data key");
let ciphertext = data_key.ciphertext.clone();
let plaintext = data_key.plaintext.clone().expect("no plaintext");
// Cycle through every status-changing method the fix touches.
client.disable_key(key_id, None).await.expect("disable");
client.enable_key(key_id, None).await.expect("enable");
client
.schedule_key_deletion(key_id, 7, None)
.await
.expect("schedule deletion");
client.cancel_key_deletion(key_id, None).await.expect("cancel deletion");
// Pre-fix, each of those regenerated the master key, so this unwrap fails with an AEAD
// error. Post-fix, the original material is preserved and the DEK still decrypts.
let decrypt_request = DecryptRequest::new(ciphertext).with_context("bucket".to_string(), "b".to_string());
let decrypted = client
.decrypt(&decrypt_request, None)
.await
.expect("DEK must still decrypt after status transitions");
assert_eq!(decrypted, plaintext, "master key material must survive status transitions");
}
#[tokio::test]
async fn test_encryption_operations() {
let (client, _temp_dir) = create_test_client().await;
let key_id = "test-key";
client
.create_key(key_id, "AES_256", None)
.await
.expect("Failed to create key");
let plaintext = b"Hello, World!";
let encrypt_request = EncryptRequest::new(key_id.to_string(), plaintext.to_vec());
// Encrypt
let encrypt_response = client.encrypt(&encrypt_request, None).await.expect("Failed to encrypt");
assert!(!encrypt_response.ciphertext.is_empty());
assert_eq!(encrypt_response.key_id, key_id);
// Note: Direct decryption of encrypt() results is not implemented in this simple version
// In a real implementation, encrypt() would create a different envelope format
}
#[tokio::test]
async fn test_encrypted_master_key_storage_uses_explicit_protection_and_salt() {
let (client, _temp_dir) = create_test_client().await;
client
.create_key("encrypted-key", "AES_256", None)
.await
.expect("Failed to create encrypted key");
let salt = fs::read(LocalKmsClient::master_key_salt_path(&client.config))
.await
.expect("master key salt should exist");
assert_eq!(salt.len(), LOCAL_KMS_MASTER_KEY_SALT_LEN);
let stored: StoredMasterKey = serde_json::from_slice(
&fs::read(client.master_key_path("encrypted-key"))
.await
.expect("stored key should exist"),
)
.expect("stored encrypted key should deserialize");
assert_eq!(stored.at_rest_protection, StoredKeyProtection::EncryptedMasterKey);
assert_eq!(stored.nonce.len(), 12);
}
#[tokio::test]
async fn test_plaintext_dev_only_storage_is_explicit_and_loadable() {
let (client, _temp_dir) = create_dev_mode_client().await;
client
.create_key("plaintext-key", "AES_256", None)
.await
.expect("Failed to create plaintext-dev-only key");
let stored: StoredMasterKey = serde_json::from_slice(
&fs::read(client.master_key_path("plaintext-key"))
.await
.expect("stored key should exist"),
)
.expect("stored plaintext key should deserialize");
assert_eq!(stored.at_rest_protection, StoredKeyProtection::PlaintextDevOnly);
assert!(stored.nonce.is_empty(), "plaintext-dev-only keys should not store a nonce");
let key_info = client
.describe_key("plaintext-key", None)
.await
.expect("plaintext-dev-only key should remain readable");
assert_eq!(key_info.key_id, "plaintext-key");
}
#[tokio::test]
async fn test_encrypted_key_requires_master_key_to_load() {
let (client, temp_dir) = create_test_client().await;
client
.create_key("encrypted-key", "AES_256", None)
.await
.expect("Failed to create encrypted key");
let config = LocalConfig {
key_dir: temp_dir.path().to_path_buf(),
master_key: None,
file_permissions: Some(0o600),
};
let client_without_master = LocalKmsClient::new(config)
.await
.expect("client without master key should still initialize in dev-mode tests");
let err = client_without_master
.describe_key("encrypted-key", None)
.await
.expect_err("encrypted key should require a master key to read");
assert!(err.to_string().contains("requires a configured master key"));
}
#[tokio::test]
async fn test_load_master_key_accepts_legacy_rfc3339_timestamp() {
let (client, _temp_dir) = create_dev_mode_client().await;
let stored_key = serde_json::json!({
"key_id": "legacy-key",
"version": 1u32,
"algorithm": "AES_256",
"usage": "EncryptDecrypt",
"status": "Active",
"description": serde_json::Value::Null,
"metadata": HashMap::<String, String>::new(),
"created_at": "2024-01-01T00:00:00+00:00",
"rotated_at": serde_json::Value::Null,
"created_by": "legacy-test",
"encrypted_key_material": BASE64.encode([7u8; 32]),
"nonce": Vec::<u8>::new()
});
let key_path = client.master_key_path("legacy-key");
fs::write(&key_path, serde_json::to_vec_pretty(&stored_key).expect("serialize test key"))
.await
.expect("write legacy key");
let key_info = client.load_master_key("legacy-key").await.expect("legacy key should load");
assert_eq!(key_info.key_id, "legacy-key");
assert_eq!(key_info.created_at.time_zone().iana_name(), Some("UTC"));
}
#[tokio::test]
async fn test_load_master_key_accepts_legacy_encrypted_record_without_protection_field() {
let (client, temp_dir) = create_test_client().await;
client
.create_key("legacy-encrypted-key", "AES_256", None)
.await
.expect("Failed to create encrypted key");
let key_path = client.master_key_path("legacy-encrypted-key");
let mut stored_json: serde_json::Value =
serde_json::from_slice(&fs::read(&key_path).await.expect("stored key should exist"))
.expect("stored key should deserialize");
stored_json
.as_object_mut()
.expect("stored key should be an object")
.remove("at_rest_protection");
fs::write(
&key_path,
serde_json::to_vec_pretty(&stored_json).expect("legacy record should serialize"),
)
.await
.expect("legacy record should be writable");
let legacy_client = LocalKmsClient::new(LocalConfig {
key_dir: temp_dir.path().to_path_buf(),
master_key: Some("test-master-key".to_string()),
file_permissions: Some(0o600),
})
.await
.expect("legacy client should initialize");
let key_info = legacy_client
.describe_key("legacy-encrypted-key", None)
.await
.expect("legacy encrypted record should remain readable");
assert_eq!(key_info.key_id, "legacy-encrypted-key");
}
}