refactor(admin): route kms handlers via app context (#1965)

This commit is contained in:
安正超
2026-02-26 10:32:16 +08:00
committed by GitHub
parent 1c01c3d73a
commit 4b82cc20bb
3 changed files with 67 additions and 37 deletions
+14 -26
View File
@@ -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<rustfs_kms::KmsServiceManager> {
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();
+20 -10
View File
@@ -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<String, String> {
params
}
fn kms_service_manager_from_context() -> Option<std::sync::Arc<rustfs_kms::KmsServiceManager>> {
resolve_kms_runtime_service_manager()
}
async fn kms_encryption_service_from_context() -> Option<std::sync::Arc<rustfs_kms::ObjectEncryptionService>> {
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<AdminOperation>) -> 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::<u32>().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::<u32>().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(),
+33 -1
View File
@@ -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<KmsServiceManager>;
}
/// KMS runtime interface for application-layer and admin handler integration.
pub trait KmsRuntimeInterface: Send + Sync {
fn service_manager(&self) -> Option<Arc<KmsServiceManager>>;
}
/// 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<Arc<KmsServiceManager>> {
get_global_kms_service_manager()
}
}
/// Default notify interface adapter.
#[derive(Default)]
pub struct NotifyHandle;
@@ -217,6 +232,7 @@ pub struct AppContext {
object_store: Arc<ECStore>,
iam: Arc<dyn IamInterface>,
kms: Arc<dyn KmsInterface>,
kms_runtime: Arc<dyn KmsRuntimeInterface>,
notify: Arc<dyn NotifyInterface>,
bucket_metadata: Arc<dyn BucketMetadataInterface>,
endpoints: Arc<dyn EndpointsInterface>,
@@ -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<dyn KmsRuntimeInterface> {
self.kms_runtime.clone()
}
pub fn notify(&self) -> Arc<dyn NotifyInterface> {
self.notify.clone()
}
@@ -295,6 +316,17 @@ pub fn default_notify_interface() -> Arc<dyn NotifyInterface> {
Arc::new(NotifyHandle)
}
pub fn default_kms_runtime_interface() -> Arc<dyn KmsRuntimeInterface> {
Arc::new(KmsRuntimeHandle)
}
/// Resolve KMS runtime service manager using AppContext-first precedence.
pub fn resolve_kms_runtime_service_manager() -> Option<Arc<KmsServiceManager>> {
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<dyn BucketMetadataInterface> {
Arc::new(BucketMetadataHandle)
}