diff --git a/rustfs/src/admin/handlers/kms_audit.rs b/rustfs/src/admin/handlers/kms_audit.rs new file mode 100644 index 000000000..35b9f6c75 --- /dev/null +++ b/rustfs/src/admin/handlers/kms_audit.rs @@ -0,0 +1,643 @@ +// 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. + +//! Audit adapter for the KMS admin API. +//! +//! Maps the KMS-side audit contract onto the server's existing audit pipeline: +//! a [`KmsAuditRecord`] becomes an [`AuditEntry`] and travels to the same +//! targets as S3 audit entries, so KMS management activity lands in whatever +//! SIEM a deployment already operates instead of a KMS-private channel. +//! +//! Two emission paths meet here: +//! +//! * [`KmsAdminAuditSink`] converts the records [`rustfs_kms::KmsManager`] +//! builds for the operations it serves, covering their success and failure +//! outcomes. +//! * The handlers emit directly for what the KMS layer never sees: a request +//! rejected by the authorization gate, and the endpoints that have no +//! context-aware KMS entry point (data-key derivation and service control). +//! +//! # Redaction +//! +//! No field copied into an entry can hold key material. The recorded failure +//! reason is the [`error_class`] vocabulary, never the error message, so a +//! backend that echoes request data into an error cannot leak it here; the +//! caller-supplied encryption context arrives already reduced by +//! [`rustfs_kms::redact_encryption_context`]. + +use crate::server::RemoteAddr; +use crate::storage::access::request_context_from_extensions; +use crate::storage::helper::spawn_background_with_context; +use hashbrown::HashMap; +use rustfs_audit::entity::{ApiDetailsBuilder, AuditEntry, AuditEntryBuilder}; +use rustfs_audit::global::AuditLogger; +use rustfs_credentials::Credentials; +use rustfs_kms::audit::error_class; +use rustfs_kms::types::OperationContext; +use rustfs_kms::{KmsAuditOperation, KmsAuditOutcome, KmsAuditRecord, KmsAuditSink, KmsError}; +use rustfs_s3_types::EventName; +use rustfs_targets::get_request_user_agent; +use s3s::S3Result; +use serde_json::Value; +use std::collections::BTreeMap; +use std::time::{Duration, Instant}; + +/// Audit entry schema version, shared with the S3 request path so a consumer +/// parses KMS entries with the parser it already has. +const AUDIT_ENTRY_VERSION: &str = "1.0"; + +/// `trigger` value marking an entry as produced by the KMS admin API. +const AUDIT_TRIGGER: &str = "kms-admin"; + +/// `type` value letting a consumer separate KMS entries from S3 ones without +/// enumerating event names. +const AUDIT_ENTRY_TYPE: &str = "kms"; + +/// Shared empty context for the emission paths that carry none. +static EMPTY_ENCRYPTION_CONTEXT: BTreeMap = BTreeMap::new(); + +/// [`OperationContext::additional_context`] key carrying the canonical request +/// id, promoted to [`AuditEntry::request_id`] by the conversion below. +const CONTEXT_REQUEST_ID: &str = "requestID"; + +/// [`OperationContext::additional_context`] key carrying the account a +/// temporary or service credential belongs to. +const CONTEXT_PARENT_USER: &str = "parentUser"; + +/// Failure class recorded when the authorization gate rejects a request. The +/// gate answers before any KMS error exists, so the class is named here rather +/// than derived from one. +const ERROR_CLASS_ACCESS_DENIED: &str = "access_denied"; + +/// KMS admin operations audited by the handler itself. +/// +/// The KMS manager audits everything it serves, but two groups never reach it +/// with a caller context: data-key derivation goes through +/// `ObjectEncryptionService`, which has no context-aware entry point, and the +/// service-control endpoints act on the service rather than on a key. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum KmsAdminOperation { + /// Derivation of a data key from a master key. + GenerateDataKey, + /// Installation of a KMS configuration. + Configure, + /// Replacement of the running configuration. + Reconfigure, + /// Start or restart of the KMS service. + Start, + /// Stop of the KMS service. + Stop, +} + +impl KmsAdminOperation { + /// Stable operation name for audit consumers. + fn as_str(self) -> &'static str { + match self { + Self::GenerateDataKey => "GenerateDataKey", + Self::Configure => "Configure", + Self::Reconfigure => "Reconfigure", + Self::Start => "Start", + Self::Stop => "Stop", + } + } + + /// Event name published to audit consumers. + fn event(self) -> EventName { + match self { + // Deriving a data key uses the master key without changing it, + // which is what the access event already denotes. + Self::GenerateDataKey => EventName::KmsKeyAccessed, + Self::Configure | Self::Reconfigure => EventName::KmsServiceConfigured, + Self::Start => EventName::KmsServiceStarted, + Self::Stop => EventName::KmsServiceStopped, + } + } +} + +/// Per-request audit state for a KMS admin endpoint. +/// +/// Built once the caller is authenticated and before the authorization gate +/// runs, so a denial is attributable to the identity that was rejected. +pub(super) struct KmsAdminAudit { + context: OperationContext, + started: Instant, +} + +impl KmsAdminAudit { + /// Build the audit state for an authenticated admin request. + /// + /// Takes the request in pieces because every KMS handler moves + /// `credentials` out of it before reaching this point. + pub(super) fn from_request(extensions: &http::Extensions, headers: &http::HeaderMap, cred: &Credentials) -> Self { + // The peer address the authorization gate was given, so the audit trail + // and the policy decision agree on where the request came from. + let source_ip = extensions + .get::>() + .and_then(|opt| opt.map(|addr| addr.0.ip().to_string())); + let request_id = request_context_from_extensions(extensions).map(|context| context.request_id); + + Self::new( + &cred.access_key, + &cred.parent_user, + source_ip, + Some(get_request_user_agent(headers)), + request_id, + ) + } + + /// Assemble the operation context from already-extracted request values. + /// + /// Separate from [`Self::from_request`] so the mapping onto an audit entry + /// is testable without an HTTP request. + fn new( + access_key: &str, + parent_user: &str, + source_ip: Option, + user_agent: Option, + request_id: Option, + ) -> Self { + let mut context = OperationContext::new(access_key.to_string()); + context.source_ip = source_ip; + context.user_agent = user_agent.filter(|agent| !agent.is_empty()); + if let Some(request_id) = request_id.filter(|id| !id.is_empty()) { + context = context.with_context(CONTEXT_REQUEST_ID.to_string(), request_id); + } + if !parent_user.is_empty() && parent_user != access_key { + context = context.with_context(CONTEXT_PARENT_USER.to_string(), parent_user.to_string()); + } + + Self { + context, + started: Instant::now(), + } + } + + /// The context handed to the KMS manager's `*_with_context` methods, so the + /// record it builds carries the same identity as a handler-side entry. + pub(super) fn context(&self) -> &OperationContext { + &self.context + } + + /// Audit an authorization denial and propagate it unchanged. + /// + /// A denied request never reaches the KMS layer, so this is the only place + /// that can record it. Wrapping the gate's result keeps the record and the + /// rejection on the same path: there is no way to return the error without + /// producing the entry. + pub(super) fn gate(&self, result: S3Result, operation: KmsAuditOperation, key_id: Option<&str>) -> S3Result { + if result.is_err() { + self.emit(operation.as_str(), operation.event_name(), key_id, Some(ERROR_CLASS_ACCESS_DENIED)); + } + result + } + + /// Audit an authorization denial for an operation the handler owns. + pub(super) fn gate_admin(&self, result: S3Result, operation: KmsAdminOperation, key_id: Option<&str>) -> S3Result { + if result.is_err() { + self.emit(operation.as_str(), operation.event(), key_id, Some(ERROR_CLASS_ACCESS_DENIED)); + } + result + } + + /// Audit the outcome of an operation the handler served itself. + pub(super) fn finish(&self, operation: KmsAdminOperation, key_id: Option<&str>, error: Option<&KmsError>) { + self.emit(operation.as_str(), operation.event(), key_id, error.map(error_class)); + } + + /// Audit the outcome of a handler-served operation whose failure is not a + /// [`KmsError`], such as a service-control call rejected by the manager's + /// own state machine. + pub(super) fn finish_with_class(&self, operation: KmsAdminOperation, error_class: Option<&'static str>) { + self.emit(operation.as_str(), operation.event(), None, error_class); + } + + fn emit(&self, operation: &str, event: EventName, key_id: Option<&str>, error_class: Option<&str>) { + dispatch(build_entry(KmsAuditFields { + operation, + event, + context: &self.context, + key_id, + key_version: None, + outcome: match error_class { + Some(_) => KmsAuditOutcome::Failure, + None => KmsAuditOutcome::Success, + }, + error_class, + backend: None, + retry_count: None, + latency: self.started.elapsed(), + encryption_context: &EMPTY_ENCRYPTION_CONTEXT, + })); + } +} + +/// Sends the KMS manager's audit records to the server's audit pipeline. +/// +/// Installed once at KMS service assembly. Delivery follows the pipeline's +/// established best-effort semantics: the KMS operation has already completed +/// when a record arrives, and nothing here can change its result. +pub struct KmsAdminAuditSink; + +impl KmsAuditSink for KmsAdminAuditSink { + fn emit(&self, record: KmsAuditRecord) { + dispatch(entry_from_record(&record)); + } +} + +/// Convert a KMS audit record into an audit entry. +fn entry_from_record(record: &KmsAuditRecord) -> AuditEntry { + build_entry(KmsAuditFields { + operation: record.operation.as_str(), + event: record.event, + context: &operation_context_of(record), + key_id: record.key_id.as_deref(), + key_version: record.key_version, + outcome: record.outcome, + error_class: record.error_class, + backend: Some(record.backend), + retry_count: record.retry_count, + latency: record.latency, + encryption_context: &record.encryption_context, + }) +} + +/// Rebuild the caller context a record was created from. The record flattens +/// it, so the shared entry builder gets it back in one shape. +fn operation_context_of(record: &KmsAuditRecord) -> OperationContext { + let mut context = OperationContext::new(record.principal.clone()); + context.operation_id = record.operation_id; + context.source_ip = record.source_ip.clone(); + context.user_agent = record.user_agent.clone(); + context.additional_context = record + .context + .iter() + .map(|(key, value)| (key.clone(), value.clone())) + .collect(); + context +} + +/// Everything one audited KMS operation contributes to an entry. +/// +/// A struct rather than an argument list because both emission paths fill the +/// same shape, and a positional mix-up between the operation name, the error +/// class and the backend would be invisible at the call site. +struct KmsAuditFields<'a> { + operation: &'a str, + event: EventName, + context: &'a OperationContext, + key_id: Option<&'a str>, + key_version: Option, + outcome: KmsAuditOutcome, + error_class: Option<&'a str>, + backend: Option<&'a str>, + retry_count: Option, + latency: Duration, + encryption_context: &'a BTreeMap, +} + +fn build_entry(fields: KmsAuditFields<'_>) -> AuditEntry { + let context = fields.context; + let api = ApiDetailsBuilder::new() + .name(fields.operation) + .status(fields.outcome.as_str()) + .time_to_response(format!("{:.2?}", fields.latency)) + .time_to_response_in_ns(fields.latency.as_nanos().to_string()) + .build(); + + let mut builder = AuditEntryBuilder::new(AUDIT_ENTRY_VERSION, fields.event, AUDIT_TRIGGER, api) + .entry_type(AUDIT_ENTRY_TYPE) + .access_key(&context.principal) + .tags(entry_tags(&fields)); + + if let Some(request_id) = context.additional_context.get(CONTEXT_REQUEST_ID) { + builder = builder.request_id(request_id); + } + if let Some(parent_user) = context.additional_context.get(CONTEXT_PARENT_USER) { + builder = builder.parent_user(parent_user); + } + if let Some(source_ip) = context.source_ip.as_deref() { + builder = builder.remote_host(source_ip); + } + if let Some(user_agent) = context.user_agent.as_deref() { + builder = builder.user_agent(user_agent); + } + // The class, not the message: an error rendered by a backend can quote the + // request that produced it, and an audit entry is a poor place to find out. + if let Some(error_class) = fields.error_class { + builder = builder.error(error_class); + } + + builder.build() +} + +fn entry_tags(fields: &KmsAuditFields<'_>) -> HashMap { + let context = fields.context; + let mut tags = HashMap::new(); + tags.insert("kmsOperation".to_string(), Value::String(fields.operation.to_string())); + tags.insert("kmsOutcome".to_string(), Value::String(fields.outcome.as_str().to_string())); + tags.insert("operationId".to_string(), Value::String(context.operation_id.to_string())); + + if let Some(key_id) = fields.key_id.filter(|value| !value.is_empty()) { + tags.insert("keyId".to_string(), Value::String(key_id.to_string())); + } + if let Some(key_version) = fields.key_version { + tags.insert("keyVersion".to_string(), Value::Number(key_version.into())); + } + if let Some(backend) = fields.backend { + tags.insert("kmsBackend".to_string(), Value::String(backend.to_string())); + } + if let Some(retry_count) = fields.retry_count { + tags.insert("retryCount".to_string(), Value::Number(retry_count.into())); + } + if !fields.encryption_context.is_empty() { + let redacted = fields + .encryption_context + .iter() + .map(|(key, value)| (key.clone(), Value::String(value.clone()))) + .collect::>(); + tags.insert("encryptionContext".to_string(), Value::Object(redacted)); + } + + // Correlation values that are not promoted to a dedicated field stay + // nested, so a future context key can never shadow a reserved tag. + let extra = context + .additional_context + .iter() + .filter(|(key, _)| key.as_str() != CONTEXT_REQUEST_ID && key.as_str() != CONTEXT_PARENT_USER) + .map(|(key, value)| (key.clone(), Value::String(value.clone()))) + .collect::>(); + if !extra.is_empty() { + tags.insert("kmsContext".to_string(), Value::Object(extra)); + } + + tags +} + +/// Hand a built entry to the audit pipeline. +/// +/// Off the caller's task for the same reason the S3 path defers its entry: the +/// operation is finished, and audit delivery must not extend its latency. With +/// no audit system configured the dispatch is a no-op. +fn dispatch(entry: AuditEntry) { + spawn_background_with_context(None, async move { + AuditLogger::log(entry).await; + }); +} + +#[cfg(test)] +mod tests { + use super::*; + use base64::Engine; + use rustfs_kms::backends::local::LocalKmsBackend; + use rustfs_kms::config::KmsConfig; + use rustfs_kms::types::{CreateKeyRequest, DescribeKeyRequest, GenerateDataKeyRequest, KeySpec}; + use rustfs_kms::{KmsAuditOperation, KmsManager}; + use std::sync::{Arc, Mutex}; + + /// Captures the entries a real `KmsManager` produces, exercising the same + /// sink the server installs. + #[derive(Default)] + struct CapturingSink { + entries: Mutex>, + } + + impl KmsAuditSink for CapturingSink { + fn emit(&self, record: KmsAuditRecord) { + self.entries + .lock() + .expect("capturing sink should not be poisoned") + .push(entry_from_record(&record)); + } + } + + impl CapturingSink { + fn drain(&self) -> Vec { + std::mem::take(&mut *self.entries.lock().expect("capturing sink should not be poisoned")) + } + } + + fn request_audit() -> KmsAdminAudit { + KmsAdminAudit::new( + "AKIAOPERATOR", + "root-account", + Some("10.1.2.3".to_string()), + Some("rustfs-admin/1".to_string()), + Some("req-kms-1".to_string()), + ) + } + + fn tag<'a>(entry: &'a AuditEntry, name: &str) -> Option<&'a Value> { + entry.tags.as_ref().and_then(|tags| tags.get(name)) + } + + async fn local_manager(temp_dir: &tempfile::TempDir, sink: Arc) -> KmsManager { + let config = KmsConfig::local(temp_dir.path().to_path_buf()).with_insecure_development_defaults(); + let backend = Arc::new( + LocalKmsBackend::new(config.clone()) + .await + .expect("local backend should build"), + ); + KmsManager::new(backend, config).with_audit_sink(sink) + } + + /// The identity, correlation id and outcome the KMS manager records must + /// all survive the mapping onto an audit entry. + #[tokio::test] + async fn manager_records_become_attributable_audit_entries() { + let temp_dir = tempfile::tempdir().expect("temp dir"); + let sink = Arc::new(CapturingSink::default()); + let manager = local_manager(&temp_dir, sink.clone()).await; + let audit = request_audit(); + + let key_id = manager + .create_key_with_context( + CreateKeyRequest { + key_name: Some("audited-key".to_string()), + ..Default::default() + }, + audit.context(), + ) + .await + .expect("key should be created") + .key_id; + + let entries = sink.drain(); + let entry = entries.first().expect("create should produce one entry"); + + assert_eq!(entry.event, EventName::KmsKeyCreated); + assert_eq!(entry.request_id.as_deref(), Some("req-kms-1")); + assert_eq!(entry.access_key.as_deref(), Some("AKIAOPERATOR")); + assert_eq!(entry.parent_user.as_deref(), Some("root-account")); + assert_eq!(entry.remote_host.as_deref(), Some("10.1.2.3")); + assert_eq!(entry.user_agent.as_deref(), Some("rustfs-admin/1")); + assert_eq!(entry.api.status.as_deref(), Some("success")); + assert_eq!(entry.error, None); + assert_eq!(tag(entry, "keyId"), Some(&Value::String(key_id.clone()))); + assert_eq!(tag(entry, "kmsBackend"), Some(&Value::String("local".to_string()))); + + // A failure must carry the class, and only the class. + manager + .describe_key_with_context( + DescribeKeyRequest { + key_id: "no-such-key".to_string(), + }, + audit.context(), + ) + .await + .expect_err("describe of a missing key should fail"); + + let entries = sink.drain(); + let entry = entries.first().expect("describe should produce one entry"); + assert_eq!(entry.api.status.as_deref(), Some("failure")); + assert_eq!(entry.error.as_deref(), Some("key_not_found")); + assert_eq!(entry.request_id.as_deref(), Some("req-kms-1")); + } + + /// A denied request never reaches the KMS layer, so the handler-side entry + /// is the only record of it. + #[test] + fn a_denied_request_is_audited_with_the_identity_that_was_rejected() { + let audit = request_audit(); + // The delete endpoint schedules deletion, which is the operation the + // KMS manager audits for a served request; a denial must be recorded + // under the same name so both outcomes of one endpoint correlate. + let entry = build_entry(KmsAuditFields { + operation: KmsAuditOperation::ScheduleKeyDeletion.as_str(), + event: KmsAuditOperation::ScheduleKeyDeletion.event_name(), + context: audit.context(), + key_id: Some("key-denied"), + key_version: None, + outcome: KmsAuditOutcome::Failure, + error_class: Some(ERROR_CLASS_ACCESS_DENIED), + backend: None, + retry_count: None, + latency: Duration::from_millis(1), + encryption_context: &EMPTY_ENCRYPTION_CONTEXT, + }); + + assert_eq!(entry.event, EventName::KmsKeyDeletionScheduled); + assert_eq!(entry.error.as_deref(), Some("access_denied")); + assert_eq!(entry.api.status.as_deref(), Some("failure")); + assert_eq!(entry.access_key.as_deref(), Some("AKIAOPERATOR")); + assert_eq!(entry.request_id.as_deref(), Some("req-kms-1")); + assert_eq!(tag(&entry, "keyId"), Some(&Value::String("key-denied".to_string()))); + } + + /// Every handler-owned operation must map to a KMS event, or its entries + /// would be indistinguishable from S3 traffic to a consumer filtering on + /// the event name. + #[test] + fn handler_owned_operations_map_to_kms_events() { + for operation in [ + KmsAdminOperation::GenerateDataKey, + KmsAdminOperation::Configure, + KmsAdminOperation::Reconfigure, + KmsAdminOperation::Start, + KmsAdminOperation::Stop, + ] { + assert!( + operation.event().is_kms(), + "{} must map to a KMS event, got {}", + operation.as_str(), + operation.event() + ); + } + } + + /// The data-key endpoint returns plaintext key material to its caller. The + /// audit entry for the same operation must not reproduce any of it, in raw + /// or base64 form. + #[tokio::test] + async fn a_data_key_audit_entry_carries_no_key_material() { + let temp_dir = tempfile::tempdir().expect("temp dir"); + let sink = Arc::new(CapturingSink::default()); + let manager = local_manager(&temp_dir, sink.clone()).await; + let audit = request_audit(); + + let key_id = manager + .create_key_with_context( + CreateKeyRequest { + key_name: Some("data-key-source".to_string()), + ..Default::default() + }, + audit.context(), + ) + .await + .expect("key should be created") + .key_id; + + let secret = "super-secret-grant-token"; + let response = manager + .generate_data_key(GenerateDataKeyRequest { + key_id: key_id.clone(), + key_spec: KeySpec::Aes256, + encryption_context: std::collections::HashMap::from([ + ("bucket".to_string(), "photos".to_string()), + ("grant_token".to_string(), secret.to_string()), + ]), + }) + .await + .expect("data key should be generated"); + + // What the endpoint hands back, and therefore what must not reappear. + let plaintext_b64 = base64::prelude::BASE64_STANDARD.encode(&response.plaintext_key); + let ciphertext_b64 = base64::prelude::BASE64_STANDARD.encode(&response.ciphertext_blob); + assert!(!response.plaintext_key.is_empty(), "the test must drive real key material"); + + let redacted = rustfs_kms::redact_encryption_context(&std::collections::HashMap::from([ + ("bucket".to_string(), "photos".to_string()), + ("grant_token".to_string(), secret.to_string()), + ])); + let entry = build_entry(KmsAuditFields { + operation: KmsAdminOperation::GenerateDataKey.as_str(), + event: KmsAdminOperation::GenerateDataKey.event(), + context: audit.context(), + key_id: Some(&key_id), + key_version: None, + outcome: KmsAuditOutcome::Success, + error_class: None, + backend: Some("local"), + retry_count: None, + latency: Duration::from_millis(1), + encryption_context: &redacted, + }); + + let rendered = serde_json::to_string(&entry).expect("audit entry should serialize"); + for material in [plaintext_b64.as_str(), ciphertext_b64.as_str(), secret] { + assert!( + !rendered.contains(material), + "audit entry must not reproduce key material or secret context: {rendered}" + ); + } + // Raw bytes cannot appear either, whatever encoding a future field uses. + assert!( + !rendered + .as_bytes() + .windows(response.plaintext_key.len()) + .any(|window| window == response.plaintext_key.as_slice()), + "audit entry must not reproduce raw key material" + ); + + // The location keys stay readable, which is why the context is audited. + let encryption_context = tag(&entry, "encryptionContext").expect("encryption context should be recorded"); + assert_eq!(encryption_context.get("bucket"), Some(&Value::String("photos".to_string()))); + assert!( + encryption_context + .get("grant_token") + .and_then(Value::as_str) + .is_some_and(|value| value.starts_with("sha256:")), + "a non-allowlisted context value must be digested" + ); + } +} diff --git a/rustfs/src/admin/handlers/kms_dynamic.rs b/rustfs/src/admin/handlers/kms_dynamic.rs index 526079d28..ba58f557e 100644 --- a/rustfs/src/admin/handlers/kms_dynamic.rs +++ b/rustfs/src/admin/handlers/kms_dynamic.rs @@ -14,6 +14,7 @@ //! KMS dynamic configuration admin API handlers +use super::kms_audit::{KmsAdminAudit, KmsAdminOperation}; use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::runtime_sources::{ @@ -531,15 +532,21 @@ impl Operation for ConfigureKmsHandler { 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, - kms_configure_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate_admin( + validate_admin_request( + &req.headers, + &cred, + owner, + false, + kms_configure_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await, + KmsAdminOperation::Configure, + None, + )?; let body = req .input @@ -615,6 +622,7 @@ impl Operation for ConfigureKmsHandler { ); let unconverged = broadcast_kms_config_reload().await; let (success, message) = local_success_with_peer_report("KMS configured successfully", &unconverged); + audit.finish(KmsAdminOperation::Configure, None, None); (success, message, status) } Err(e) => { @@ -628,6 +636,7 @@ impl Operation for ConfigureKmsHandler { error = %e, "admin kms dynamic state" ); + audit.finish(KmsAdminOperation::Configure, None, Some(&e)); let status = service_manager.get_status().await; (false, error_msg, status) } @@ -675,15 +684,21 @@ impl Operation for StartKmsHandler { 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, - kms_service_control_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate_admin( + validate_admin_request( + &req.headers, + &cred, + owner, + false, + kms_service_control_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await, + KmsAdminOperation::Start, + None, + )?; let body = req .input @@ -735,6 +750,7 @@ impl Operation for StartKmsHandler { status = ?status, "admin kms dynamic state" ); + audit.finish(KmsAdminOperation::Start, None, None); (true, "KMS service started successfully".to_string(), status) } Ok(rustfs_kms::KmsStartOutcome::Restarted) => { @@ -748,6 +764,7 @@ impl Operation for StartKmsHandler { status = ?status, "admin kms dynamic state" ); + audit.finish(KmsAdminOperation::Start, None, None); (true, "KMS service restarted successfully".to_string(), status) } Ok(rustfs_kms::KmsStartOutcome::AlreadyRunning) => { @@ -760,6 +777,9 @@ impl Operation for StartKmsHandler { state = "already_running", "admin kms dynamic state" ); + // A refusal, not an outage: recorded as a failed attempt so the + // trail shows the request without inventing a KMS error for it. + audit.finish_with_class(KmsAdminOperation::Start, Some("invalid_operation")); (false, "KMS service is already running. Use force=true to restart.".to_string(), status) } Err(e) => { @@ -773,6 +793,7 @@ impl Operation for StartKmsHandler { error = %e, "admin kms dynamic state" ); + audit.finish(KmsAdminOperation::Start, None, Some(&e)); let status = service_manager.get_status().await; (false, error_msg, status) } @@ -820,15 +841,21 @@ impl Operation for StopKmsHandler { 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, - kms_service_control_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate_admin( + validate_admin_request( + &req.headers, + &cred, + owner, + false, + kms_service_control_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await, + KmsAdminOperation::Stop, + None, + )?; info!( component = LOG_COMPONENT_ADMIN, @@ -853,6 +880,7 @@ impl Operation for StopKmsHandler { status = ?status, "admin kms dynamic state" ); + audit.finish(KmsAdminOperation::Stop, None, None); (true, "KMS service stopped successfully".to_string(), status) } Err(e) => { @@ -866,6 +894,7 @@ impl Operation for StopKmsHandler { error = %e, "admin kms dynamic state" ); + audit.finish(KmsAdminOperation::Stop, None, Some(&e)); let status = service_manager.get_status().await; (false, error_msg, status) } @@ -1004,15 +1033,21 @@ impl Operation for ReconfigureKmsHandler { 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, - kms_configure_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate_admin( + validate_admin_request( + &req.headers, + &cred, + owner, + false, + kms_configure_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await, + KmsAdminOperation::Reconfigure, + None, + )?; let body = req .input @@ -1089,6 +1124,7 @@ impl Operation for ReconfigureKmsHandler { let unconverged = broadcast_kms_config_reload().await; let (success, message) = local_success_with_peer_report("KMS reconfigured and restarted successfully", &unconverged); + audit.finish(KmsAdminOperation::Reconfigure, None, None); (success, message, status) } Err(e) => { @@ -1102,6 +1138,7 @@ impl Operation for ReconfigureKmsHandler { error = %e, "admin kms dynamic state" ); + audit.finish(KmsAdminOperation::Reconfigure, None, Some(&e)); let status = service_manager.get_status().await; (false, error_msg, status) } diff --git a/rustfs/src/admin/handlers/kms_key_lifecycle.rs b/rustfs/src/admin/handlers/kms_key_lifecycle.rs index 24c7a3c34..32fbdf43c 100644 --- a/rustfs/src/admin/handlers/kms_key_lifecycle.rs +++ b/rustfs/src/admin/handlers/kms_key_lifecycle.rs @@ -14,6 +14,7 @@ //! KMS key lifecycle admin API handlers: enable, disable and rotate. +use super::kms_audit::KmsAdminAudit; use super::kms_keys::{extract_query_params, scoped_key_id}; use crate::admin::auth::validate_admin_request_with_kms_key; use crate::admin::router::{AdminOperation, Operation, S3Router}; @@ -24,8 +25,8 @@ use hyper::{HeaderMap, Method, StatusCode}; use matchit::Params; use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; use rustfs_kms::{ - KmsError, KmsManager, - types::{DescribeKeyRequest, KeyMetadata}, + KmsAuditOperation, KmsError, KmsManager, + types::{DescribeKeyRequest, KeyMetadata, OperationContext}, }; use rustfs_policy::policy::action::{Action, KmsAction}; use s3s::header::CONTENT_TYPE; @@ -60,6 +61,15 @@ enum LifecycleOperation { } impl LifecycleOperation { + /// The KMS audit operation this endpoint is recorded as. + fn audit_operation(self) -> KmsAuditOperation { + match self { + Self::Enable => KmsAuditOperation::EnableKey, + Self::Disable => KmsAuditOperation::DisableKey, + Self::Rotate => KmsAuditOperation::RotateKey, + } + } + fn action(self) -> &'static str { match self { Self::Enable => "enable_key", @@ -148,11 +158,12 @@ async fn execute_lifecycle( manager: &KmsManager, key_id: &str, operation: LifecycleOperation, + context: &OperationContext, ) -> (StatusCode, KmsKeyLifecycleResponse) { let result = match operation { - LifecycleOperation::Enable => manager.enable_key(key_id).await, - LifecycleOperation::Disable => manager.disable_key(key_id).await, - LifecycleOperation::Rotate => manager.rotate_key(key_id).await, + LifecycleOperation::Enable => manager.enable_key_with_context(key_id, context).await, + LifecycleOperation::Disable => manager.disable_key_with_context(key_id, context).await, + LifecycleOperation::Rotate => manager.rotate_key_with_context(key_id, context).await, }; match result { @@ -256,21 +267,23 @@ async fn handle_lifecycle_request( // on, so it has to be resolved before the gate runs. A read failure is // surfaced only afterwards so an unauthorized caller still sees AccessDenied. let body = req.input.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE).await; + let scoped = body.as_ref().ok().and_then(|body| scoped_key_id(body, &req.uri)); + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); - validate_admin_request_with_kms_key( - &req.headers, - &cred, - owner, - false, - actions, - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - body.as_ref() - .ok() - .and_then(|body| scoped_key_id(body, &req.uri)) - .as_deref() - .unwrap_or_default(), - ) - .await?; + audit.gate( + validate_admin_request_with_kms_key( + &req.headers, + &cred, + owner, + false, + actions, + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + scoped.as_deref().unwrap_or_default(), + ) + .await, + operation.audit_operation(), + scoped.as_deref(), + )?; let body = body.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?; @@ -300,7 +313,7 @@ async fn handle_lifecycle_request( return unavailable_response("kms service is not running", request.key_id); }; - let (status, response) = execute_lifecycle(&manager, &request.key_id, operation).await; + let (status, response) = execute_lifecycle(&manager, &request.key_id, operation, audit.context()).await; json_response(status, &response) } @@ -383,18 +396,21 @@ mod tests { let manager = local_manager(&temp_dir).await; let key_id = create_key(&manager, "lifecycle-round-trip").await; - let (status, response) = execute_lifecycle(&manager, &key_id, LifecycleOperation::Disable).await; + let (status, response) = + execute_lifecycle(&manager, &key_id, LifecycleOperation::Disable, &OperationContext::internal()).await; assert_eq!(status, StatusCode::OK); assert!(response.success); assert_eq!(key_state(&response), KeyState::Disabled); - let (status, response) = execute_lifecycle(&manager, &key_id, LifecycleOperation::Enable).await; + let (status, response) = + execute_lifecycle(&manager, &key_id, LifecycleOperation::Enable, &OperationContext::internal()).await; assert_eq!(status, StatusCode::OK); assert!(response.success); assert_eq!(key_state(&response), KeyState::Enabled); // Enabling an already enabled key stays idempotent. - let (status, response) = execute_lifecycle(&manager, &key_id, LifecycleOperation::Enable).await; + let (status, response) = + execute_lifecycle(&manager, &key_id, LifecycleOperation::Enable, &OperationContext::internal()).await; assert_eq!(status, StatusCode::OK); assert_eq!(key_state(&response), KeyState::Enabled); } @@ -405,7 +421,8 @@ mod tests { let manager = local_manager(&temp_dir).await; let key_id = create_key(&manager, "rotate-unsupported").await; - let (status, response) = execute_lifecycle(&manager, &key_id, LifecycleOperation::Rotate).await; + let (status, response) = + execute_lifecycle(&manager, &key_id, LifecycleOperation::Rotate, &OperationContext::internal()).await; assert_eq!(status, StatusCode::NOT_IMPLEMENTED); assert_ne!(status, StatusCode::NOT_FOUND); assert!(!response.success); @@ -416,9 +433,10 @@ mod tests { ); // A disabled key must not change the picture: rotation stays rejected. - let (status, _) = execute_lifecycle(&manager, &key_id, LifecycleOperation::Disable).await; + let (status, _) = execute_lifecycle(&manager, &key_id, LifecycleOperation::Disable, &OperationContext::internal()).await; assert_eq!(status, StatusCode::OK); - let (status, response) = execute_lifecycle(&manager, &key_id, LifecycleOperation::Rotate).await; + let (status, response) = + execute_lifecycle(&manager, &key_id, LifecycleOperation::Rotate, &OperationContext::internal()).await; assert_ne!(status, StatusCode::OK); assert!(!response.success); } @@ -428,7 +446,8 @@ mod tests { let temp_dir = tempfile::tempdir().expect("temp dir"); let manager = local_manager(&temp_dir).await; - let (status, response) = execute_lifecycle(&manager, "no-such-key", LifecycleOperation::Enable).await; + let (status, response) = + execute_lifecycle(&manager, "no-such-key", LifecycleOperation::Enable, &OperationContext::internal()).await; assert_eq!(status, StatusCode::NOT_FOUND); assert!(!response.success); } diff --git a/rustfs/src/admin/handlers/kms_keys.rs b/rustfs/src/admin/handlers/kms_keys.rs index fc5317536..72f9cfc97 100644 --- a/rustfs/src/admin/handlers/kms_keys.rs +++ b/rustfs/src/admin/handlers/kms_keys.rs @@ -14,6 +14,7 @@ //! KMS key management admin API handlers +use super::kms_audit::{KmsAdminAudit, KmsAdminOperation}; use crate::admin::auth::{validate_admin_request, validate_admin_request_with_kms_key}; use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::runtime_sources::{current_kms_runtime_service_manager, current_or_init_kms_runtime_service_manager}; @@ -23,7 +24,7 @@ use base64::Engine; use hyper::{HeaderMap, Method, StatusCode}; use matchit::Params; use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; -use rustfs_kms::{KmsError, types::*}; +use rustfs_kms::{KmsAuditOperation, KmsError, types::*}; use rustfs_policy::policy::action::{Action, KmsAction}; use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; @@ -215,15 +216,21 @@ impl Operation for CreateKeyHandler { 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, - kms_create_key_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate( + validate_admin_request( + &req.headers, + &cred, + owner, + false, + kms_create_key_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await, + KmsAuditOperation::CreateKey, + None, + )?; let body = req .input @@ -258,7 +265,7 @@ impl Operation for CreateKeyHandler { policy: None, }; - match service.create_key(kms_request).await { + match service.create_key_with_context(kms_request, audit.context()).await { Ok(response) => { let api_response = CreateKeyApiResponse { key_id: response.key_id, @@ -303,17 +310,22 @@ impl Operation for DescribeKeyHandler { check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; let requested_key_id = extract_key_id(&req.uri); + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); - validate_admin_request_with_kms_key( - &req.headers, - &cred, - owner, - false, - kms_describe_key_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - requested_key_id.as_deref().unwrap_or_default(), - ) - .await?; + audit.gate( + validate_admin_request_with_kms_key( + &req.headers, + &cred, + owner, + false, + kms_describe_key_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + requested_key_id.as_deref().unwrap_or_default(), + ) + .await, + KmsAuditOperation::DescribeKey, + requested_key_id.as_deref(), + )?; let Some(key_id) = requested_key_id else { return Err(s3_error!(InvalidRequest, "missing required parameter: 'keyId'")); @@ -325,7 +337,7 @@ impl Operation for DescribeKeyHandler { let request = DescribeKeyRequest { key_id: key_id.clone() }; - match service.describe_key(request).await { + match service.describe_key_with_context(request, audit.context()).await { Ok(response) => { let api_response = DescribeKeyApiResponse { key_metadata: response.key_metadata, @@ -707,15 +719,21 @@ impl Operation for ListKeysHandler { 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, - kms_list_keys_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate( + validate_admin_request( + &req.headers, + &cred, + owner, + false, + kms_list_keys_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await, + KmsAuditOperation::ListKeys, + None, + )?; let query_params = extract_query_params(&req.uri); let limit = query_params.get("limit").and_then(|s| s.parse::().ok()).unwrap_or(100); @@ -732,7 +750,7 @@ impl Operation for ListKeysHandler { usage_filter: None, }; - match service.list_keys(request).await { + match service.list_keys_with_context(request, audit.context()).await { Ok(response) => { let api_response = ListKeysApiResponse { keys: response.keys, @@ -790,16 +808,22 @@ impl Operation for GenerateDataKeyHandler { .map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e)) }); - validate_admin_request_with_kms_key( - &req.headers, - &cred, - owner, - false, - kms_generate_data_key_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - parsed.as_ref().map(|request| request.key_id.as_str()).unwrap_or_default(), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate_admin( + validate_admin_request_with_kms_key( + &req.headers, + &cred, + owner, + false, + kms_generate_data_key_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + parsed.as_ref().map(|request| request.key_id.as_str()).unwrap_or_default(), + ) + .await, + KmsAdminOperation::GenerateDataKey, + parsed.as_ref().ok().map(|request| request.key_id.as_str()), + )?; let request = parsed?; @@ -807,13 +831,20 @@ impl Operation for GenerateDataKeyHandler { return Err(s3_error!(InternalError, "KMS service not initialized")); }; + let key_id = request.key_id.clone(); 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 { + // Audited here rather than in the KMS layer: data-key derivation goes + // through the encryption service, which has no context-aware entry + // point, so this is the only place that knows who asked. + let result = service.generate_data_key(kms_request).await; + audit.finish(KmsAdminOperation::GenerateDataKey, Some(&key_id), result.as_ref().err()); + + match result { Ok(response) => { let api_response = GenerateDataKeyApiResponse { key_id: response.key_id, @@ -858,15 +889,21 @@ impl Operation for CreateKmsKeyHandler { 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, - kms_create_key_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate( + validate_admin_request( + &req.headers, + &cred, + owner, + false, + kms_create_key_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await, + KmsAuditOperation::CreateKey, + None, + )?; let body = req .input @@ -925,7 +962,7 @@ impl Operation for CreateKmsKeyHandler { policy: None, }; - match manager.create_key(kms_request).await { + match manager.create_key_with_context(kms_request, audit.context()).await { Ok(kms_response) => { info!( event = EVENT_ADMIN_KMS_KEYS_STATE, @@ -1013,21 +1050,23 @@ impl Operation for DeleteKmsKeyHandler { // targets; the read failure is surfaced afterwards to keep AccessDenied // ahead of input errors for callers that are not authorized at all. let body = req.input.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE).await; + let scoped = body.as_ref().ok().and_then(|body| scoped_key_id(body, &req.uri)); + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); - validate_admin_request_with_kms_key( - &req.headers, - &cred, - owner, - false, - kms_delete_key_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - body.as_ref() - .ok() - .and_then(|body| scoped_key_id(body, &req.uri)) - .as_deref() - .unwrap_or_default(), - ) - .await?; + audit.gate( + validate_admin_request_with_kms_key( + &req.headers, + &cred, + owner, + false, + kms_delete_key_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + scoped.as_deref().unwrap_or_default(), + ) + .await, + KmsAuditOperation::ScheduleKeyDeletion, + scoped.as_deref(), + )?; let body = body.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?; @@ -1094,7 +1133,7 @@ impl Operation for DeleteKmsKeyHandler { force_immediate: request.force_immediate, }; - match manager.delete_key(kms_request).await { + match manager.delete_key_with_context(kms_request, audit.context()).await { Ok(kms_response) => { info!( event = EVENT_ADMIN_KMS_KEYS_STATE, @@ -1192,21 +1231,23 @@ impl Operation for CancelKmsKeyDeletionHandler { // Same ordering as the delete endpoint: the target key is resolved before // the gate, the read failure only after it. let body = req.input.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE).await; + let scoped = body.as_ref().ok().and_then(|body| scoped_key_id(body, &req.uri)); + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); - validate_admin_request_with_kms_key( - &req.headers, - &cred, - owner, - false, - kms_delete_key_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - body.as_ref() - .ok() - .and_then(|body| scoped_key_id(body, &req.uri)) - .as_deref() - .unwrap_or_default(), - ) - .await?; + audit.gate( + validate_admin_request_with_kms_key( + &req.headers, + &cred, + owner, + false, + kms_delete_key_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + scoped.as_deref().unwrap_or_default(), + ) + .await, + KmsAuditOperation::CancelKeyDeletion, + scoped.as_deref(), + )?; let body = body.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?; @@ -1262,7 +1303,7 @@ impl Operation for CancelKmsKeyDeletionHandler { key_id: request.key_id.clone(), }; - match manager.cancel_key_deletion(kms_request).await { + match manager.cancel_key_deletion_with_context(kms_request, audit.context()).await { Ok(kms_response) => { info!( event = EVENT_ADMIN_KMS_KEYS_STATE, @@ -1340,15 +1381,21 @@ impl Operation for ListKmsKeysHandler { 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, - kms_list_keys_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate( + validate_admin_request( + &req.headers, + &cred, + owner, + false, + kms_list_keys_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await, + KmsAuditOperation::ListKeys, + None, + )?; let query_params = extract_query_params(&req.uri); let limit = query_params.get("limit").and_then(|s| s.parse::().ok()).unwrap_or(100); @@ -1391,7 +1438,7 @@ impl Operation for ListKmsKeysHandler { usage_filter: None, }; - match manager.list_keys(kms_request).await { + match manager.list_keys_with_context(kms_request, audit.context()).await { Ok(kms_response) => { info!( event = EVENT_ADMIN_KMS_KEYS_STATE, @@ -1468,16 +1515,22 @@ impl Operation for DescribeKmsKeyHandler { let (cred, owner) = check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?; - validate_admin_request_with_kms_key( - &req.headers, - &cred, - owner, - false, - kms_describe_key_actions(), - req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), - params.get("key_id").unwrap_or_default(), - ) - .await?; + let audit = KmsAdminAudit::from_request(&req.extensions, &req.headers, &cred); + + audit.gate( + validate_admin_request_with_kms_key( + &req.headers, + &cred, + owner, + false, + kms_describe_key_actions(), + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + params.get("key_id").unwrap_or_default(), + ) + .await, + KmsAuditOperation::DescribeKey, + params.get("key_id"), + )?; let Some(key_id) = params.get("key_id") else { let response = DescribeKmsKeyResponse { @@ -1522,7 +1575,7 @@ impl Operation for DescribeKmsKeyHandler { key_id: key_id.to_string(), }; - match manager.describe_key(kms_request).await { + match manager.describe_key_with_context(kms_request, audit.context()).await { Ok(kms_response) => { info!( event = EVENT_ADMIN_KMS_KEYS_STATE, diff --git a/rustfs/src/admin/handlers/mod.rs b/rustfs/src/admin/handlers/mod.rs index 17347d5e3..477e44904 100644 --- a/rustfs/src/admin/handlers/mod.rs +++ b/rustfs/src/admin/handlers/mod.rs @@ -32,6 +32,7 @@ pub mod ilm_transition; pub mod inspect_archive; pub mod is_admin; pub mod kms; +pub mod kms_audit; pub mod kms_dynamic; pub mod kms_key_lifecycle; pub mod kms_keys; diff --git a/rustfs/src/init.rs b/rustfs/src/init.rs index 1d7774b40..454c8af03 100644 --- a/rustfs/src/init.rs +++ b/rustfs/src/init.rs @@ -478,6 +478,12 @@ pub async fn init_kms_system(config: &config::Config) -> std::io::Result<()> { service_manager .set_deletion_reference_checker(std::sync::Arc::new(crate::kms_deletion_gate::BucketEncryptionReferenceChecker)); + // Route KMS management records into the server's audit pipeline. Installed + // before the service can start so every service version built afterwards + // carries it; with no audit target configured the records are dropped and + // KMS operations are unaffected. + service_manager.set_audit_sink(std::sync::Arc::new(crate::admin::handlers::kms_audit::KmsAdminAuditSink)); + // If KMS is enabled in configuration, configure and start the service if config.kms_enable { info!( diff --git a/rustfs/src/storage/access.rs b/rustfs/src/storage/access.rs index cfa03f88d..12eb016f5 100644 --- a/rustfs/src/storage/access.rs +++ b/rustfs/src/storage/access.rs @@ -78,10 +78,16 @@ fn ext_req_info_mut(ext: &mut http::Extensions) -> S3Result<&mut ReqInfo> { /// Extract the canonical `RequestContext` from a request, checking both /// the request extensions directly and the `ReqInfo.request_context` field. pub(crate) fn request_context_from_req(req: &S3Request) -> Option { - req.extensions + request_context_from_extensions(&req.extensions) +} + +/// Same lookup against the extensions alone, for callers that have already +/// moved a field out of the request and can no longer borrow it whole. +pub(crate) fn request_context_from_extensions(extensions: &http::Extensions) -> Option { + extensions .get::() .cloned() - .or_else(|| req.extensions.get::().and_then(|ri| ri.request_context.clone())) + .or_else(|| extensions.get::().and_then(|ri| ri.request_context.clone())) } #[derive(Clone, Debug)] diff --git a/scripts/check_logging_guardrails.sh b/scripts/check_logging_guardrails.sh index 5f22c1f65..12fd1fa4d 100755 --- a/scripts/check_logging_guardrails.sh +++ b/scripts/check_logging_guardrails.sh @@ -14,7 +14,10 @@ checked_files=( "rustfs/src/admin/router.rs" "rustfs/src/admin/handlers/table_catalog.rs" "rustfs/src/admin/handlers/service_account.rs" + "rustfs/src/admin/handlers/kms_audit.rs" "rustfs/src/admin/handlers/kms_dynamic.rs" + "rustfs/src/admin/handlers/kms_keys.rs" + "rustfs/src/admin/handlers/kms_key_lifecycle.rs" "rustfs/src/admin/handlers/site_replication.rs" "rustfs/src/admin/handlers/group.rs" "rustfs/src/admin/handlers/quota.rs"