mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-13 16:46:55 +00:00
refactor(sse): decouple ecstore and harden KMS lifecycle (#5435)
* refactor(sse): decouple encryption from ecstore * feat(kms): enhance KMS service manager with runtime state and persistence support * feat(kms): add local key export functionality for SSE-S3 migration tests * fix(kms): keep local key export narrowly scoped * fix(sse): validate copy source customer algorithm --------- Co-authored-by: Zhengchao An <anzhengchao@gmail.com>
This commit is contained in:
@@ -36,6 +36,7 @@ use std::path::{Component, Path, PathBuf};
|
||||
use std::time::Duration;
|
||||
use tokio::fs;
|
||||
use tracing::{debug, warn};
|
||||
use zeroize::Zeroizing;
|
||||
|
||||
/// Reject key identifiers that would not name a single file directly inside the key
|
||||
/// directory.
|
||||
@@ -146,6 +147,45 @@ impl LocalKmsClient {
|
||||
Ok(client)
|
||||
}
|
||||
|
||||
/// Open a Local KMS key directory without creating or modifying any files.
|
||||
///
|
||||
/// This constructor is restricted to explicit key-export tooling. Normal
|
||||
/// backend operation must use [`Self::new`].
|
||||
pub async fn new_for_key_export(config: LocalConfig) -> Result<Self> {
|
||||
if !fs::try_exists(&config.key_dir).await? {
|
||||
return Err(KmsError::configuration_error("Local KMS key directory does not exist"));
|
||||
}
|
||||
|
||||
let (master_cipher, legacy_master_cipher) = if let Some(ref master_key) = config.master_key {
|
||||
let legacy_key = Self::derive_legacy_master_key(master_key)?;
|
||||
let legacy_master_cipher = Aes256Gcm::new(&legacy_key);
|
||||
let salt_path = Self::master_key_salt_path(&config);
|
||||
let master_cipher = if fs::try_exists(&salt_path).await? {
|
||||
let salt = fs::read(&salt_path).await?;
|
||||
let salt: [u8; LOCAL_KMS_MASTER_KEY_SALT_LEN] = salt.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
|
||||
))
|
||||
})?;
|
||||
Aes256Gcm::new(&Self::derive_master_key(master_key, &salt)?)
|
||||
} else {
|
||||
Aes256Gcm::new(&legacy_key)
|
||||
};
|
||||
(Some(master_cipher), Some(legacy_master_cipher))
|
||||
} else {
|
||||
(None, None)
|
||||
};
|
||||
|
||||
Ok(Self {
|
||||
config,
|
||||
master_cipher,
|
||||
legacy_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(
|
||||
@@ -435,6 +475,20 @@ impl LocalKmsClient {
|
||||
Ok(key_material)
|
||||
}
|
||||
|
||||
/// Decrypt an AES-256 Local KMS key for explicit migration tooling.
|
||||
///
|
||||
/// The returned buffer is zeroized on drop. Callers must treat the value as
|
||||
/// plaintext key material and avoid logging or persisting it.
|
||||
pub async fn decrypt_key_material_for_export(&self, key_id: &str) -> Result<Zeroizing<[u8; 32]>> {
|
||||
let (stored_key, key_material) = self.decode_stored_key(key_id).await?;
|
||||
if stored_key.algorithm != "AES_256" {
|
||||
return Err(KmsError::unsupported_algorithm(stored_key.algorithm));
|
||||
}
|
||||
let actual = key_material.len();
|
||||
let key_material = key_material.try_into().map_err(|_| KmsError::invalid_key_size(32, actual))?;
|
||||
Ok(Zeroizing::new(key_material))
|
||||
}
|
||||
|
||||
async fn validate_existing_keys(&self) -> Result<()> {
|
||||
let mut entries = fs::read_dir(&self.config.key_dir).await?;
|
||||
while let Some(entry) = entries.next_entry().await? {
|
||||
@@ -1208,6 +1262,52 @@ mod tests {
|
||||
assert!(matches!(wrong_master_error, KmsError::CryptographicError { .. }));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn key_export_uses_existing_local_decryption_path_without_writing_files() {
|
||||
let (client, _temp_dir) = create_test_client().await;
|
||||
let key_id = "export-key";
|
||||
client
|
||||
.create_key(key_id, "AES_256", None)
|
||||
.await
|
||||
.expect("create encrypted key");
|
||||
let expected = client.get_key_material(key_id).await.expect("load expected key material");
|
||||
let salt_path = LocalKmsClient::master_key_salt_path(&client.config);
|
||||
let salt_before = fs::read(&salt_path).await.expect("read existing salt");
|
||||
|
||||
let export_client = LocalKmsClient::new_for_key_export(client.config.clone())
|
||||
.await
|
||||
.expect("open read-only export client");
|
||||
let exported = export_client
|
||||
.decrypt_key_material_for_export(key_id)
|
||||
.await
|
||||
.expect("decrypt key for export");
|
||||
|
||||
assert_eq!(exported.as_ref(), expected.as_slice());
|
||||
assert_eq!(fs::read(&salt_path).await.expect("read unchanged salt"), salt_before);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn key_export_accepts_plaintext_dev_only_key_without_master_key() {
|
||||
let (client, _temp_dir) = create_dev_mode_client().await;
|
||||
let key_id = "plaintext-export-key";
|
||||
client
|
||||
.create_key(key_id, "AES_256", None)
|
||||
.await
|
||||
.expect("create plaintext-dev-only key");
|
||||
let expected = client.get_key_material(key_id).await.expect("load expected key material");
|
||||
|
||||
let export_client = LocalKmsClient::new_for_key_export(client.config.clone())
|
||||
.await
|
||||
.expect("open read-only export client");
|
||||
let exported = export_client
|
||||
.decrypt_key_material_for_export(key_id)
|
||||
.await
|
||||
.expect("export plaintext-dev-only key");
|
||||
|
||||
assert_eq!(exported.as_ref(), expected.as_slice());
|
||||
assert!(!LocalKmsClient::master_key_salt_path(&client.config).exists());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_plaintext_dev_only_storage_is_explicit_and_loadable() {
|
||||
let (client, _temp_dir) = create_dev_mode_client().await;
|
||||
|
||||
@@ -89,7 +89,7 @@ pub use error::{KmsError, KmsUnavailableError, Result};
|
||||
pub use manager::KmsManager;
|
||||
pub use service::{DataKey, ObjectEncryptionService};
|
||||
pub use service_manager::{
|
||||
KmsServiceManager, KmsServiceStatus, get_global_encryption_service, get_global_kms_service_manager,
|
||||
KmsServiceManager, KmsServiceStatus, KmsStartOutcome, get_global_encryption_service, get_global_kms_service_manager,
|
||||
init_global_kms_service_manager,
|
||||
};
|
||||
pub use types::*;
|
||||
|
||||
@@ -21,12 +21,13 @@ use crate::manager::KmsManager;
|
||||
use crate::service::ObjectEncryptionService;
|
||||
use arc_swap::ArcSwap;
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::future::Future;
|
||||
use std::sync::{
|
||||
Arc, OnceLock,
|
||||
atomic::{AtomicU64, Ordering},
|
||||
};
|
||||
use subtle::ConstantTimeEq;
|
||||
use tokio::sync::{Mutex, RwLock};
|
||||
use tokio::sync::Mutex;
|
||||
use tracing::{debug, error, info, warn};
|
||||
|
||||
const LOG_COMPONENT_KMS: &str = "kms";
|
||||
@@ -93,6 +94,13 @@ pub enum KmsServiceStatus {
|
||||
Error(String),
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum KmsStartOutcome {
|
||||
Started,
|
||||
Restarted,
|
||||
AlreadyRunning,
|
||||
}
|
||||
|
||||
/// Service version information for zero-downtime reconfiguration
|
||||
#[derive(Clone)]
|
||||
struct ServiceVersion {
|
||||
@@ -104,16 +112,17 @@ struct ServiceVersion {
|
||||
manager: Arc<KmsManager>,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct RuntimeState {
|
||||
config: Option<KmsConfig>,
|
||||
status: KmsServiceStatus,
|
||||
current_service: Option<ServiceVersion>,
|
||||
}
|
||||
|
||||
/// Dynamic KMS service manager with versioned services for zero-downtime reconfiguration
|
||||
pub struct KmsServiceManager {
|
||||
/// Current service version (if running)
|
||||
/// Uses ArcSwap for atomic, lock-free service switching
|
||||
/// This allows instant atomic updates without blocking readers
|
||||
current_service: ArcSwap<Option<ServiceVersion>>,
|
||||
/// Current configuration
|
||||
config: Arc<RwLock<Option<KmsConfig>>>,
|
||||
/// Current status
|
||||
status: Arc<RwLock<KmsServiceStatus>>,
|
||||
/// Atomically published configuration, status, and current service.
|
||||
state: ArcSwap<RuntimeState>,
|
||||
/// Version counter (monotonically increasing)
|
||||
version_counter: Arc<AtomicU64>,
|
||||
/// Mutex to protect lifecycle operations (start, stop, reconfigure)
|
||||
@@ -125,9 +134,11 @@ impl KmsServiceManager {
|
||||
/// Create a new KMS service manager (not configured)
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
current_service: ArcSwap::from_pointee(None),
|
||||
config: Arc::new(RwLock::new(None)),
|
||||
status: Arc::new(RwLock::new(KmsServiceStatus::NotConfigured)),
|
||||
state: ArcSwap::from_pointee(RuntimeState {
|
||||
config: None,
|
||||
status: KmsServiceStatus::NotConfigured,
|
||||
current_service: None,
|
||||
}),
|
||||
version_counter: Arc::new(AtomicU64::new(0)),
|
||||
lifecycle_mutex: Arc::new(Mutex::new(())),
|
||||
}
|
||||
@@ -135,44 +146,67 @@ impl KmsServiceManager {
|
||||
|
||||
/// Get current service status
|
||||
pub async fn get_status(&self) -> KmsServiceStatus {
|
||||
self.status.read().await.clone()
|
||||
self.state.load().status.clone()
|
||||
}
|
||||
|
||||
/// Get current configuration (if any)
|
||||
pub async fn get_config(&self) -> Option<KmsConfig> {
|
||||
self.config.read().await.clone()
|
||||
self.state.load().config.clone()
|
||||
}
|
||||
|
||||
/// Get configuration for status and management responses without static key material.
|
||||
pub async fn get_redacted_config(&self) -> Option<KmsConfig> {
|
||||
let mut config = self.config.read().await.clone()?;
|
||||
let mut config = self.state.load().config.clone()?;
|
||||
Self::redact_config(&mut config);
|
||||
Some(config)
|
||||
}
|
||||
|
||||
/// Get status and redacted configuration from the same published snapshot.
|
||||
pub async fn get_redacted_state(&self) -> (KmsServiceStatus, Option<KmsConfig>) {
|
||||
let state = self.state.load();
|
||||
let mut config = state.config.clone();
|
||||
if let Some(config) = &mut config {
|
||||
Self::redact_config(config);
|
||||
}
|
||||
(state.status.clone(), config)
|
||||
}
|
||||
|
||||
fn redact_config(config: &mut KmsConfig) {
|
||||
if let BackendConfig::Static(static_config) = &mut config.backend_config {
|
||||
use zeroize::Zeroize;
|
||||
static_config.secret_key.zeroize();
|
||||
}
|
||||
Some(config)
|
||||
}
|
||||
|
||||
/// Configure KMS with new configuration
|
||||
pub async fn configure(&self, new_config: KmsConfig) -> Result<()> {
|
||||
let _guard = self.lifecycle_mutex.lock().await;
|
||||
self.configure_with_persistence(new_config, || async { Ok(()) }).await
|
||||
}
|
||||
|
||||
/// Configure KMS and publish the in-memory state only after persistence succeeds.
|
||||
///
|
||||
/// The persistence callback runs under the lifecycle lock and must not call
|
||||
/// another lifecycle method on this manager.
|
||||
pub async fn configure_with_persistence<Persist, PersistFuture>(&self, new_config: KmsConfig, persist: Persist) -> Result<()>
|
||||
where
|
||||
Persist: FnOnce() -> PersistFuture,
|
||||
PersistFuture: Future<Output = Result<()>>,
|
||||
{
|
||||
new_config.validate()?;
|
||||
{
|
||||
let config = self.config.read().await;
|
||||
validate_local_transition(config.as_ref(), &new_config)?;
|
||||
}
|
||||
|
||||
// Update configuration
|
||||
{
|
||||
let mut config = self.config.write().await;
|
||||
*config = Some(new_config.clone());
|
||||
}
|
||||
|
||||
// Update status
|
||||
{
|
||||
let mut status = self.status.write().await;
|
||||
*status = KmsServiceStatus::Configured;
|
||||
let _guard = self.lifecycle_mutex.lock().await;
|
||||
let current = self.state.load_full();
|
||||
validate_local_transition(current.config.as_ref(), &new_config)?;
|
||||
if current.current_service.is_some() {
|
||||
return Err(KmsError::configuration_error(
|
||||
"Cannot configure KMS while it is running; use reconfigure instead",
|
||||
));
|
||||
}
|
||||
persist().await?;
|
||||
self.state.store(Arc::new(RuntimeState {
|
||||
config: Some(new_config),
|
||||
status: KmsServiceStatus::Configured,
|
||||
current_service: None,
|
||||
}));
|
||||
|
||||
debug!(
|
||||
event = EVENT_KMS_SERVICE_STATE,
|
||||
@@ -190,19 +224,35 @@ impl KmsServiceManager {
|
||||
self.start_internal().await
|
||||
}
|
||||
|
||||
/// Start or restart KMS with the running-state decision serialized with the lifecycle action.
|
||||
pub async fn start_or_restart(&self, force: bool) -> Result<KmsStartOutcome> {
|
||||
let _guard = self.lifecycle_mutex.lock().await;
|
||||
let running = self.state.load().current_service.is_some();
|
||||
if running && !force {
|
||||
return Ok(KmsStartOutcome::AlreadyRunning);
|
||||
}
|
||||
self.start_internal().await?;
|
||||
Ok(if running {
|
||||
KmsStartOutcome::Restarted
|
||||
} else {
|
||||
KmsStartOutcome::Started
|
||||
})
|
||||
}
|
||||
|
||||
/// Internal start implementation (called within lifecycle mutex)
|
||||
async fn start_internal(&self) -> Result<()> {
|
||||
let config = {
|
||||
let config_guard = self.config.read().await;
|
||||
match config_guard.as_ref() {
|
||||
Some(config) => config.clone(),
|
||||
None => {
|
||||
let err_msg = "Cannot start KMS: no configuration provided";
|
||||
error!("{}", err_msg);
|
||||
let mut status = self.status.write().await;
|
||||
*status = KmsServiceStatus::Error(err_msg.to_string());
|
||||
return Err(KmsError::configuration_error(err_msg));
|
||||
}
|
||||
let state = self.state.load_full();
|
||||
let config = match state.config.as_ref() {
|
||||
Some(config) => config.clone(),
|
||||
None => {
|
||||
let err_msg = "Cannot start KMS: no configuration provided";
|
||||
error!("{}", err_msg);
|
||||
self.state.store(Arc::new(RuntimeState {
|
||||
config: None,
|
||||
status: KmsServiceStatus::Error(err_msg.to_string()),
|
||||
current_service: None,
|
||||
}));
|
||||
return Err(KmsError::configuration_error(err_msg));
|
||||
}
|
||||
};
|
||||
|
||||
@@ -215,17 +265,9 @@ impl KmsServiceManager {
|
||||
"KMS service starting"
|
||||
);
|
||||
|
||||
match self.create_service_version(&config).await {
|
||||
match self.create_healthy_service_version(&config).await {
|
||||
Ok(service_version) => {
|
||||
// Atomically update to new service version (lock-free, instant)
|
||||
// ArcSwap::store() is a true atomic operation using CAS
|
||||
self.current_service.store(Arc::new(Some(service_version)));
|
||||
|
||||
// Update status
|
||||
{
|
||||
let mut status = self.status.write().await;
|
||||
*status = KmsServiceStatus::Running;
|
||||
}
|
||||
self.publish_running(config, service_version);
|
||||
|
||||
debug!(
|
||||
event = EVENT_KMS_SERVICE_STATE,
|
||||
@@ -239,13 +281,24 @@ impl KmsServiceManager {
|
||||
Err(e) => {
|
||||
let err_msg = format!("Failed to create KMS backend: {e}");
|
||||
error!("{}", err_msg);
|
||||
let mut status = self.status.write().await;
|
||||
*status = KmsServiceStatus::Error(err_msg.clone());
|
||||
if state.current_service.is_none() {
|
||||
self.state.store(Arc::new(RuntimeState {
|
||||
config: state.config.clone(),
|
||||
status: KmsServiceStatus::Error(err_msg.clone()),
|
||||
current_service: None,
|
||||
}));
|
||||
}
|
||||
Err(KmsError::backend_error(&err_msg))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Replace the running service without exposing a stopped interval.
|
||||
pub async fn restart(&self) -> Result<()> {
|
||||
let _guard = self.lifecycle_mutex.lock().await;
|
||||
self.start_internal().await
|
||||
}
|
||||
|
||||
/// Stop KMS service
|
||||
///
|
||||
/// Note: This stops accepting new operations, but existing operations using
|
||||
@@ -267,15 +320,16 @@ impl KmsServiceManager {
|
||||
|
||||
// Atomically clear current service version (lock-free, instant)
|
||||
// Note: Existing Arc references will keep the service alive until operations complete
|
||||
self.current_service.store(Arc::new(None));
|
||||
|
||||
// Update status (keep configuration)
|
||||
{
|
||||
let mut status = self.status.write().await;
|
||||
if !matches!(*status, KmsServiceStatus::NotConfigured) {
|
||||
*status = KmsServiceStatus::Configured;
|
||||
}
|
||||
}
|
||||
let state = self.state.load_full();
|
||||
self.state.store(Arc::new(RuntimeState {
|
||||
config: state.config.clone(),
|
||||
status: if state.config.is_some() {
|
||||
KmsServiceStatus::Configured
|
||||
} else {
|
||||
KmsServiceStatus::NotConfigured
|
||||
},
|
||||
current_service: None,
|
||||
}));
|
||||
|
||||
debug!(
|
||||
event = EVENT_KMS_SERVICE_STATE,
|
||||
@@ -298,6 +352,22 @@ impl KmsServiceManager {
|
||||
/// This ensures zero downtime during reconfiguration, even for long-running
|
||||
/// operations like encrypting large files.
|
||||
pub async fn reconfigure(&self, new_config: KmsConfig) -> Result<()> {
|
||||
self.reconfigure_with_persistence(new_config, || async { Ok(()) }).await
|
||||
}
|
||||
|
||||
/// Reconfigure KMS after the candidate is healthy and persistence succeeds.
|
||||
///
|
||||
/// The persistence callback runs under the lifecycle lock and must not call
|
||||
/// another lifecycle method on this manager.
|
||||
pub async fn reconfigure_with_persistence<Persist, PersistFuture>(
|
||||
&self,
|
||||
new_config: KmsConfig,
|
||||
persist: Persist,
|
||||
) -> Result<()>
|
||||
where
|
||||
Persist: FnOnce() -> PersistFuture,
|
||||
PersistFuture: Future<Output = Result<()>>,
|
||||
{
|
||||
let _guard = self.lifecycle_mutex.lock().await;
|
||||
|
||||
debug!(
|
||||
@@ -308,33 +378,18 @@ impl KmsServiceManager {
|
||||
"KMS service reconfiguring"
|
||||
);
|
||||
new_config.validate()?;
|
||||
{
|
||||
let config = self.config.read().await;
|
||||
validate_local_transition(config.as_ref(), &new_config)?;
|
||||
}
|
||||
validate_local_transition(self.state.load().config.as_ref(), &new_config)?;
|
||||
|
||||
// Create new service version without stopping old one
|
||||
// This allows existing operations to continue while new operations use new service
|
||||
match self.create_service_version(&new_config).await {
|
||||
match self.create_healthy_service_version(&new_config).await {
|
||||
Ok(new_service_version) => {
|
||||
// Get old version for logging (lock-free read)
|
||||
let old_version = self.current_service.load().as_ref().as_ref().map(|sv| sv.version);
|
||||
let old_version = self.state.load().current_service.as_ref().map(|sv| sv.version);
|
||||
|
||||
{
|
||||
let mut config = self.config.write().await;
|
||||
*config = Some(new_config);
|
||||
}
|
||||
persist().await?;
|
||||
|
||||
// Atomically switch to new service version (lock-free, instant CAS operation)
|
||||
// This is a true atomic operation - no waiting for locks, instant switch
|
||||
// Old service will be dropped when no more Arc references exist
|
||||
self.current_service.store(Arc::new(Some(new_service_version.clone())));
|
||||
|
||||
// Update status
|
||||
{
|
||||
let mut status = self.status.write().await;
|
||||
*status = KmsServiceStatus::Running;
|
||||
}
|
||||
self.publish_running(new_config, new_service_version.clone());
|
||||
|
||||
if let Some(old_ver) = old_version {
|
||||
info!(
|
||||
@@ -371,7 +426,7 @@ impl KmsServiceManager {
|
||||
/// Returns the manager from the current service version.
|
||||
/// Uses lock-free atomic load for optimal performance.
|
||||
pub async fn get_manager(&self) -> Option<Arc<KmsManager>> {
|
||||
self.current_service.load().as_ref().as_ref().map(|sv| sv.manager.clone())
|
||||
self.state.load().current_service.as_ref().map(|sv| sv.manager.clone())
|
||||
}
|
||||
|
||||
/// Get encryption service (if running)
|
||||
@@ -381,7 +436,7 @@ impl KmsServiceManager {
|
||||
/// This ensures new operations always use the latest service version,
|
||||
/// while existing operations continue using their Arc references.
|
||||
pub async fn get_encryption_service(&self) -> Option<Arc<ObjectEncryptionService>> {
|
||||
self.current_service.load().as_ref().as_ref().map(|sv| sv.service.clone())
|
||||
self.state.load().current_service.as_ref().map(|sv| sv.service.clone())
|
||||
}
|
||||
|
||||
/// Get current service version number
|
||||
@@ -389,14 +444,16 @@ impl KmsServiceManager {
|
||||
/// Useful for monitoring and debugging.
|
||||
/// Uses lock-free atomic load.
|
||||
pub async fn get_service_version(&self) -> Option<u64> {
|
||||
self.current_service.load().as_ref().as_ref().map(|sv| sv.version)
|
||||
self.state.load().current_service.as_ref().map(|sv| sv.version)
|
||||
}
|
||||
|
||||
/// Health check for the KMS service
|
||||
pub async fn health_check(&self) -> Result<bool> {
|
||||
let manager = self.get_manager().await;
|
||||
match manager {
|
||||
Some(manager) => {
|
||||
let checked_state = self.state.load_full();
|
||||
match checked_state.current_service.as_ref() {
|
||||
Some(service_version) => {
|
||||
let manager = service_version.manager.clone();
|
||||
let checked_version = service_version.version;
|
||||
// Perform health check on the backend
|
||||
match manager.health_check().await {
|
||||
Ok(healthy) => {
|
||||
@@ -407,9 +464,8 @@ impl KmsServiceManager {
|
||||
}
|
||||
Err(e) => {
|
||||
error!("KMS health check error: {}", e);
|
||||
// Update status to error
|
||||
let mut status = self.status.write().await;
|
||||
*status = KmsServiceStatus::Error(format!("Health check failed: {e}"));
|
||||
let _guard = self.lifecycle_mutex.lock().await;
|
||||
self.mark_health_error_if_current(checked_version, &e);
|
||||
Err(e)
|
||||
}
|
||||
}
|
||||
@@ -468,6 +524,33 @@ impl KmsServiceManager {
|
||||
manager: kms_manager,
|
||||
})
|
||||
}
|
||||
|
||||
async fn create_healthy_service_version(&self, config: &KmsConfig) -> Result<ServiceVersion> {
|
||||
let service_version = self.create_service_version(config).await?;
|
||||
if !service_version.manager.health_check().await? {
|
||||
return Err(KmsError::backend_error("KMS backend health check failed"));
|
||||
}
|
||||
Ok(service_version)
|
||||
}
|
||||
|
||||
fn publish_running(&self, config: KmsConfig, service_version: ServiceVersion) {
|
||||
self.state.store(Arc::new(RuntimeState {
|
||||
config: Some(config),
|
||||
status: KmsServiceStatus::Running,
|
||||
current_service: Some(service_version),
|
||||
}));
|
||||
}
|
||||
|
||||
fn mark_health_error_if_current(&self, checked_version: u64, error: &KmsError) {
|
||||
let current = self.state.load_full();
|
||||
if current.current_service.as_ref().map(|version| version.version) == Some(checked_version) {
|
||||
self.state.store(Arc::new(RuntimeState {
|
||||
config: current.config.clone(),
|
||||
status: KmsServiceStatus::Error(format!("Health check failed: {error}")),
|
||||
current_service: current.current_service.clone(),
|
||||
}));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for KmsServiceManager {
|
||||
@@ -500,6 +583,11 @@ pub async fn get_global_encryption_service() -> Option<Arc<ObjectEncryptionServi
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
|
||||
|
||||
fn static_config(key_id: &str, fill: u8) -> KmsConfig {
|
||||
KmsConfig::static_kms(key_id.to_string(), BASE64_STANDARD.encode([fill; 32]))
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn configure_rejects_insecure_development_defaults_before_state_update() {
|
||||
@@ -517,8 +605,6 @@ mod tests {
|
||||
|
||||
#[tokio::test]
|
||||
async fn redacted_config_omits_static_key_material() {
|
||||
use base64::Engine as _;
|
||||
|
||||
let manager = KmsServiceManager::new();
|
||||
let encoded_key = base64::engine::general_purpose::STANDARD.encode([0x5au8; 32]);
|
||||
manager
|
||||
@@ -533,6 +619,136 @@ mod tests {
|
||||
assert!(static_config.secret_key.is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn configure_persistence_failure_leaves_state_unchanged() {
|
||||
let manager = KmsServiceManager::new();
|
||||
|
||||
let result = manager
|
||||
.configure_with_persistence(static_config("key-a", 0x11), || async { Err(KmsError::backend_error("persist failed")) })
|
||||
.await;
|
||||
|
||||
assert!(result.is_err());
|
||||
assert_eq!(manager.get_status().await, KmsServiceStatus::NotConfigured);
|
||||
assert!(manager.get_config().await.is_none());
|
||||
assert!(manager.get_encryption_service().await.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn configure_rejects_running_service_without_changing_snapshot() {
|
||||
let manager = KmsServiceManager::new();
|
||||
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
|
||||
manager.start().await.expect("start");
|
||||
let version = manager.get_service_version().await;
|
||||
|
||||
let result = manager.configure(static_config("key-b", 0x22)).await;
|
||||
|
||||
assert!(result.is_err());
|
||||
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
|
||||
assert_eq!(manager.get_service_version().await, version);
|
||||
assert_eq!(
|
||||
manager.get_config().await.and_then(|config| config.default_key_id),
|
||||
Some("key-a".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn reconfigure_persistence_failure_keeps_old_running_snapshot() {
|
||||
let manager = KmsServiceManager::new();
|
||||
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
|
||||
manager.start().await.expect("start");
|
||||
let old_version = manager.get_service_version().await;
|
||||
let old_service = manager.get_encryption_service().await.expect("old service");
|
||||
|
||||
let result = manager
|
||||
.reconfigure_with_persistence(static_config("key-b", 0x22), || async {
|
||||
Err(KmsError::backend_error("persist failed"))
|
||||
})
|
||||
.await;
|
||||
|
||||
assert!(result.is_err());
|
||||
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
|
||||
assert_eq!(manager.get_service_version().await, old_version);
|
||||
assert_eq!(
|
||||
manager.get_config().await.and_then(|config| config.default_key_id),
|
||||
Some("key-a".to_string())
|
||||
);
|
||||
assert!(Arc::ptr_eq(
|
||||
&old_service,
|
||||
&manager.get_encryption_service().await.expect("old service remains")
|
||||
));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn reconfigure_candidate_failure_keeps_old_running_snapshot() {
|
||||
let manager = KmsServiceManager::new();
|
||||
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
|
||||
manager.start().await.expect("start");
|
||||
let old_version = manager.get_service_version().await;
|
||||
let invalid_parent = tempfile::NamedTempFile::new().expect("temporary file");
|
||||
let invalid_config = KmsConfig::local(invalid_parent.path().join("keys")).with_insecure_development_defaults();
|
||||
|
||||
let result = manager.reconfigure(invalid_config).await;
|
||||
|
||||
assert!(result.is_err());
|
||||
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
|
||||
assert_eq!(manager.get_service_version().await, old_version);
|
||||
assert_eq!(
|
||||
manager.get_config().await.and_then(|config| config.default_key_id),
|
||||
Some("key-a".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn restart_never_unpublishes_the_running_service() {
|
||||
let manager = Arc::new(KmsServiceManager::new());
|
||||
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
|
||||
manager.start().await.expect("start");
|
||||
let old_version = manager.get_service_version().await.expect("old version");
|
||||
let restarting = {
|
||||
let manager = manager.clone();
|
||||
tokio::spawn(async move { manager.restart().await })
|
||||
};
|
||||
|
||||
while !restarting.is_finished() {
|
||||
assert!(manager.get_encryption_service().await.is_some());
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
restarting.await.expect("restart task").expect("restart");
|
||||
|
||||
assert!(manager.get_encryption_service().await.is_some());
|
||||
assert!(manager.get_service_version().await.expect("new version") > old_version);
|
||||
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn start_or_restart_decides_under_the_lifecycle_lock() {
|
||||
let manager = KmsServiceManager::new();
|
||||
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
|
||||
|
||||
assert_eq!(manager.start_or_restart(false).await.expect("initial start"), KmsStartOutcome::Started);
|
||||
let first_version = manager.get_service_version().await.expect("first version");
|
||||
assert_eq!(
|
||||
manager.start_or_restart(false).await.expect("already running"),
|
||||
KmsStartOutcome::AlreadyRunning
|
||||
);
|
||||
assert_eq!(manager.get_service_version().await, Some(first_version));
|
||||
assert_eq!(manager.start_or_restart(true).await.expect("forced restart"), KmsStartOutcome::Restarted);
|
||||
assert!(manager.get_service_version().await.expect("restarted version") > first_version);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stale_health_failure_cannot_poison_new_service_status() {
|
||||
let manager = KmsServiceManager::new();
|
||||
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
|
||||
manager.start().await.expect("start");
|
||||
let old_version = manager.get_service_version().await.expect("old version");
|
||||
manager.restart().await.expect("restart");
|
||||
|
||||
manager.mark_health_error_if_current(old_version, &KmsError::backend_error("stale failure"));
|
||||
|
||||
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn forbidden_local_master_key_change_preserves_running_config_and_service() {
|
||||
use crate::types::{CreateKeyRequest, KeyUsage};
|
||||
|
||||
Reference in New Issue
Block a user