diff --git a/rustfs/src/admin/handlers/kms_dynamic.rs b/rustfs/src/admin/handlers/kms_dynamic.rs index d8860769e..b866316f7 100644 --- a/rustfs/src/admin/handlers/kms_dynamic.rs +++ b/rustfs/src/admin/handlers/kms_dynamic.rs @@ -16,6 +16,7 @@ use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; +use crate::app::context::resolve_kms_runtime_service_manager; use crate::auth::{check_key_valid, get_session_token}; use crate::server::{ADMIN_PREFIX, RemoteAddr}; use hyper::{Method, StatusCode}; @@ -25,7 +26,7 @@ use rustfs_ecstore::config::com::{read_config, save_config}; use rustfs_ecstore::new_object_layer_fn; use rustfs_kms::{ ConfigureKmsRequest, ConfigureKmsResponse, KmsConfig, KmsConfigSummary, KmsServiceStatus, KmsStatusResponse, StartKmsRequest, - StartKmsResponse, StopKmsResponse, get_global_kms_service_manager, + StartKmsResponse, StopKmsResponse, }; use rustfs_policy::policy::action::{Action, AdminAction}; use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; @@ -34,6 +35,13 @@ use tracing::{error, info, warn}; /// Path to store KMS configuration in the cluster metadata const KMS_CONFIG_PATH: &str = "config/kms_config.json"; +fn kms_service_manager_from_context() -> std::sync::Arc { + resolve_kms_runtime_service_manager().unwrap_or_else(|| { + warn!("KMS service manager not initialized, initializing now as fallback"); + rustfs_kms::init_global_kms_service_manager() + }) +} + /// Save KMS configuration to cluster storage async fn save_kms_config(config: &KmsConfig) -> Result<(), String> { let Some(store) = new_object_layer_fn() else { @@ -160,11 +168,7 @@ impl Operation for ConfigureKmsHandler { info!("Configuring KMS with request: {:?}", configure_request); - let service_manager = get_global_kms_service_manager().unwrap_or_else(|| { - warn!("KMS service manager not initialized, initializing now as fallback"); - // Initialize the service manager as a fallback - rustfs_kms::init_global_kms_service_manager() - }); + let service_manager = kms_service_manager_from_context(); // Convert request to KmsConfig let kms_config = configure_request.to_kms_config(); @@ -256,11 +260,7 @@ impl Operation for StartKmsHandler { info!("Starting KMS service with force: {:?}", start_request.force); - let service_manager = get_global_kms_service_manager().unwrap_or_else(|| { - warn!("KMS service manager not initialized, initializing now as fallback"); - // Initialize the service manager as a fallback - rustfs_kms::init_global_kms_service_manager() - }); + let service_manager = kms_service_manager_from_context(); // Check if already running and force flag let current_status = service_manager.get_status().await; @@ -372,11 +372,7 @@ impl Operation for StopKmsHandler { info!("Stopping KMS service"); - let service_manager = get_global_kms_service_manager().unwrap_or_else(|| { - warn!("KMS service manager not initialized, initializing now as fallback"); - // Initialize the service manager as a fallback - rustfs_kms::init_global_kms_service_manager() - }); + let service_manager = kms_service_manager_from_context(); let (success, message, status) = match service_manager.stop().await { Ok(()) => { @@ -438,11 +434,7 @@ impl Operation for GetKmsStatusHandler { info!("Getting KMS service status"); - let service_manager = get_global_kms_service_manager().unwrap_or_else(|| { - warn!("KMS service manager not initialized, initializing now as fallback"); - // Initialize the service manager as a fallback - rustfs_kms::init_global_kms_service_manager() - }); + let service_manager = kms_service_manager_from_context(); let status = service_manager.get_status().await; let config = service_manager.get_config().await; @@ -531,11 +523,7 @@ impl Operation for ReconfigureKmsHandler { info!("Reconfiguring KMS with request: {:?}", configure_request); - let service_manager = get_global_kms_service_manager().unwrap_or_else(|| { - warn!("KMS service manager not initialized, initializing now as fallback"); - // Initialize the service manager as a fallback - rustfs_kms::init_global_kms_service_manager() - }); + let service_manager = kms_service_manager_from_context(); // Convert request to KmsConfig let kms_config = configure_request.to_kms_config(); diff --git a/rustfs/src/admin/handlers/kms_keys.rs b/rustfs/src/admin/handlers/kms_keys.rs index e0039c190..2351a698e 100644 --- a/rustfs/src/admin/handlers/kms_keys.rs +++ b/rustfs/src/admin/handlers/kms_keys.rs @@ -16,13 +16,14 @@ use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; +use crate::app::context::resolve_kms_runtime_service_manager; use crate::auth::{check_key_valid, get_session_token}; use crate::server::{ADMIN_PREFIX, RemoteAddr}; use base64::Engine; use hyper::{HeaderMap, Method, StatusCode}; use matchit::Params; use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; -use rustfs_kms::{KmsError, get_global_encryption_service, get_global_kms_service_manager, types::*}; +use rustfs_kms::{KmsError, init_global_kms_service_manager, types::*}; use rustfs_policy::policy::action::{Action, AdminAction}; use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; @@ -101,6 +102,15 @@ fn extract_query_params(uri: &hyper::Uri) -> HashMap { params } +fn kms_service_manager_from_context() -> Option> { + resolve_kms_runtime_service_manager() +} + +async fn kms_encryption_service_from_context() -> Option> { + let manager = kms_service_manager_from_context().unwrap_or_else(init_global_kms_service_manager); + manager.get_encryption_service().await +} + pub fn register_kms_key_route(r: &mut S3Router) -> std::io::Result<()> { r.insert( Method::POST, @@ -174,7 +184,7 @@ impl Operation for CreateKeyHandler { serde_json::from_slice(&body).map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))? }; - let Some(service) = get_global_encryption_service().await else { + let Some(service) = kms_encryption_service_from_context().await else { return Err(s3_error!(InternalError, "KMS service not initialized")); }; @@ -242,7 +252,7 @@ impl Operation for DescribeKeyHandler { return Err(s3_error!(InvalidRequest, "missing keyId parameter")); }; - let Some(service) = get_global_encryption_service().await else { + let Some(service) = kms_encryption_service_from_context().await else { return Err(s3_error!(InternalError, "KMS service not initialized")); }; @@ -297,7 +307,7 @@ impl Operation for ListKeysHandler { let limit = query_params.get("limit").and_then(|s| s.parse::().ok()).unwrap_or(100); let marker = query_params.get("marker").cloned(); - let Some(service) = get_global_encryption_service().await else { + let Some(service) = kms_encryption_service_from_context().await else { return Err(s3_error!(InternalError, "KMS service not initialized")); }; @@ -364,7 +374,7 @@ impl Operation for GenerateDataKeyHandler { let request: GenerateDataKeyApiRequest = serde_json::from_slice(&body).map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))?; - let Some(service) = get_global_encryption_service().await else { + let Some(service) = kms_encryption_service_from_context().await else { return Err(s3_error!(InternalError, "KMS service not initialized")); }; @@ -437,7 +447,7 @@ impl Operation for CreateKmsKeyHandler { serde_json::from_slice(&body).map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))? }; - let Some(service_manager) = get_global_kms_service_manager() else { + let Some(service_manager) = kms_service_manager_from_context() else { let response = CreateKmsKeyResponse { success: false, message: "KMS service manager not initialized".to_string(), @@ -590,7 +600,7 @@ impl Operation for DeleteKmsKeyHandler { serde_json::from_slice(&body).map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))? }; - let Some(service_manager) = get_global_kms_service_manager() else { + let Some(service_manager) = kms_service_manager_from_context() else { let response = DeleteKmsKeyResponse { success: false, message: "KMS service manager not initialized".to_string(), @@ -730,7 +740,7 @@ impl Operation for CancelKmsKeyDeletionHandler { serde_json::from_slice(&body).map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))? }; - let Some(service_manager) = get_global_kms_service_manager() else { + let Some(service_manager) = kms_service_manager_from_context() else { let response = CancelKmsKeyDeletionResponse { success: false, message: "KMS service manager not initialized".to_string(), @@ -837,7 +847,7 @@ impl Operation for ListKmsKeysHandler { let limit = query_params.get("limit").and_then(|s| s.parse::().ok()).unwrap_or(100); let marker = query_params.get("marker").cloned(); - let Some(service_manager) = get_global_kms_service_manager() else { + let Some(service_manager) = kms_service_manager_from_context() else { let response = ListKmsKeysResponse { success: false, message: "KMS service manager not initialized".to_string(), @@ -958,7 +968,7 @@ impl Operation for DescribeKmsKeyHandler { return Ok(S3Response::with_headers((StatusCode::BAD_REQUEST, Body::from(data)), headers)); }; - let Some(service_manager) = get_global_kms_service_manager() else { + let Some(service_manager) = kms_service_manager_from_context() else { let response = DescribeKmsKeyResponse { success: false, message: "KMS service manager not initialized".to_string(), diff --git a/rustfs/src/app/context.rs b/rustfs/src/app/context.rs index f2ec0b01e..e0d106a47 100644 --- a/rustfs/src/app/context.rs +++ b/rustfs/src/app/context.rs @@ -27,7 +27,7 @@ use rustfs_ecstore::global::get_global_region; use rustfs_ecstore::store::ECStore; use rustfs_ecstore::tier::tier::TierConfigMgr; use rustfs_iam::{store::object::ObjectStore, sys::IamSys}; -use rustfs_kms::KmsServiceManager; +use rustfs_kms::{KmsServiceManager, get_global_kms_service_manager}; use rustfs_notify::{EventArgs, NotificationError, notifier_global}; use rustfs_targets::{EventName, arn::TargetID}; use std::sync::{Arc, OnceLock}; @@ -44,6 +44,11 @@ pub trait KmsInterface: Send + Sync { fn handle(&self) -> Arc; } +/// KMS runtime interface for application-layer and admin handler integration. +pub trait KmsRuntimeInterface: Send + Sync { + fn service_manager(&self) -> Option>; +} + /// Notify interface for application-layer use-cases. #[async_trait] pub trait NotifyInterface: Send + Sync { @@ -127,6 +132,16 @@ impl KmsInterface for KmsHandle { } } +/// Default KMS runtime interface adapter. +#[derive(Default)] +pub struct KmsRuntimeHandle; + +impl KmsRuntimeInterface for KmsRuntimeHandle { + fn service_manager(&self) -> Option> { + get_global_kms_service_manager() + } +} + /// Default notify interface adapter. #[derive(Default)] pub struct NotifyHandle; @@ -217,6 +232,7 @@ pub struct AppContext { object_store: Arc, iam: Arc, kms: Arc, + kms_runtime: Arc, notify: Arc, bucket_metadata: Arc, endpoints: Arc, @@ -232,6 +248,7 @@ impl AppContext { object_store, iam, kms, + kms_runtime: default_kms_runtime_interface(), notify: default_notify_interface(), bucket_metadata: default_bucket_metadata_interface(), endpoints: default_endpoints_interface(), @@ -262,6 +279,10 @@ impl AppContext { self.kms.clone() } + pub fn kms_runtime(&self) -> Arc { + self.kms_runtime.clone() + } + pub fn notify(&self) -> Arc { self.notify.clone() } @@ -295,6 +316,17 @@ pub fn default_notify_interface() -> Arc { Arc::new(NotifyHandle) } +pub fn default_kms_runtime_interface() -> Arc { + Arc::new(KmsRuntimeHandle) +} + +/// Resolve KMS runtime service manager using AppContext-first precedence. +pub fn resolve_kms_runtime_service_manager() -> Option> { + get_global_app_context() + .and_then(|context| context.kms_runtime().service_manager()) + .or_else(|| default_kms_runtime_interface().service_manager()) +} + pub fn default_bucket_metadata_interface() -> Arc { Arc::new(BucketMetadataHandle) }