diff --git a/rustfs/src/admin/handlers/kms.rs b/rustfs/src/admin/handlers/kms.rs index 1f3b126f1..e6bf1daa4 100644 --- a/rustfs/src/admin/handlers/kms.rs +++ b/rustfs/src/admin/handlers/kms.rs @@ -19,58 +19,16 @@ use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::auth::{check_key_valid, get_session_token}; use crate::server::RemoteAddr; -use base64::Engine; use hyper::{HeaderMap, StatusCode}; use matchit::Params; -use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; -use rustfs_kms::{get_global_encryption_service, types::*}; +use rustfs_kms::get_global_encryption_service; use rustfs_policy::policy::action::{Action, AdminAction}; use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; use serde::{Deserialize, Serialize}; use serde_json; -use std::collections::HashMap; use tracing::{error, info, warn}; -#[derive(Debug, Serialize, Deserialize)] -pub struct CreateKeyApiRequest { - pub key_usage: Option, - pub description: Option, - pub tags: Option>, -} - -#[derive(Debug, Serialize, Deserialize)] -pub struct CreateKeyApiResponse { - pub key_id: String, - pub key_metadata: KeyMetadata, -} - -#[derive(Debug, Serialize, Deserialize)] -pub struct DescribeKeyApiResponse { - pub key_metadata: KeyMetadata, -} - -#[derive(Debug, Serialize, Deserialize)] -pub struct ListKeysApiResponse { - pub keys: Vec, - pub truncated: bool, - pub next_marker: Option, -} - -#[derive(Debug, Serialize, Deserialize)] -pub struct GenerateDataKeyApiRequest { - pub key_id: String, - pub key_spec: KeySpec, - pub encryption_context: Option>, -} - -#[derive(Debug, Serialize, Deserialize)] -pub struct GenerateDataKeyApiResponse { - pub key_id: String, - pub plaintext_key: String, // Base64 encoded - pub ciphertext_blob: String, // Base64 encoded -} - #[derive(Debug, Serialize, Deserialize)] pub struct KmsStatusResponse { pub backend_type: String, @@ -95,292 +53,13 @@ pub struct KmsConfigResponse { pub default_key_id: Option, } -fn extract_query_params(uri: &hyper::Uri) -> HashMap { - let mut params = HashMap::new(); - if let Some(query) = uri.query() { - query.split('&').for_each(|pair| { - if let Some((key, value)) = pair.split_once('=') { - params.insert( - urlencoding::decode(key).unwrap_or_default().into_owned(), - urlencoding::decode(value).unwrap_or_default().into_owned(), - ); - } - }); - } - params -} - pub fn register_kms_route(r: &mut S3Router) -> std::io::Result<()> { kms_management::register_kms_management_route(r)?; kms_dynamic::register_kms_dynamic_route(r)?; kms_keys::register_kms_key_route(r)?; - Ok(()) } -/// Create a new KMS master key -pub struct CreateKeyHandler {} - -#[async_trait::async_trait] -impl Operation for CreateKeyHandler { - async fn call(&self, mut req: S3Request, _params: Params<'_, '_>) -> S3Result> { - let Some(cred) = req.credentials else { - return Err(s3_error!(InvalidRequest, "authentication required")); - }; - - let (cred, owner) = - check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; - - validate_admin_request( - &req.headers, - &cred, - owner, - false, - vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], // TODO: Add specific KMS action - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; - - let body = req - .input - .store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE) - .await - .map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?; - - let request: CreateKeyApiRequest = if body.is_empty() { - CreateKeyApiRequest { - key_usage: Some(KeyUsage::EncryptDecrypt), - description: None, - tags: None, - } - } else { - serde_json::from_slice(&body).map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))? - }; - - let Some(service) = get_global_encryption_service().await else { - return Err(s3_error!(InternalError, "KMS service not initialized")); - }; - - // Extract key name from tags if provided - let tags = request.tags.unwrap_or_default(); - let key_name = tags.get("name").cloned(); - - let kms_request = CreateKeyRequest { - key_name, - key_usage: request.key_usage.unwrap_or(KeyUsage::EncryptDecrypt), - description: request.description, - tags, - origin: Some("AWS_KMS".to_string()), - policy: None, - }; - - match service.create_key(kms_request).await { - Ok(response) => { - let api_response = CreateKeyApiResponse { - key_id: response.key_id, - key_metadata: response.key_metadata, - }; - - let data = serde_json::to_vec(&api_response) - .map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; - - let mut headers = HeaderMap::new(); - headers.insert(CONTENT_TYPE, "application/json".parse().unwrap()); - - Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) - } - Err(e) => { - error!("Failed to create KMS key: {}", e); - Err(s3_error!(InternalError, "failed to create key: {}", e)) - } - } - } -} - -/// Describe a KMS key -pub struct DescribeKeyHandler {} - -#[async_trait::async_trait] -impl Operation for DescribeKeyHandler { - async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { - let Some(cred) = req.credentials else { - return Err(s3_error!(InvalidRequest, "authentication required")); - }; - - let (cred, owner) = - check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; - - validate_admin_request( - &req.headers, - &cred, - owner, - false, - vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; - - let query_params = extract_query_params(&req.uri); - let Some(key_id) = query_params.get("keyId") else { - return Err(s3_error!(InvalidRequest, "missing keyId parameter")); - }; - - let Some(service) = get_global_encryption_service().await else { - return Err(s3_error!(InternalError, "KMS service not initialized")); - }; - - let request = DescribeKeyRequest { key_id: key_id.clone() }; - - match service.describe_key(request).await { - Ok(response) => { - let api_response = DescribeKeyApiResponse { - key_metadata: response.key_metadata, - }; - - let data = serde_json::to_vec(&api_response) - .map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; - - let mut headers = HeaderMap::new(); - headers.insert(CONTENT_TYPE, "application/json".parse().unwrap()); - - Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) - } - Err(e) => { - error!("Failed to describe KMS key {}: {}", key_id, e); - Err(s3_error!(InternalError, "failed to describe key: {}", e)) - } - } - } -} - -/// List KMS keys -pub struct ListKeysHandler {} - -#[async_trait::async_trait] -impl Operation for ListKeysHandler { - async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { - let Some(cred) = req.credentials else { - return Err(s3_error!(InvalidRequest, "authentication required")); - }; - - let (cred, owner) = - check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; - - validate_admin_request( - &req.headers, - &cred, - owner, - false, - vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; - - let query_params = extract_query_params(&req.uri); - 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 { - return Err(s3_error!(InternalError, "KMS service not initialized")); - }; - - let request = ListKeysRequest { - limit: Some(limit), - marker, - status_filter: None, - usage_filter: None, - }; - - match service.list_keys(request).await { - Ok(response) => { - let api_response = ListKeysApiResponse { - keys: response.keys, - truncated: response.truncated, - next_marker: response.next_marker, - }; - - let data = serde_json::to_vec(&api_response) - .map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; - - let mut headers = HeaderMap::new(); - headers.insert(CONTENT_TYPE, "application/json".parse().unwrap()); - - Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) - } - Err(e) => { - error!("Failed to list KMS keys: {}", e); - Err(s3_error!(InternalError, "failed to list keys: {}", e)) - } - } - } -} - -/// Generate data encryption key -pub struct GenerateDataKeyHandler {} - -#[async_trait::async_trait] -impl Operation for GenerateDataKeyHandler { - async fn call(&self, mut req: S3Request, _params: Params<'_, '_>) -> S3Result> { - let Some(cred) = req.credentials else { - return Err(s3_error!(InvalidRequest, "authentication required")); - }; - - let (cred, owner) = - check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; - - validate_admin_request( - &req.headers, - &cred, - owner, - false, - vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; - - let body = req - .input - .store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE) - .await - .map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?; - - 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 { - return Err(s3_error!(InternalError, "KMS service not initialized")); - }; - - let kms_request = GenerateDataKeyRequest { - key_id: request.key_id, - key_spec: request.key_spec, - encryption_context: request.encryption_context.unwrap_or_default(), - }; - - match service.generate_data_key(kms_request).await { - Ok(response) => { - let api_response = GenerateDataKeyApiResponse { - key_id: response.key_id, - plaintext_key: base64::prelude::BASE64_STANDARD.encode(&response.plaintext_key), - ciphertext_blob: base64::prelude::BASE64_STANDARD.encode(&response.ciphertext_blob), - }; - - let data = serde_json::to_vec(&api_response) - .map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; - - let mut headers = HeaderMap::new(); - headers.insert(CONTENT_TYPE, "application/json".parse().unwrap()); - - Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) - } - Err(e) => { - error!("Failed to generate data key: {}", e); - Err(s3_error!(InternalError, "failed to generate data key: {}", e)) - } - } - } -} - /// Get KMS service status pub struct KmsStatusHandler {} diff --git a/rustfs/src/admin/handlers/kms_keys.rs b/rustfs/src/admin/handlers/kms_keys.rs index 558a4da96..e0039c190 100644 --- a/rustfs/src/admin/handlers/kms_keys.rs +++ b/rustfs/src/admin/handlers/kms_keys.rs @@ -18,10 +18,11 @@ use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; 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_kms_service_manager, types::*}; +use rustfs_kms::{KmsError, get_global_encryption_service, get_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}; @@ -46,6 +47,45 @@ pub struct CreateKmsKeyResponse { pub key_metadata: Option, } +#[derive(Debug, Serialize, Deserialize)] +pub struct CreateKeyApiRequest { + pub key_usage: Option, + pub description: Option, + pub tags: Option>, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct CreateKeyApiResponse { + pub key_id: String, + pub key_metadata: KeyMetadata, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct DescribeKeyApiResponse { + pub key_metadata: KeyMetadata, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct ListKeysApiResponse { + pub keys: Vec, + pub truncated: bool, + pub next_marker: Option, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GenerateDataKeyApiRequest { + pub key_id: String, + pub key_spec: KeySpec, + pub encryption_context: Option>, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GenerateDataKeyApiResponse { + pub key_id: String, + pub plaintext_key: String, // Base64 encoded + pub ciphertext_blob: String, // Base64 encoded +} + fn extract_query_params(uri: &hyper::Uri) -> HashMap { let mut params = HashMap::new(); if let Some(query) = uri.query() { @@ -95,6 +135,269 @@ pub fn register_kms_key_route(r: &mut S3Router) -> std::io::Resu Ok(()) } +/// Create a new KMS master key (legacy endpoint) +pub struct CreateKeyHandler {} + +#[async_trait::async_trait] +impl Operation for CreateKeyHandler { + async fn call(&self, mut req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let Some(cred) = req.credentials else { + return Err(s3_error!(InvalidRequest, "authentication required")); + }; + + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; + + validate_admin_request( + &req.headers, + &cred, + owner, + false, + vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], // TODO: Add specific KMS action + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await?; + + let body = req + .input + .store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE) + .await + .map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?; + + let request: CreateKeyApiRequest = if body.is_empty() { + CreateKeyApiRequest { + key_usage: Some(KeyUsage::EncryptDecrypt), + description: None, + tags: None, + } + } else { + serde_json::from_slice(&body).map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))? + }; + + let Some(service) = get_global_encryption_service().await else { + return Err(s3_error!(InternalError, "KMS service not initialized")); + }; + + // Extract key name from tags if provided + let tags = request.tags.unwrap_or_default(); + let key_name = tags.get("name").cloned(); + + let kms_request = CreateKeyRequest { + key_name, + key_usage: request.key_usage.unwrap_or(KeyUsage::EncryptDecrypt), + description: request.description, + tags, + origin: Some("AWS_KMS".to_string()), + policy: None, + }; + + match service.create_key(kms_request).await { + Ok(response) => { + let api_response = CreateKeyApiResponse { + key_id: response.key_id, + key_metadata: response.key_metadata, + }; + + let data = serde_json::to_vec(&api_response) + .map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; + + let mut headers = HeaderMap::new(); + headers.insert(CONTENT_TYPE, "application/json".parse().unwrap()); + + Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) + } + Err(e) => { + error!("Failed to create KMS key: {}", e); + Err(s3_error!(InternalError, "failed to create key: {}", e)) + } + } + } +} + +/// Describe a KMS key (legacy endpoint) +pub struct DescribeKeyHandler {} + +#[async_trait::async_trait] +impl Operation for DescribeKeyHandler { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let Some(cred) = req.credentials else { + return Err(s3_error!(InvalidRequest, "authentication required")); + }; + + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; + + validate_admin_request( + &req.headers, + &cred, + owner, + false, + vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await?; + + let query_params = extract_query_params(&req.uri); + let Some(key_id) = query_params.get("keyId") else { + return Err(s3_error!(InvalidRequest, "missing keyId parameter")); + }; + + let Some(service) = get_global_encryption_service().await else { + return Err(s3_error!(InternalError, "KMS service not initialized")); + }; + + let request = DescribeKeyRequest { key_id: key_id.clone() }; + + match service.describe_key(request).await { + Ok(response) => { + let api_response = DescribeKeyApiResponse { + key_metadata: response.key_metadata, + }; + + let data = serde_json::to_vec(&api_response) + .map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; + + let mut headers = HeaderMap::new(); + headers.insert(CONTENT_TYPE, "application/json".parse().unwrap()); + + Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) + } + Err(e) => { + error!("Failed to describe KMS key {}: {}", key_id, e); + Err(s3_error!(InternalError, "failed to describe key: {}", e)) + } + } + } +} + +/// List KMS keys (legacy endpoint) +pub struct ListKeysHandler {} + +#[async_trait::async_trait] +impl Operation for ListKeysHandler { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let Some(cred) = req.credentials else { + return Err(s3_error!(InvalidRequest, "authentication required")); + }; + + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; + + validate_admin_request( + &req.headers, + &cred, + owner, + false, + vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await?; + + let query_params = extract_query_params(&req.uri); + 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 { + return Err(s3_error!(InternalError, "KMS service not initialized")); + }; + + let request = ListKeysRequest { + limit: Some(limit), + marker, + status_filter: None, + usage_filter: None, + }; + + match service.list_keys(request).await { + Ok(response) => { + let api_response = ListKeysApiResponse { + keys: response.keys, + truncated: response.truncated, + next_marker: response.next_marker, + }; + + let data = serde_json::to_vec(&api_response) + .map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; + + let mut headers = HeaderMap::new(); + headers.insert(CONTENT_TYPE, "application/json".parse().unwrap()); + + Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) + } + Err(e) => { + error!("Failed to list KMS keys: {}", e); + Err(s3_error!(InternalError, "failed to list keys: {}", e)) + } + } + } +} + +/// Generate data encryption key (legacy endpoint) +pub struct GenerateDataKeyHandler {} + +#[async_trait::async_trait] +impl Operation for GenerateDataKeyHandler { + async fn call(&self, mut req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let Some(cred) = req.credentials else { + return Err(s3_error!(InvalidRequest, "authentication required")); + }; + + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; + + validate_admin_request( + &req.headers, + &cred, + owner, + false, + vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await?; + + let body = req + .input + .store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE) + .await + .map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?; + + 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 { + return Err(s3_error!(InternalError, "KMS service not initialized")); + }; + + let kms_request = GenerateDataKeyRequest { + key_id: request.key_id, + key_spec: request.key_spec, + encryption_context: request.encryption_context.unwrap_or_default(), + }; + + match service.generate_data_key(kms_request).await { + Ok(response) => { + let api_response = GenerateDataKeyApiResponse { + key_id: response.key_id, + plaintext_key: base64::prelude::BASE64_STANDARD.encode(&response.plaintext_key), + ciphertext_blob: base64::prelude::BASE64_STANDARD.encode(&response.ciphertext_blob), + }; + + let data = serde_json::to_vec(&api_response) + .map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; + + let mut headers = HeaderMap::new(); + headers.insert(CONTENT_TYPE, "application/json".parse().unwrap()); + + Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) + } + Err(e) => { + error!("Failed to generate data key: {}", e); + Err(s3_error!(InternalError, "failed to generate data key: {}", e)) + } + } + } +} + /// Create a new KMS key pub struct CreateKmsKeyHandler; diff --git a/rustfs/src/admin/handlers/kms_management.rs b/rustfs/src/admin/handlers/kms_management.rs index ed39f0cbc..08d9f9299 100644 --- a/rustfs/src/admin/handlers/kms_management.rs +++ b/rustfs/src/admin/handlers/kms_management.rs @@ -14,10 +14,8 @@ //! KMS management route registration. -use super::kms::{ - CreateKeyHandler, DescribeKeyHandler, GenerateDataKeyHandler, KmsClearCacheHandler, KmsConfigHandler, KmsStatusHandler, - ListKeysHandler, -}; +use super::kms::{KmsClearCacheHandler, KmsConfigHandler, KmsStatusHandler}; +use super::kms_keys::{CreateKeyHandler, DescribeKeyHandler, GenerateDataKeyHandler, ListKeysHandler}; use crate::admin::router::{AdminOperation, S3Router}; use crate::server::ADMIN_PREFIX; use hyper::Method;