From 5ea9a1fd8f3c1b39154b27fd419cde76abc68f3b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Wed, 22 Jul 2026 15:19:24 +0800 Subject: [PATCH] fix(sse): separate SSE-S3 and KMS key providers --- crates/e2e_test/src/multipart_auth_test.rs | 28 +- crates/ecstore/src/api/mod.rs | 4 + crates/ecstore/src/lib.rs | 1 + crates/ecstore/src/object_api/readers.rs | 130 +++- crates/ecstore/src/sse/mod.rs | 46 ++ .../tests/minio_generated_read_test.rs | 4 +- crates/kms/Cargo.toml | 2 +- crates/kms/src/encryption/dek.rs | 35 + crates/kms/src/encryption/mod.rs | 2 +- crates/kms/src/lib.rs | 5 + crates/kms/src/managed_context.rs | 110 +++ crates/utils/src/http/header_compat.rs | 56 ++ docs/architecture/compat-cleanup-register.md | 1 + rustfs/src/storage/sse.rs | 713 +++++++++++++----- rustfs/src/storage/storage_api.rs | 4 + 15 files changed, 895 insertions(+), 246 deletions(-) create mode 100644 crates/ecstore/src/sse/mod.rs create mode 100644 crates/kms/src/managed_context.rs diff --git a/crates/e2e_test/src/multipart_auth_test.rs b/crates/e2e_test/src/multipart_auth_test.rs index feabb0ead..7366c287d 100644 --- a/crates/e2e_test/src/multipart_auth_test.rs +++ b/crates/e2e_test/src/multipart_auth_test.rs @@ -1068,7 +1068,8 @@ async fn test_anonymous_post_object_uses_bucket_default_sse_s3() -> Result<(), B #[tokio::test] #[serial] -async fn test_anonymous_post_object_uses_bucket_default_sse_kms() -> Result<(), Box> { +async fn test_anonymous_post_object_bucket_default_sse_kms_requires_kms() -> Result<(), Box> +{ init_logging(); let mut env = RustFSTestEnvironment::new().await?; @@ -1128,14 +1129,23 @@ async fn test_anonymous_post_object_uses_bucket_default_sse_kms() -> Result<(), .send() .await?; - assert_eq!(post_resp.status(), reqwest::StatusCode::NO_CONTENT); - - let head = admin_client.head_object().bucket(bucket).key(object_key).send().await?; - assert_eq!(head.server_side_encryption().map(|value| value.as_str()), Some("aws:kms")); - - let uploaded = admin_client.get_object().bucket(bucket).key(object_key).send().await?; - let uploaded = uploaded.body.collect().await?.into_bytes(); - assert_eq!(uploaded.as_ref(), expected_body.as_slice()); + let status = post_resp.status(); + let response_body = post_resp.text().await?; + assert_eq!(status, reqwest::StatusCode::BAD_REQUEST); + assert!( + response_body.contains("configured KMS service"), + "unexpected response body: {response_body}" + ); + assert!( + admin_client + .head_object() + .bucket(bucket) + .key(object_key) + .send() + .await + .is_err(), + "failed SSE-KMS POST must not create an object" + ); Ok(()) } diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index c7f186c91..dccde0e82 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -398,6 +398,10 @@ pub mod set_disk { pub use crate::set_disk::{DEFAULT_READ_BUFFER_SIZE, SetDisks, get_lock_acquire_timeout, is_valid_storage_class}; } +pub mod sse { + pub use crate::sse::{ManagedDekProvider, ManagedSseScheme, managed_dek_provider}; +} + pub mod store_list { pub use crate::store::list_objects::{ListPathOptions, max_keys_plus_one}; } diff --git a/crates/ecstore/src/lib.rs b/crates/ecstore/src/lib.rs index a30223c03..4191eea9c 100644 --- a/crates/ecstore/src/lib.rs +++ b/crates/ecstore/src/lib.rs @@ -49,6 +49,7 @@ mod object_api; mod runtime; mod services; mod set_disk; +mod sse; mod storage_api_contracts; mod store; diff --git a/crates/ecstore/src/object_api/readers.rs b/crates/ecstore/src/object_api/readers.rs index 586ce0fa6..29923f0ca 100644 --- a/crates/ecstore/src/object_api/readers.rs +++ b/crates/ecstore/src/object_api/readers.rs @@ -13,6 +13,7 @@ // limitations under the License. use super::*; +use crate::sse::{ManagedDekProvider, ManagedSseScheme, managed_dek_provider as classify_managed_dek_provider}; #[cfg(feature = "rio-v2")] use aes_gcm::aead::Payload; use aes_gcm::{ @@ -26,7 +27,11 @@ use chacha20poly1305::ChaCha20Poly1305; use hmac::{Hmac, Mac}; use md5::{Digest, Md5}; use rustfs_kms::types::ObjectEncryptionContext; -use rustfs_utils::http::{SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER}; +#[cfg(feature = "rio-v2")] +use rustfs_kms::{MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER, RUSTFS_ENCRYPTION_CONTEXT_HEADER, decode_managed_kms_context}; +use rustfs_utils::http::{ + AMZ_SERVER_SIDE_ENCRYPTION, SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER, get_consistent_metadata_value, +}; use rustfs_utils::path::path_join_buf; #[cfg(feature = "rio-v2")] use serde::Deserialize; @@ -39,11 +44,11 @@ use crate::io_support::rio::Index; const INTERNAL_ENCRYPTION_KEY_ID_HEADER: &str = "x-rustfs-encryption-key-id"; const INTERNAL_ENCRYPTION_KEY_HEADER: &str = "x-rustfs-encryption-key"; -const INTERNAL_ENCRYPTION_CONTEXT_HEADER: &str = "x-rustfs-encryption-context"; const INTERNAL_ENCRYPTION_IV_HEADER: &str = "x-rustfs-encryption-iv"; const INTERNAL_ENCRYPTION_ORIGINAL_SIZE_HEADER: &str = "x-rustfs-encryption-original-size"; const SSEC_ORIGINAL_SIZE_HEADER: &str = "x-amz-server-side-encryption-customer-original-size"; const DEFAULT_SSE_ALGORITHM: &str = "AES256"; +const SSE_KMS_ALGORITHM: &str = "aws:kms"; #[cfg(feature = "rio-v2")] const DARE_PAYLOAD_SIZE: i64 = 64 * 1024; #[cfg(feature = "rio-v2")] @@ -60,8 +65,6 @@ const MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER: &str = "X-Minio-Internal-Serv #[cfg(feature = "rio-v2")] const MINIO_INTERNAL_ENCRYPTION_KMS_DATA_KEY_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-S3-Kms-Sealed-Key"; #[cfg(feature = "rio-v2")] -const MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-Context"; -#[cfg(feature = "rio-v2")] const MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-Sealed-Key"; #[cfg(feature = "rio-v2")] const MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM: &str = "DAREv2-HMAC-SHA256"; @@ -1381,31 +1384,30 @@ async fn resolve_managed_material(bucket: &str, object: &str, metadata: &HashMap let kms_key_id = metadata_get(&normalized_metadata, INTERNAL_ENCRYPTION_KEY_ID_HEADER).unwrap_or("default"); #[cfg(feature = "rio-v2")] - let kms_context = metadata_get(&normalized_metadata, INTERNAL_ENCRYPTION_CONTEXT_HEADER) - .map(|value| { - serde_json::from_str::>(value) - .map_err(|e| Error::other(format!("failed to parse managed KMS context: {e}"))) - }) - .transpose()?; + let kms_context = decode_managed_kms_context(metadata).map_err(|err| Error::other(err.to_string()))?; #[cfg(not(feature = "rio-v2"))] let kms_context: Option> = None; let object_context = build_object_encryption_context(bucket, object, kms_context.as_ref()); - let decrypted_key = if let Some(service) = crate::runtime::sources::object_encryption_service().await { - #[cfg(feature = "rio-v2")] - let data_key = if is_legacy_rustfs_managed_metadata(&normalized_metadata) { - service.decrypt_legacy_data_key(&encrypted_dek).await - } else { - service.decrypt_data_key(&encrypted_dek, &object_context).await - }; - #[cfg(not(feature = "rio-v2"))] - let data_key = service.decrypt_data_key(&encrypted_dek, &object_context).await; + let decrypted_key = match managed_dek_provider(metadata, &encrypted_dek)? { + ManagedDekProvider::LocalSseS3 => decrypt_local_sse_dek(&encrypted_dek, kms_key_id, &object_context)?, + ManagedDekProvider::Kms => { + let service = crate::runtime::sources::object_encryption_service() + .await + .ok_or_else(|| Error::other("KMS encryption service is required to decrypt this object"))?; + #[cfg(feature = "rio-v2")] + let data_key = if is_legacy_rustfs_managed_metadata(&normalized_metadata) { + service.decrypt_legacy_data_key(&encrypted_dek).await + } else { + service.decrypt_data_key(&encrypted_dek, &object_context).await + }; + #[cfg(not(feature = "rio-v2"))] + let data_key = service.decrypt_data_key(&encrypted_dek, &object_context).await; - data_key - .map_err(|e| Error::other(format!("failed to decrypt managed data key: {e}")))? - .plaintext_key - } else { - decrypt_local_sse_dek(&encrypted_dek, kms_key_id, &object_context)? + data_key + .map_err(|e| Error::other(format!("failed to decrypt managed data key: {e}")))? + .plaintext_key + } }; #[cfg(feature = "rio-v2")] @@ -1436,6 +1438,19 @@ async fn resolve_managed_material(bucket: &str, object: &str, metadata: &HashMap }) } +fn managed_dek_provider(metadata: &HashMap, encrypted_dek: &[u8]) -> Result { + let algorithm = get_consistent_metadata_value(metadata, AMZ_SERVER_SIDE_ENCRYPTION) + .map_err(|_| Error::other(format!("conflicting managed encryption metadata for {AMZ_SERVER_SIDE_ENCRYPTION}")))?; + let scheme = match algorithm { + Some(SSE_KMS_ALGORITHM) => ManagedSseScheme::SseKms, + Some(DEFAULT_SSE_ALGORITHM) | None => ManagedSseScheme::SseS3, + Some(algorithm) => return Err(Error::other(format!("unsupported stored server-side encryption {algorithm}"))), + }; + // RUSTFS_COMPAT_TODO(rustfs-5063): Keep legacy SSE-S3 KMS envelopes readable. Remove after SSE-S3 migration rewraps every referenced legacy DEK. + let has_kms_envelope = rustfs_kms::is_data_key_envelope(encrypted_dek); + Ok(classify_managed_dek_provider(scheme, has_kms_envelope)) +} + fn normalize_managed_metadata(metadata: &HashMap) -> HashMap { #[cfg(feature = "rio-v2")] { @@ -1460,13 +1475,13 @@ fn normalize_managed_metadata(metadata: &HashMap) -> HashMap>(&decoded) && let Ok(encoded) = serde_json::to_string(&context) { - normalized.insert(INTERNAL_ENCRYPTION_CONTEXT_HEADER.to_string(), encoded); + normalized.insert(RUSTFS_ENCRYPTION_CONTEXT_HEADER.to_string(), encoded); } normalized @@ -1624,15 +1639,21 @@ fn marshal_minio_kms_context(context: &HashMap) -> Vec { } fn local_sse_master_key() -> Result<[u8; 32]> { + #[cfg(test)] if let Some(key) = decode_master_key_env("__RUSTFS_SSE_SIMPLE_CMK")? { return Ok(key); } if let Some(key) = decode_master_key_env("RUSTFS_SSE_S3_MASTER_KEY")? { + if key == [0u8; 32] { + return Err(Error::other("RUSTFS_SSE_S3_MASTER_KEY must not be an all-zero key")); + } return Ok(key); } - Ok([0u8; 32]) + Err(Error::other( + "SSE-S3 requires RUSTFS_SSE_S3_MASTER_KEY to decrypt locally managed objects", + )) } fn decode_master_key_env(name: &str) -> Result> { @@ -2108,6 +2129,57 @@ mod tests { format!("{}:{}", BASE64_STANDARD.encode(nonce), BASE64_STANDARD.encode(ciphertext)) } + #[test] + fn managed_dek_provider_routes_by_algorithm_and_persisted_envelope() { + let local_dek = encrypt_managed_dek_for_test([0x24; 32], [0x42; 32]); + let kms_dek = serde_json::to_vec(&serde_json::json!({ + "key_id": "legacy-data-key", + "master_key_id": "legacy-master-key", + "key_spec": "AES_256", + "encrypted_key": [1, 2, 3, 4], + "nonce": [5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16], + "encryption_context": {}, + "created_at": "2024-01-01T00:00:00+00:00" + })) + .expect("legacy KMS envelope should serialize"); + + assert_eq!( + managed_dek_provider( + &HashMap::from([(AMZ_SERVER_SIDE_ENCRYPTION.to_string(), DEFAULT_SSE_ALGORITHM.to_string())]), + local_dek.as_bytes(), + ) + .expect("SSE-S3 local DEK should be classified"), + ManagedDekProvider::LocalSseS3 + ); + assert_eq!( + managed_dek_provider( + &HashMap::from([(AMZ_SERVER_SIDE_ENCRYPTION.to_string(), SSE_KMS_ALGORITHM.to_string())]), + local_dek.as_bytes(), + ) + .expect("SSE-KMS should be classified by its stored algorithm"), + ManagedDekProvider::Kms + ); + assert_eq!( + managed_dek_provider( + &HashMap::from([(AMZ_SERVER_SIDE_ENCRYPTION.to_string(), DEFAULT_SSE_ALGORITHM.to_string())]), + &kms_dek, + ) + .expect("legacy SSE-S3 KMS envelope should be classified"), + ManagedDekProvider::Kms + ); + assert_eq!( + managed_dek_provider(&HashMap::new(), &kms_dek).expect("legacy KMS envelope without algorithm should be classified"), + ManagedDekProvider::Kms + ); + assert!( + managed_dek_provider( + &HashMap::from([(AMZ_SERVER_SIDE_ENCRYPTION.to_string(), "unsupported".to_string())]), + local_dek.as_bytes(), + ) + .is_err() + ); + } + #[cfg(feature = "rio-v2")] fn seal_managed_s3_object_key_for_test( bucket: &str, @@ -2622,12 +2694,12 @@ mod tests { async_with_vars( [ ("__RUSTFS_SSE_SIMPLE_CMK", None::), - ("RUSTFS_SSE_S3_MASTER_KEY", Some(BASE64_STANDARD.encode([0u8; 32]))), + ("RUSTFS_SSE_S3_MASTER_KEY", Some(BASE64_STANDARD.encode([0x33u8; 32]))), ], async { let plaintext = b"managed-local-fallback".to_vec(); let data_key = [0x22; 32]; - let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0u8; 32]); + let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0x33; 32]); let bucket = "bucket"; let object = "managed-local-fallback"; diff --git a/crates/ecstore/src/sse/mod.rs b/crates/ecstore/src/sse/mod.rs new file mode 100644 index 000000000..3467fa2af --- /dev/null +++ b/crates/ecstore/src/sse/mod.rs @@ -0,0 +1,46 @@ +// 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. + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ManagedSseScheme { + SseS3, + SseKms, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ManagedDekProvider { + LocalSseS3, + Kms, +} + +pub fn managed_dek_provider(scheme: ManagedSseScheme, has_kms_envelope: bool) -> ManagedDekProvider { + match scheme { + ManagedSseScheme::SseKms => ManagedDekProvider::Kms, + ManagedSseScheme::SseS3 if has_kms_envelope => ManagedDekProvider::Kms, + ManagedSseScheme::SseS3 => ManagedDekProvider::LocalSseS3, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn managed_sse_routing_is_determined_by_scheme_and_envelope() { + assert_eq!(managed_dek_provider(ManagedSseScheme::SseS3, false), ManagedDekProvider::LocalSseS3); + assert_eq!(managed_dek_provider(ManagedSseScheme::SseS3, true), ManagedDekProvider::Kms); + assert_eq!(managed_dek_provider(ManagedSseScheme::SseKms, false), ManagedDekProvider::Kms); + assert_eq!(managed_dek_provider(ManagedSseScheme::SseKms, true), ManagedDekProvider::Kms); + } +} diff --git a/crates/ecstore/tests/minio_generated_read_test.rs b/crates/ecstore/tests/minio_generated_read_test.rs index 875aa9be9..621646a81 100644 --- a/crates/ecstore/tests/minio_generated_read_test.rs +++ b/crates/ecstore/tests/minio_generated_read_test.rs @@ -133,8 +133,8 @@ async fn read_fixture_plaintext(encrypted: Vec, object_info: ObjectInfo, kms async_with_vars( [ - ("__RUSTFS_SSE_SIMPLE_CMK", Some(kms_key_b64)), - ("RUSTFS_SSE_S3_MASTER_KEY", None::), + ("RUSTFS_SSE_S3_MASTER_KEY", Some(kms_key_b64)), + ("__RUSTFS_SSE_SIMPLE_CMK", None::), ], async move { let (mut reader, offset, length) = GetObjectReader::new( diff --git a/crates/kms/Cargo.toml b/crates/kms/Cargo.toml index aada8547b..e098bc51b 100644 --- a/crates/kms/Cargo.toml +++ b/crates/kms/Cargo.toml @@ -56,7 +56,7 @@ moka = { workspace = true, features = ["future"] } # Additional dependencies md5 = { workspace = true } arc-swap = { workspace = true } -rustfs-utils = { workspace = true } +rustfs-utils = { workspace = true, features = ["http"] } rustfs-security-governance = { workspace = true } # HTTP client for Vault diff --git a/crates/kms/src/encryption/dek.rs b/crates/kms/src/encryption/dek.rs index a2753b65d..3af6e410e 100644 --- a/crates/kms/src/encryption/dek.rs +++ b/crates/kms/src/encryption/dek.rs @@ -44,6 +44,23 @@ pub struct DataKeyEnvelope { pub created_at: Zoned, } +/// Return whether bytes contain a complete RustFS KMS data-key envelope. +pub fn is_data_key_envelope(ciphertext: &[u8]) -> bool { + const MAX_ENVELOPE_SIZE: usize = 64 * 1024; + + if ciphertext.is_empty() || ciphertext.len() > MAX_ENVELOPE_SIZE { + return false; + } + + serde_json::from_slice::(ciphertext).is_ok_and(|envelope| { + !envelope.key_id.trim().is_empty() + && !envelope.master_key_id.trim().is_empty() + && envelope.key_spec == "AES_256" + && !envelope.encrypted_key.is_empty() + && (envelope.nonce.is_empty() || envelope.nonce.len() == 12) + }) +} + /// Trait for encrypting and decrypting data encryption keys (DEK) /// /// This trait abstracts the encryption operations used to protect @@ -293,6 +310,24 @@ mod tests { assert_eq!(deserialized.key_id, envelope.key_id); assert_eq!(deserialized.master_key_id, envelope.master_key_id); assert_eq!(deserialized.encrypted_key, envelope.encrypted_key); + assert!(is_data_key_envelope(&serialized)); + assert!(!is_data_key_envelope(b"not-a-kms-envelope")); + let mut boundary_envelope = serialized; + boundary_envelope.resize(64 * 1024, b' '); + assert!(is_data_key_envelope(&boundary_envelope)); + boundary_envelope.push(b' '); + assert!(!is_data_key_envelope(&boundary_envelope)); + + let mut invalid = serde_json::to_value(&envelope).expect("Envelope should convert to JSON"); + invalid["key_spec"] = serde_json::Value::String("AES_128".to_string()); + assert!(!is_data_key_envelope( + &serde_json::to_vec(&invalid).expect("Invalid envelope should serialize") + )); + invalid["key_spec"] = serde_json::Value::String("AES_256".to_string()); + invalid["unknown"] = serde_json::Value::Bool(true); + assert!(is_data_key_envelope( + &serde_json::to_vec(&invalid).expect("Envelope with unknown field should serialize") + )); } #[tokio::test] diff --git a/crates/kms/src/encryption/mod.rs b/crates/kms/src/encryption/mod.rs index 3c1ee8711..d18a1df34 100644 --- a/crates/kms/src/encryption/mod.rs +++ b/crates/kms/src/encryption/mod.rs @@ -17,4 +17,4 @@ pub mod ciphers; pub mod dek; -pub use dek::{AesDekCrypto, DataKeyEnvelope, DekCrypto, generate_key_material}; +pub use dek::{AesDekCrypto, DataKeyEnvelope, DekCrypto, generate_key_material, is_data_key_envelope}; diff --git a/crates/kms/src/lib.rs b/crates/kms/src/lib.rs index 29cfd25ac..b273bbf22 100644 --- a/crates/kms/src/lib.rs +++ b/crates/kms/src/lib.rs @@ -70,6 +70,7 @@ mod cache; pub mod config; mod encryption; mod error; +mod managed_context; pub mod manager; pub mod service; pub mod service_manager; @@ -83,7 +84,11 @@ pub use api_types::{ TagKeyRequest, TagKeyResponse, UntagKeyRequest, UntagKeyResponse, UpdateKeyDescriptionRequest, UpdateKeyDescriptionResponse, }; pub use config::*; +pub use encryption::is_data_key_envelope; pub use error::{KmsError, Result}; +pub use managed_context::{ + MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER, RUSTFS_ENCRYPTION_CONTEXT_HEADER, decode_managed_kms_context, +}; pub use manager::KmsManager; pub use service::{DataKey, ObjectEncryptionService}; pub use service_manager::{ diff --git a/crates/kms/src/managed_context.rs b/crates/kms/src/managed_context.rs new file mode 100644 index 000000000..9005e181a --- /dev/null +++ b/crates/kms/src/managed_context.rs @@ -0,0 +1,110 @@ +// 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. + +use crate::{KmsError, Result}; +use base64::{Engine, engine::general_purpose::STANDARD as BASE64_STANDARD}; +use rustfs_utils::http::get_consistent_metadata_value; +use std::collections::HashMap; + +pub const RUSTFS_ENCRYPTION_CONTEXT_HEADER: &str = "x-rustfs-encryption-context"; +pub const MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-Context"; + +pub fn decode_managed_kms_context(metadata: &HashMap) -> Result>> { + let minio_context = consistent_value(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER)? + .map(|context| { + let decoded = BASE64_STANDARD + .decode(context) + .map_err(|err| KmsError::serialization_error(format!("Failed to decode MinIO KMS context: {err}")))?; + serde_json::from_slice(&decoded) + .map_err(|err| KmsError::serialization_error(format!("Failed to parse MinIO KMS context: {err}"))) + }) + .transpose()?; + let rustfs_context = consistent_value(metadata, RUSTFS_ENCRYPTION_CONTEXT_HEADER)? + .map(|context| { + serde_json::from_str(context) + .map_err(|err| KmsError::serialization_error(format!("Failed to parse RustFS KMS context: {err}"))) + }) + .transpose()?; + + match (minio_context, rustfs_context) { + (Some(minio), Some(rustfs)) if minio != rustfs => { + Err(KmsError::context_mismatch("Conflicting RustFS and MinIO KMS contexts")) + } + (Some(context), _) | (_, Some(context)) => Ok(Some(context)), + (None, None) => Ok(None), + } +} + +fn consistent_value<'a>(metadata: &'a HashMap, name: &str) -> Result> { + get_consistent_metadata_value(metadata, name) + .map_err(|_| KmsError::validation_error(format!("Conflicting managed encryption metadata for {name}"))) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn decode_context_accepts_compatible_headers_and_rejects_conflicts() { + let expected = HashMap::from([("tenant".to_string(), "alpha".to_string())]); + let metadata = HashMap::from([ + ( + RUSTFS_ENCRYPTION_CONTEXT_HEADER.to_string(), + serde_json::to_string(&expected).expect("RustFS KMS context should serialize"), + ), + ( + MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER.to_string(), + BASE64_STANDARD.encode(serde_json::to_vec(&expected).expect("MinIO KMS context should serialize")), + ), + ]); + assert_eq!( + decode_managed_kms_context(&metadata).expect("matching KMS contexts should parse"), + Some(expected.clone()) + ); + assert_eq!( + decode_managed_kms_context(&HashMap::from([( + RUSTFS_ENCRYPTION_CONTEXT_HEADER.to_string(), + serde_json::to_string(&expected).expect("legacy RustFS KMS context should serialize"), + )])) + .expect("legacy RustFS KMS context should parse"), + Some(expected) + ); + + let conflicting = HashMap::from([ + ( + RUSTFS_ENCRYPTION_CONTEXT_HEADER.to_string(), + serde_json::to_string(&HashMap::from([("tenant", "alpha")])).expect("RustFS KMS context should serialize"), + ), + ( + MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER.to_string(), + BASE64_STANDARD.encode( + serde_json::to_vec(&HashMap::from([("tenant", "beta")])).expect("MinIO KMS context should serialize"), + ), + ), + ]); + assert!(decode_managed_kms_context(&conflicting).is_err()); + + assert!( + decode_managed_kms_context(&HashMap::from([( + MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER.to_string(), + "not-base64".to_string(), + )])) + .is_err() + ); + assert!( + decode_managed_kms_context(&HashMap::from([(RUSTFS_ENCRYPTION_CONTEXT_HEADER.to_string(), "not-json".to_string(),)])) + .is_err() + ); + } +} diff --git a/crates/utils/src/http/header_compat.rs b/crates/utils/src/http/header_compat.rs index 6dda300e6..394bf2764 100644 --- a/crates/utils/src/http/header_compat.rs +++ b/crates/utils/src/http/header_compat.rs @@ -20,12 +20,26 @@ use http::{HeaderMap, HeaderValue}; use std::borrow::Cow; +use std::collections::HashMap; +use std::error::Error; +use std::fmt::{Display, Formatter}; const RUSTFS_PREFIX: &str = "x-rustfs-"; const MINIO_PREFIX: &str = "x-minio-"; const MINIO_ENCRYPTION_PREFIX: &str = "x-minio-encryption-"; const RUSTFS_ENCRYPTION_PREFIX: &str = "x-rustfs-encryption-"; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ConflictingMetadataValue; + +impl Display for ConflictingMetadataValue { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + f.write_str("metadata contains conflicting values for the same case-insensitive key") + } +} + +impl Error for ConflictingMetadataValue {} + // Suffix constants (part after x-rustfs- or x-minio-). Use with get_header/insert_header. pub const SUFFIX_FORCE_DELETE: &str = "force-delete"; pub const SUFFIX_INCLUDE_DELETED: &str = "include-deleted"; @@ -84,6 +98,24 @@ pub fn get_header_map(map: &std::collections::HashMap, suffix: & map.get(&rk).cloned().or_else(|| map.get(&mk).cloned()) } +/// Returns the value when every case-insensitive occurrence of `name` agrees. +pub fn get_consistent_metadata_value<'a>( + metadata: &'a HashMap, + name: &str, +) -> Result, ConflictingMetadataValue> { + let mut value = None; + for (candidate, candidate_value) in metadata { + if !candidate.eq_ignore_ascii_case(name) { + continue; + } + if value.is_some_and(|existing| existing != candidate_value) { + return Err(ConflictingMetadataValue); + } + value = Some(candidate_value.as_str()); + } + Ok(value) +} + /// Insert into HashMap with both x-rustfs-{suffix} and x-minio-{suffix}. pub fn insert_header_map(map: &mut std::collections::HashMap, suffix: &str, value: impl Into) { let v = value.into(); @@ -120,4 +152,28 @@ mod tests { headers2.insert("X-Rustfs-Force-Delete", HeaderValue::from_static("true")); assert_eq!(get_header(&headers2, SUFFIX_FORCE_DELETE).as_deref(), Some("true")); } + + #[test] + fn test_get_consistent_metadata_value() { + let mut metadata = HashMap::new(); + assert_eq!( + get_consistent_metadata_value(&metadata, "x-test").expect("missing value should be valid"), + None + ); + + metadata.insert("X-Test".to_string(), "value".to_string()); + assert_eq!( + get_consistent_metadata_value(&metadata, "x-test").expect("single value should be valid"), + Some("value") + ); + + metadata.insert("x-test".to_string(), "value".to_string()); + assert_eq!( + get_consistent_metadata_value(&metadata, "x-test").expect("matching values should be valid"), + Some("value") + ); + + metadata.insert("x-TEST".to_string(), "other".to_string()); + assert!(get_consistent_metadata_value(&metadata, "x-test").is_err()); + } } diff --git a/docs/architecture/compat-cleanup-register.md b/docs/architecture/compat-cleanup-register.md index 11da713d5..c7a67f88e 100644 --- a/docs/architecture/compat-cleanup-register.md +++ b/docs/architecture/compat-cleanup-register.md @@ -12,6 +12,7 @@ for later deletion. ## Open Items +- `rustfs-5063` legacy SSE-S3 KMS envelopes: releases before provider routing was separated could wrap an AES256 object's DEK with KMS, so readers identify the persisted KMS envelope and retain KMS decryption for those objects. Remove this compatibility path after the SSE-S3 migration has rewrapped every referenced legacy DEK with the local SSE-S3 key provider. - `#4648` walk-dir stream completion capability: old clients can append fallback output to an already-used metacache writer after a terminal body error, so servers emit terminal walk errors only to clients that sign the `walk_dir_stream_completion=error-v1` query capability and its request-body digest. Remove the legacy clean-EOF path after the minimum supported RustFS peer version always advertises this capability. - `heal-rpc-auth-v2` internode gRPC authentication: servers temporarily accept legacy prefix signatures so old peers remain available during rolling upgrades. Remove the legacy fallback after the minimum supported RustFS peer version sends v2 authentication on every internode gRPC request. - `heal-status-rpc-v1` node heal status capability: new peers treat an unimplemented BackgroundHealStatus RPC as an explicitly incomplete rolling-upgrade response. Remove the fallback after the minimum supported RustFS peer version implements BackgroundHealStatus. diff --git a/rustfs/src/storage/sse.rs b/rustfs/src/storage/sse.rs index 3c438a672..faac3f765 100644 --- a/rustfs/src/storage/sse.rs +++ b/rustfs/src/storage/sse.rs @@ -27,8 +27,8 @@ //! - `sse_decryption()` - Unified decryption entry point //! //! ### Managed SSE (SSE-S3 / SSE-KMS) -//! - Keys are managed by the server-side KMS service -//! - Data keys are generated and encrypted by KMS +//! - SSE-S3 data keys are protected by the local SSE-S3 master key +//! - SSE-KMS data keys are generated and protected by KMS //! - Encryption metadata is stored in object metadata //! //! ### Customer-Provided Keys (SSE-C) @@ -87,8 +87,14 @@ use http::{HeaderMap, HeaderValue}; use rand::Rng; #[cfg(feature = "rio-v2")] use rand::RngExt; -use rustfs_kms::{DataKey, types::ObjectEncryptionContext}; -use rustfs_utils::get_env_opt_str; +use rustfs_kms::{ + DataKey, MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER, RUSTFS_ENCRYPTION_CONTEXT_HEADER, decode_managed_kms_context, + types::ObjectEncryptionContext, +}; +use rustfs_utils::{ + get_env_opt_str, + http::{AMZ_SERVER_SIDE_ENCRYPTION, get_consistent_metadata_value}, +}; use s3s::S3ErrorCode; use s3s::dto::ServerSideEncryption; #[cfg(feature = "rio-v2")] @@ -97,6 +103,8 @@ use std::collections::HashMap; use std::sync::{Arc, LazyLock, RwLock}; use tracing::{debug, error}; +use super::storage_api::ecstore_sse::{ManagedDekProvider, ManagedSseScheme, managed_dek_provider}; + const LOG_COMPONENT_STORAGE: &str = "storage"; const LOG_SUBSYSTEM_SSE: &str = "sse"; @@ -114,7 +122,6 @@ const MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER: &str = "X-Minio-Internal-S const MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-Kms-Sealed-Key"; const MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-S3-Kms-Key-Id"; const MINIO_INTERNAL_ENCRYPTION_KMS_DATA_KEY_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-S3-Kms-Sealed-Key"; -const MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-Context"; #[cfg(feature = "rio-v2")] const MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM: &str = "DAREv2-HMAC-SHA256"; #[cfg(feature = "rio-v2")] @@ -812,18 +819,6 @@ fn encode_minio_kms_context(context: &HashMap) -> Result) -> Result>, ApiError> { - let Some(encoded) = metadata.get(MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER) else { - return Ok(None); - }; - let decoded = BASE64_STANDARD - .decode(encoded) - .map_err(|e| ApiError::from(StorageError::other(format!("Failed to decode MinIO KMS context: {e}"))))?; - serde_json::from_slice(&decoded) - .map(Some) - .map_err(|e| ApiError::from(StorageError::other(format!("Failed to parse MinIO KMS context: {e}")))) -} - #[cfg(feature = "rio-v2")] fn is_supported_sealed_object_key_cipher(cipher: u8) -> bool { matches!(cipher, DARE_CIPHER_AES_256_GCM | DARE_CIPHER_CHACHA20_POLY1305) @@ -1096,7 +1091,7 @@ pub fn encryption_material_to_metadata(material: &EncryptionMaterial) -> Result< && !kms_context.is_empty() { if let Ok(serialized) = serde_json::to_string(kms_context) { - metadata.insert("x-rustfs-encryption-context".to_string(), serialized); + metadata.insert(RUSTFS_ENCRYPTION_CONTEXT_HEADER.to_string(), serialized); } if matches!(material.sse_type, SSEType::SseKms) && let Ok(encoded) = encode_minio_kms_context(kms_context) @@ -1563,41 +1558,31 @@ async fn apply_managed_encryption_material( ssekms_context: Option>, content_size: i64, ) -> Result { - if !is_managed_sse(&server_side_encryption) { - return Err(ApiError::from(StorageError::other(format!( - "Unsupported server-side encryption: {}", - server_side_encryption.as_str() - )))); - } - let encryption_type = match server_side_encryption.as_str() { - "AES256" => SSEType::SseS3, - "aws:kms" => SSEType::SseKms, - _ => SSEType::SseS3, + ServerSideEncryption::AES256 => SSEType::SseS3, + ServerSideEncryption::AWS_KMS => SSEType::SseKms, + _ => { + return Err(ApiError::from(StorageError::other(format!( + "Unsupported server-side encryption: {}", + server_side_encryption.as_str() + )))); + } }; - // Determine KMS key ID to use for internal key wrapping. - let mut kms_key_candidate = kms_key_id.clone(); - if kms_key_candidate.is_none() { - // Try to get default key from KMS service (if available) - if let Some(service) = runtime_sources::current_encryption_service().await { - kms_key_candidate = service.get_default_key_id().cloned(); + let provider = get_sse_dek_provider(encryption_type).await?; + let kms_key_to_use = match encryption_type { + SSEType::SseS3 => "default".to_string(), + SSEType::SseKms => match kms_key_id { + Some(kms_key_id) => kms_key_id, + None => provider + .default_kms_key_id() + .ok_or_else(|| sse_not_configured("SSE-KMS requires a configured KMS key"))?, + }, + SSEType::SseC => { + return Err(ApiError::from(StorageError::other("SSE-C cannot use managed encryption material"))); } - } - - let kms_key_to_use = match (encryption_type, kms_key_candidate.clone()) { - (SSEType::SseS3, Some(kms_key_id)) => kms_key_id, - (SSEType::SseS3, None) => "default".to_string(), - (SSEType::SseKms, Some(kms_key_id)) => kms_key_id, - (SSEType::SseKms, None) => { - return Err(ApiError::from(StorageError::other( - "No KMS key available for managed server-side encryption (required for SSE-KMS)", - ))); - } - _ => unreachable!("managed SSE branch only supports SSE-S3 or SSE-KMS"), }; - let provider = get_sse_dek_provider().await?; let object_context = build_object_encryption_context(bucket, key, ssekms_context.as_ref()); let (data_key, encrypted_data_key) = provider .generate_sse_dek(&object_context, &kms_key_to_use) @@ -1637,20 +1622,31 @@ async fn apply_managed_decryption_material( key: &str, metadata: &HashMap, ) -> Result, ApiError> { - #[cfg(not(feature = "rio-v2"))] - let _ = (bucket, key); - if !contains_managed_encryption_metadata(metadata) || !metadata.contains_key("x-amz-server-side-encryption") { + if !contains_managed_encryption_metadata(metadata) { return Ok(None); } - // Safe: presence is guaranteed by the contains_key check above. - let server_side_encryption = metadata.get("x-amz-server-side-encryption").cloned().unwrap_or_default(); + let stored_algorithm = get_consistent_metadata_value(metadata, AMZ_SERVER_SIDE_ENCRYPTION).map_err(|_| { + ApiError::from(StorageError::other(format!( + "Conflicting managed encryption metadata for {AMZ_SERVER_SIDE_ENCRYPTION}" + ))) + })?; + let scheme = match stored_algorithm { + Some(ServerSideEncryption::AWS_KMS) => ManagedSseScheme::SseKms, + // RUSTFS_COMPAT_TODO(rustfs-5063): Legacy objects may omit this header. Remove after every referenced legacy DEK is rewrapped. + Some(ServerSideEncryption::AES256) | None => ManagedSseScheme::SseS3, + Some(algorithm) => { + return Err(ApiError::from(StorageError::other(format!( + "Unsupported stored server-side encryption: {algorithm}" + )))); + } + }; + let server_side_encryption = stored_algorithm.unwrap_or(ServerSideEncryption::AES256).to_string(); let normalized_metadata = normalize_managed_metadata(metadata); - let encryption_type = match server_side_encryption.as_str() { - ServerSideEncryption::AES256 => SSEType::SseS3, - ServerSideEncryption::AWS_KMS => SSEType::SseKms, - _ => SSEType::SseS3, + let encryption_type = match scheme { + ManagedSseScheme::SseS3 => SSEType::SseS3, + ManagedSseScheme::SseKms => SSEType::SseKms, }; #[cfg(feature = "rio-v2")] let minio_sealed_key = parse_minio_managed_sealed_key(metadata, encryption_type)?; @@ -1673,19 +1669,7 @@ async fn apply_managed_decryption_material( .cloned() .unwrap_or_else(|| "AES256".to_string()), ) - } else if let Some(service) = runtime_sources::current_encryption_service().await { - // Production mode: use service for metadata parsing - let parsed = service - .headers_to_metadata(&normalized_metadata) - .map_err(|e| ApiError::from(StorageError::other(format!("Failed to parse encryption metadata: {e}"))))?; - - if parsed.iv.len() != 12 { - return Err(ApiError::from(StorageError::other("Invalid encryption nonce length; expected 12 bytes"))); - } - - (parsed.encrypted_data_key, parsed.iv, parsed.algorithm) } else { - // Test mode: parse metadata manually let encrypted_key_b64 = normalized_metadata .get(INTERNAL_ENCRYPTION_KEY_HEADER) .ok_or_else(|| ApiError::from(StorageError::other("Missing encrypted key in metadata")))?; @@ -1718,15 +1702,20 @@ async fn apply_managed_decryption_material( .or_else(|| metadata.get("x-amz-server-side-encryption-aws-kms-key-id")) .cloned() .unwrap_or_else(|| "default".to_string()); - let kms_context = if matches!(encryption_type, SSEType::SseKms) { - decode_minio_kms_context(metadata)? + let has_kms_envelope = rustfs_kms::is_data_key_envelope(&encrypted_data_key); + // RUSTFS_COMPAT_TODO(rustfs-5063): Keep legacy SSE-S3 KMS envelopes readable. Remove after SSE-S3 migration rewraps every referenced legacy DEK. + let provider_kind = managed_dek_provider(scheme, has_kms_envelope); + let kms_context = if matches!(provider_kind, ManagedDekProvider::Kms) { + decode_managed_kms_context(metadata).map_err(|err| ApiError::from(StorageError::other(err.to_string())))? } else { None }; let object_context = build_object_encryption_context(bucket, key, kms_context.as_ref()); - - // Use factory pattern to get provider (test or production mode) - let provider = get_sse_dek_provider().await?; + let provider_type = match provider_kind { + ManagedDekProvider::LocalSseS3 => SSEType::SseS3, + ManagedDekProvider::Kms => SSEType::SseKms, + }; + let provider = get_sse_dek_provider(provider_type).await?; #[cfg(feature = "rio-v2")] let decrypted_data_key = if is_legacy_rustfs_managed_metadata(&normalized_metadata) { provider @@ -1809,6 +1798,10 @@ pub struct SsecParams { /// Abstracts the source of encryption keys (KMS, test provider, etc.) #[async_trait] pub trait SseDekProvider: Send + Sync { + fn default_kms_key_id(&self) -> Option { + None + } + /// Generate an SSE data encryption key async fn generate_sse_dek(&self, context: &ObjectEncryptionContext, kms_key_id: &str) -> Result<(DataKey, Vec), ApiError>; @@ -1833,66 +1826,59 @@ pub trait SseDekProvider: Send + Sync { } } +fn sse_not_configured(message: impl Into) -> ApiError { + ApiError { + code: S3ErrorCode::InvalidRequest, + message: message.into(), + source: None, + } +} + // ============================================================================ // Production KMS-backed DEK Provider // ============================================================================ -/// Production KMS-backed DEK provider -/// Resolves the latest ObjectEncryptionService on each call. +/// Production KMS-backed DEK provider. +/// +/// Each instance retains one service snapshot so a request cannot mix KMS +/// generations while resolving the default key and creating or decrypting a DEK. struct KmsSseDekProvider { - #[cfg(test)] - service_manager: Option>, + service: Arc, } impl KmsSseDekProvider { /// Create a new KMS-backed provider pub async fn new() -> Result { - let provider = Self { - #[cfg(test)] - service_manager: None, - }; - provider - .current_service() + let service = runtime_sources::current_encryption_service() .await .ok_or_else(|| ApiError::from(StorageError::other("KMS encryption service is not initialized")))?; - Ok(provider) + Ok(Self { service }) } #[cfg(test)] async fn new_with_service_manager(service_manager: Arc) -> Result { - let provider = Self { - service_manager: Some(service_manager), - }; - provider - .current_service() + let service = service_manager + .get_encryption_service() .await .ok_or_else(|| ApiError::from(StorageError::other("KMS encryption service is not initialized")))?; - Ok(provider) - } - - async fn current_service(&self) -> Option> { - #[cfg(test)] - if let Some(service_manager) = &self.service_manager { - return service_manager.get_encryption_service().await; - } - - runtime_sources::current_encryption_service().await + Ok(Self { service }) } } #[async_trait] impl SseDekProvider for KmsSseDekProvider { + fn default_kms_key_id(&self) -> Option { + self.service.get_default_key_id().cloned() + } + async fn generate_sse_dek( &self, context: &ObjectEncryptionContext, kms_key_id: &str, ) -> Result<(DataKey, Vec), ApiError> { let kms_key_option = Some(kms_key_id.to_string()); - let service = self - .current_service() - .await - .ok_or_else(|| ApiError::from(StorageError::other("KMS encryption service is not initialized")))?; - let (data_key, encrypted_data_key) = service + let (data_key, encrypted_data_key) = self + .service .create_data_key(&kms_key_option, context) .await .map_err(|e| ApiError::from(StorageError::other(format!("Failed to create data key: {}", e))))?; @@ -1906,11 +1892,8 @@ impl SseDekProvider for KmsSseDekProvider { _kms_key_id: &str, context: &ObjectEncryptionContext, ) -> Result<[u8; 32], ApiError> { - let service = self - .current_service() - .await - .ok_or_else(|| ApiError::from(StorageError::other("KMS encryption service is not initialized")))?; - let data_key = service + let data_key = self + .service .decrypt_data_key(encrypted_dek, context) .await .map_err(|e| ApiError::from(StorageError::other(format!("Failed to decrypt data key: {}", e))))?; @@ -1925,11 +1908,8 @@ impl SseDekProvider for KmsSseDekProvider { _kms_key_id: &str, _context: &ObjectEncryptionContext, ) -> Result<[u8; 32], ApiError> { - let service = self - .current_service() - .await - .ok_or_else(|| ApiError::from(StorageError::other("KMS encryption service is not initialized")))?; - let data_key = service + let data_key = self + .service .decrypt_legacy_data_key(encrypted_dek) .await .map_err(|e| ApiError::from(StorageError::other(format!("Failed to decrypt legacy data key: {e}"))))?; @@ -1953,11 +1933,6 @@ impl SseDekProvider for KmsSseDekProvider { /// __RUSTFS_SSE_SIMPLE_CMK= /// ``` /// -/// Example: -/// ```bash -/// export __RUSTFS_SSE_SIMPLE_CMK="AKHul86TBMMJ3+VrGlh9X3dHJsOtSXOXHOODPwmAnOo=" -/// ``` -/// /// # Key Generation /// /// Use the provided script to generate a valid key: @@ -1977,6 +1952,7 @@ pub(crate) struct TestSseDekProvider { /// Returns an error (never crashes) for a missing, non-base64, wrong-length, or /// all-zero key so callers on the request path can fail the request instead of /// taking the whole server down (backlog#806). +#[cfg(test)] fn parse_simple_sse_cmk(cmk_value: &str) -> Result<[u8; 32], ApiError> { let trimmed = cmk_value.trim(); if trimmed.is_empty() { @@ -2006,6 +1982,7 @@ impl TestSseDekProvider { Self { master_key } } + #[cfg(test)] pub fn new() -> Result { let cmk_value = std::env::var("__RUSTFS_SSE_SIMPLE_CMK").unwrap_or_default(); // A missing/invalid key must surface as a request error, never crash the @@ -2017,37 +1994,37 @@ impl TestSseDekProvider { Ok(Self { master_key }) } - /// Create a local SSE DEK provider for SSE-S3 when KMS is not configured. + /// Create a local SSE DEK provider for SSE-S3. /// Requires RUSTFS_SSE_S3_MASTER_KEY to be a valid base64-encoded 32-byte key. /// /// The failures here are server configuration problems, not internal /// faults: surface them as `InvalidRequest` (HTTP 400) so a managed-SSE /// request against an unconfigured server does not report 500 (rustfs#4844). pub fn new_for_local_sse() -> Result { - fn sse_not_configured(message: impl Into) -> ApiError { - ApiError { - code: S3ErrorCode::InvalidRequest, - message: message.into(), - source: None, - } - } - let Some(raw_value) = get_env_opt_str("RUSTFS_SSE_S3_MASTER_KEY").filter(|value| !value.trim().is_empty()) else { return Err(sse_not_configured( - "SSE-S3 requires RUSTFS_SSE_S3_MASTER_KEY to be set to a base64-encoded 32-byte key when KMS is not configured", + "SSE-S3 requires RUSTFS_SSE_S3_MASTER_KEY to be set to a base64-encoded 32-byte key", )); }; - let decoded = BASE64_STANDARD.decode(raw_value.trim()).map_err(|err| { - sse_not_configured(format!( - "RUSTFS_SSE_S3_MASTER_KEY must be valid base64 for SSE-S3 when KMS is not configured: {err}" - )) - })?; - let master_key: [u8; 32] = decoded.try_into().map_err(|_| { - sse_not_configured("RUSTFS_SSE_S3_MASTER_KEY must decode to exactly 32 bytes for SSE-S3 when KMS is not configured") - })?; + let decoded = BASE64_STANDARD + .decode(raw_value.trim()) + .map_err(|err| sse_not_configured(format!("RUSTFS_SSE_S3_MASTER_KEY must be valid base64 for SSE-S3: {err}")))?; + let master_key: [u8; 32] = decoded + .try_into() + .map_err(|_| sse_not_configured("RUSTFS_SSE_S3_MASTER_KEY must decode to exactly 32 bytes for SSE-S3"))?; + if master_key == [0u8; 32] { + return Err(sse_not_configured("RUSTFS_SSE_S3_MASTER_KEY must not be an all-zero key")); + } - tracing::info!("Using RUSTFS_SSE_S3_MASTER_KEY for SSE-S3 (KMS not configured)"); + tracing::info!( + event = "sse_key_source_loaded", + component = "storage", + subsystem = "sse_s3", + state = "ready", + key_source = "environment", + "SSE-S3 key source loaded" + ); Ok(Self { master_key }) } @@ -2148,53 +2125,101 @@ impl SseDekProvider for TestSseDekProvider { // Factory Function for SSE DEK Provider // ============================================================================ -/// Global SSE DEK provider cache -static GLOBAL_SSE_DEK_PROVIDER: LazyLock>>> = LazyLock::new(|| RwLock::new(None)); +/// Global local SSE-S3 DEK provider cache. +static GLOBAL_LOCAL_SSE_DEK_PROVIDER: LazyLock>>> = LazyLock::new(|| RwLock::new(None)); + +#[cfg(test)] +static TEST_KMS_SSE_DEK_PROVIDER: LazyLock>>> = LazyLock::new(|| RwLock::new(None)); + +#[cfg(test)] +struct TestKmsSseDekProviderGuard { + previous: Option>, +} + +#[cfg(test)] +impl Drop for TestKmsSseDekProviderGuard { + fn drop(&mut self) { + let mut slot = TEST_KMS_SSE_DEK_PROVIDER + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner); + *slot = self.previous.take(); + } +} /// Get or initialize the global SSE DEK provider /// -/// Factory function that automatically selects the appropriate provider: -/// - If `__RUSTFS_SSE_SIMPLE_CMK` environment variable exists: use SimpleSseDekProvider (test mode) -/// - Otherwise: use KmsSseDekProvider (production mode with real KMS) +/// The requested encryption type selects the provider. SSE-S3 never selects +/// KMS, and SSE-KMS never falls back to the local SSE-S3 master key. /// /// # Returns /// Arc to the global SSE DEK provider instance /// /// # Example /// ```rust,ignore -/// let provider = get_sse_dek_provider().await?; +/// let provider = get_sse_dek_provider(SSEType::SseS3).await?; /// let (data_key, encrypted_dek) = provider /// .generate_sse_dek("bucket", "key", "kms-key-id") /// .await?; /// ``` -pub async fn get_sse_dek_provider() -> Result, ApiError> { - if runtime_sources::current_encryption_service().await.is_some() { - debug!("Using KmsSseDekProvider (KMS configured)"); - return Ok(Arc::new(KmsSseDekProvider::new().await?)); +pub async fn get_sse_dek_provider(encryption_type: SSEType) -> Result, ApiError> { + if matches!(encryption_type, SSEType::SseKms) { + #[cfg(test)] + if let Some(provider) = TEST_KMS_SSE_DEK_PROVIDER + .read() + .map_err(|_| ApiError::from(StorageError::other("Failed to read test KMS SSE DEK provider")))? + .as_ref() + .cloned() + { + return Ok(provider); + } + + debug!( + event = "sse_dek_provider_selected", + component = "storage", + subsystem = "sse_kms", + provider = "kms", + "SSE DEK provider selected" + ); + return KmsSseDekProvider::new() + .await + .map(|provider| Arc::new(provider) as Arc) + .map_err(|_| sse_not_configured("SSE-KMS requires a configured KMS service")); + } + + if matches!(encryption_type, SSEType::SseC) { + return Err(ApiError::from(StorageError::other("SSE-C cannot use a managed DEK provider"))); } // Check if already initialized - if let Some(provider) = GLOBAL_SSE_DEK_PROVIDER + if let Some(provider) = GLOBAL_LOCAL_SSE_DEK_PROVIDER .read() - .map_err(|_| ApiError::from(StorageError::other("Failed to read global SSE DEK provider cache")))? + .map_err(|_| ApiError::from(StorageError::other("Failed to read local SSE-S3 DEK provider cache")))? .as_ref() .cloned() { return Ok(provider); } - // Determine provider: KMS when available, else test env, else local SSE-S3 fallback (no KMS) + #[cfg(test)] let provider: Arc = if std::env::var("__RUSTFS_SSE_SIMPLE_CMK").is_ok() { debug!("Using SimpleSseDekProvider (test mode) based on __RUSTFS_SSE_SIMPLE_CMK"); Arc::new(TestSseDekProvider::new()?) } else { - debug!("Using local SSE-S3 provider (KMS not configured)"); + debug!( + event = "sse_dek_provider_selected", + component = "storage", + subsystem = "sse_s3", + provider = "local", + "SSE DEK provider selected" + ); Arc::new(TestSseDekProvider::new_for_local_sse()?) }; + #[cfg(not(test))] + let provider: Arc = Arc::new(TestSseDekProvider::new_for_local_sse()?); - let mut slot = GLOBAL_SSE_DEK_PROVIDER + let mut slot = GLOBAL_LOCAL_SSE_DEK_PROVIDER .write() - .map_err(|_| ApiError::from(StorageError::other("Failed to update global SSE DEK provider cache")))?; + .map_err(|_| ApiError::from(StorageError::other("Failed to update local SSE-S3 DEK provider cache")))?; if let Some(existing) = slot.as_ref() { return Ok(existing.clone()); } @@ -2206,20 +2231,25 @@ pub async fn get_sse_dek_provider() -> Result, ApiError> /// Reset the global SSE DEK provider (for testing only) /// /// Note: OnceLock doesn't support reset in stable Rust. -/// Tests should set environment variables before first call to `get_sse_dek_provider()`. +/// Tests should set environment variables before the first provider lookup. #[cfg(test)] #[allow(dead_code)] pub fn reset_sse_dek_provider() { - if let Ok(mut slot) = GLOBAL_SSE_DEK_PROVIDER.write() { + if let Ok(mut slot) = GLOBAL_LOCAL_SSE_DEK_PROVIDER.write() { + *slot = None; + } + if let Ok(mut slot) = TEST_KMS_SSE_DEK_PROVIDER.write() { *slot = None; } } #[cfg(test)] -#[cfg(feature = "rio-v2")] -pub fn set_sse_dek_provider_for_test(provider: Arc) { - if let Ok(mut slot) = GLOBAL_SSE_DEK_PROVIDER.write() { - *slot = Some(provider); +fn set_sse_kms_dek_provider_for_test(provider: Arc) -> TestKmsSseDekProviderGuard { + let mut slot = TEST_KMS_SSE_DEK_PROVIDER + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner); + TestKmsSseDekProviderGuard { + previous: slot.replace(provider), } } @@ -2244,7 +2274,7 @@ pub fn strip_managed_encryption_metadata(metadata: &mut HashMap) INTERNAL_ENCRYPTION_IV_HEADER, "x-rustfs-encryption-tag", INTERNAL_ENCRYPTION_KEY_HEADER, - "x-rustfs-encryption-context", + RUSTFS_ENCRYPTION_CONTEXT_HEADER, INTERNAL_ENCRYPTION_ORIGINAL_SIZE_HEADER, MINIO_INTERNAL_ENCRYPTION_MULTIPART_HEADER, MINIO_INTERNAL_ENCRYPTION_IV_HEADER, @@ -2350,13 +2380,13 @@ fn normalize_managed_metadata(metadata: &HashMap) -> HashMap>(&decoded) && let Ok(encoded) = serde_json::to_string(&context) { - normalized.insert("x-rustfs-encryption-context".to_string(), encoded); + normalized.insert(RUSTFS_ENCRYPTION_CONTEXT_HEADER.to_string(), encoded); } normalized @@ -2510,13 +2540,13 @@ mod tests { MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER, MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER, MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER, MINIO_INTERNAL_ENCRYPTION_MULTIPART_HEADER, MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER, MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER, - PrepareEncryptionRequest, SSEC_ORIGINAL_SIZE_HEADER, SSEType, SseDekProvider, SsecParams, StorageError, - TestSseDekProvider, encryption_material_to_metadata, extract_server_side_encryption_from_headers, - extract_ssec_params_from_headers, extract_ssekms_context_from_headers, generate_ssec_nonce, is_managed_sse, - map_get_object_reader_error, mark_encrypted_multipart_metadata, normalize_managed_metadata, reset_sse_dek_provider, - resolve_effective_kms_key_id, sse_decryption, sse_encryption, sse_prepare_encryption, strip_managed_encryption_metadata, - validate_sse_headers_for_read, validate_sse_headers_for_write, validate_ssec_for_read, validate_ssec_params, - verify_ssec_key_match, + PrepareEncryptionRequest, RUSTFS_ENCRYPTION_CONTEXT_HEADER, SSEC_ORIGINAL_SIZE_HEADER, SSEType, SseDekProvider, + SsecParams, StorageError, TestSseDekProvider, encryption_material_to_metadata, + extract_server_side_encryption_from_headers, extract_ssec_params_from_headers, extract_ssekms_context_from_headers, + generate_ssec_nonce, is_managed_sse, map_get_object_reader_error, mark_encrypted_multipart_metadata, + normalize_managed_metadata, reset_sse_dek_provider, resolve_effective_kms_key_id, sse_decryption, sse_encryption, + sse_prepare_encryption, strip_managed_encryption_metadata, validate_sse_headers_for_read, validate_sse_headers_for_write, + validate_ssec_for_read, validate_ssec_params, verify_ssec_key_match, }; #[test] @@ -2542,6 +2572,7 @@ mod tests { let got = super::parse_simple_sse_cmk(&encoded).expect("valid 32-byte key must parse"); assert_eq!(got, key); } + use aes_gcm::aead::{Aead, KeyInit}; use aes_gcm::{Aes256Gcm, Key, Nonce}; use base64::{Engine, engine::general_purpose::STANDARD as BASE64_STANDARD}; @@ -2552,6 +2583,7 @@ mod tests { use s3s::S3ErrorCode; use s3s::dto::{SSECustomerAlgorithm, SSECustomerKey, SSECustomerKeyMD5, ServerSideEncryption}; use std::collections::HashMap; + use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, OnceLock}; use temp_env::async_with_vars; use tokio::sync::Mutex; @@ -2566,6 +2598,40 @@ mod tests { BASE64_STANDARD.encode([0x24u8; 32]) } + struct CountingKmsDekProvider { + generate_calls: Arc, + decrypt_calls: Arc, + plaintext_key: [u8; 32], + } + + #[async_trait::async_trait] + impl SseDekProvider for CountingKmsDekProvider { + async fn generate_sse_dek( + &self, + _context: &ObjectEncryptionContext, + _kms_key_id: &str, + ) -> Result<(rustfs_kms::DataKey, Vec), crate::error::ApiError> { + self.generate_calls.fetch_add(1, Ordering::SeqCst); + Ok(( + rustfs_kms::DataKey { + plaintext_key: self.plaintext_key, + nonce: [0x33; 12], + }, + b"kms-envelope".to_vec(), + )) + } + + async fn decrypt_sse_dek( + &self, + _encrypted_dek: &[u8], + _kms_key_id: &str, + _context: &ObjectEncryptionContext, + ) -> Result<[u8; 32], crate::error::ApiError> { + self.decrypt_calls.fetch_add(1, Ordering::SeqCst); + Ok(self.plaintext_key) + } + } + #[test] fn test_extract_ssec_params_from_headers() { let mut headers = http::HeaderMap::new(); @@ -3293,7 +3359,7 @@ mod tests { let provider = KmsSseDekProvider::new_with_service_manager(manager.clone()) .await .expect("kms provider should initialize from the configured test manager"); - super::set_sse_dek_provider_for_test(Arc::new(provider)); + let _provider_guard = super::set_sse_kms_dek_provider_for_test(Arc::new(provider)); let client_context = HashMap::from([("tenant".to_string(), "alpha".to_string())]); let request = EncryptionRequest { @@ -3437,6 +3503,217 @@ mod tests { reset_sse_dek_provider(); } + #[tokio::test] + async fn test_sse_s3_never_calls_available_kms_provider() { + let _guard = lock_sse_test_state().await; + reset_sse_dek_provider(); + let generate_calls = Arc::new(AtomicUsize::new(0)); + let decrypt_calls = Arc::new(AtomicUsize::new(0)); + let _provider_guard = super::set_sse_kms_dek_provider_for_test(Arc::new(CountingKmsDekProvider { + generate_calls: generate_calls.clone(), + decrypt_calls: decrypt_calls.clone(), + plaintext_key: [0x55; 32], + })); + let local_sse_master_key = local_sse_master_key_b64(); + + async_with_vars( + [ + ("__RUSTFS_SSE_SIMPLE_CMK", None::<&str>), + ("RUSTFS_SSE_S3_MASTER_KEY", Some(local_sse_master_key.as_str())), + ], + async { + let material = sse_encryption(EncryptionRequest { + bucket: "test-bucket", + key: "test-key", + server_side_encryption: Some(ServerSideEncryption::from_static(ServerSideEncryption::AES256)), + ssekms_key_id: None, + ssekms_context: None, + sse_customer_algorithm: None, + sse_customer_key: None, + sse_customer_key_md5: None, + content_size: 1024, + }) + .await + .expect("SSE-S3 encryption should use the local provider") + .expect("SSE-S3 encryption should return material"); + let metadata = encryption_material_to_metadata(&material).expect("SSE-S3 metadata should serialize"); + let decrypted = sse_decryption(DecryptionRequest { + bucket: "test-bucket", + key: "test-key", + metadata: &metadata, + sse_customer_key: None, + sse_customer_key_md5: None, + }) + .await + .expect("SSE-S3 decryption should use the local provider") + .expect("SSE-S3 decryption should return material"); + + assert_eq!(decrypted.key_bytes, material.key_bytes); + assert_eq!(generate_calls.load(Ordering::SeqCst), 0); + assert_eq!(decrypt_calls.load(Ordering::SeqCst), 0); + }, + ) + .await; + + reset_sse_dek_provider(); + } + + #[tokio::test] + async fn test_sse_kms_never_falls_back_to_local_sse_s3_provider() { + let _guard = lock_sse_test_state().await; + reset_sse_dek_provider(); + let local_sse_master_key = local_sse_master_key_b64(); + + async_with_vars( + [ + ("__RUSTFS_SSE_SIMPLE_CMK", None::<&str>), + ("RUSTFS_SSE_S3_MASTER_KEY", Some(local_sse_master_key.as_str())), + ], + async { + let err = sse_encryption(EncryptionRequest { + bucket: "test-bucket", + key: "test-key", + server_side_encryption: Some(ServerSideEncryption::from_static(ServerSideEncryption::AWS_KMS)), + ssekms_key_id: Some("kms-key".to_string()), + ssekms_context: None, + sse_customer_algorithm: None, + sse_customer_key: None, + sse_customer_key_md5: None, + content_size: 1024, + }) + .await + .expect_err("SSE-KMS must fail when KMS is unavailable even if the SSE-S3 key is configured"); + + assert_eq!(err.code, S3ErrorCode::InvalidRequest); + assert!(err.message.contains("configured KMS service")); + }, + ) + .await; + + reset_sse_dek_provider(); + } + + #[tokio::test] + async fn test_kms_provider_override_guard_restores_previous_provider() { + let _guard = lock_sse_test_state().await; + reset_sse_dek_provider(); + let first_generate_calls = Arc::new(AtomicUsize::new(0)); + let second_generate_calls = Arc::new(AtomicUsize::new(0)); + let first_guard = super::set_sse_kms_dek_provider_for_test(Arc::new(CountingKmsDekProvider { + generate_calls: first_generate_calls.clone(), + decrypt_calls: Arc::new(AtomicUsize::new(0)), + plaintext_key: [0x51; 32], + })); + let second_guard = super::set_sse_kms_dek_provider_for_test(Arc::new(CountingKmsDekProvider { + generate_calls: second_generate_calls.clone(), + decrypt_calls: Arc::new(AtomicUsize::new(0)), + plaintext_key: [0x52; 32], + })); + let context = ObjectEncryptionContext::new("bucket".to_string(), "object".to_string()); + + super::get_sse_dek_provider(SSEType::SseKms) + .await + .expect("second test provider should be selected") + .generate_sse_dek(&context, "key") + .await + .expect("second test provider should generate a DEK"); + drop(second_guard); + super::get_sse_dek_provider(SSEType::SseKms) + .await + .expect("first test provider should be restored") + .generate_sse_dek(&context, "key") + .await + .expect("restored test provider should generate a DEK"); + + assert_eq!(second_generate_calls.load(Ordering::SeqCst), 1); + assert_eq!(first_generate_calls.load(Ordering::SeqCst), 1); + drop(first_guard); + assert!( + super::TEST_KMS_SSE_DEK_PROVIDER + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .is_none() + ); + reset_sse_dek_provider(); + } + + #[tokio::test] + async fn test_legacy_sse_s3_kms_envelope_keeps_kms_read_path() { + let _guard = lock_sse_test_state().await; + reset_sse_dek_provider(); + let generate_calls = Arc::new(AtomicUsize::new(0)); + let decrypt_calls = Arc::new(AtomicUsize::new(0)); + let plaintext_key = [0x66; 32]; + let _provider_guard = super::set_sse_kms_dek_provider_for_test(Arc::new(CountingKmsDekProvider { + generate_calls: generate_calls.clone(), + decrypt_calls: decrypt_calls.clone(), + plaintext_key, + })); + let legacy_envelope = serde_json::to_vec(&serde_json::json!({ + "key_id": "legacy-data-key", + "master_key_id": "legacy-master-key", + "key_spec": "AES_256", + "encrypted_key": [1, 2, 3, 4], + "nonce": [5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16], + "encryption_context": {}, + "created_at": "2024-01-01T00:00:00+00:00" + })) + .expect("legacy KMS envelope should serialize"); + let metadata = HashMap::from([ + (super::AMZ_SERVER_SIDE_ENCRYPTION.to_string(), ServerSideEncryption::AES256.to_string()), + (INTERNAL_ENCRYPTION_KEY_ID_HEADER.to_string(), "legacy-master-key".to_string()), + (INTERNAL_ENCRYPTION_KEY_HEADER.to_string(), BASE64_STANDARD.encode(legacy_envelope)), + (INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode([0x44; 12])), + (INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), ServerSideEncryption::AES256.to_string()), + ]); + + for metadata in [ + metadata.clone(), + metadata + .clone() + .into_iter() + .filter(|(key, _)| !key.eq_ignore_ascii_case(super::AMZ_SERVER_SIDE_ENCRYPTION)) + .collect(), + ] { + let decrypted = sse_decryption(DecryptionRequest { + bucket: "test-bucket", + key: "test-key", + metadata: &metadata, + sse_customer_key: None, + sse_customer_key_md5: None, + }) + .await + .expect("legacy SSE-S3 KMS envelope should remain readable") + .expect("legacy SSE-S3 metadata should return material"); + + assert_eq!(decrypted.key_bytes, plaintext_key); + } + + let mut conflicting_context = metadata.clone(); + conflicting_context.insert( + RUSTFS_ENCRYPTION_CONTEXT_HEADER.to_string(), + serde_json::to_string(&HashMap::from([("tenant", "alpha")])).expect("RustFS context should serialize"), + ); + conflicting_context.insert( + MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER.to_string(), + BASE64_STANDARD + .encode(serde_json::to_vec(&HashMap::from([("tenant", "beta")])).expect("MinIO context should serialize")), + ); + sse_decryption(DecryptionRequest { + bucket: "test-bucket", + key: "test-key", + metadata: &conflicting_context, + sse_customer_key: None, + sse_customer_key_md5: None, + }) + .await + .expect_err("legacy SSE-S3 KMS envelopes must reject conflicting KMS contexts"); + + assert_eq!(generate_calls.load(Ordering::SeqCst), 0); + assert_eq!(decrypt_calls.load(Ordering::SeqCst), 2); + reset_sse_dek_provider(); + } + #[test] fn test_strip_managed_encryption_metadata() { let mut metadata = HashMap::new(); @@ -3941,35 +4218,40 @@ mod tests { async fn test_sse_encryption_fails_closed_with_invalid_local_sse_master_key() { let _guard = lock_sse_test_state().await; reset_sse_dek_provider(); - async_with_vars( - [ - ("__RUSTFS_SSE_SIMPLE_CMK", None::<&str>), - ("RUSTFS_SSE_S3_MASTER_KEY", Some("not-base64")), - ], - async { - let err = sse_encryption(EncryptionRequest { - bucket: "test-bucket", - key: "test-key", - server_side_encryption: Some(ServerSideEncryption::from_static(ServerSideEncryption::AES256)), - ssekms_key_id: None, - ssekms_context: None, - sse_customer_algorithm: None, - sse_customer_key: None, - sse_customer_key_md5: None, - content_size: 1024, - }) - .await - .expect_err("SSE-S3 should fail closed with an invalid local master key"); + for (invalid_key, expected_message) in [ + ("not-base64".to_string(), "valid base64"), + (BASE64_STANDARD.encode([0u8; 32]), "all-zero"), + ] { + async_with_vars( + [ + ("__RUSTFS_SSE_SIMPLE_CMK", None::<&str>), + ("RUSTFS_SSE_S3_MASTER_KEY", Some(invalid_key.as_str())), + ], + async { + let err = sse_encryption(EncryptionRequest { + bucket: "test-bucket", + key: "test-key", + server_side_encryption: Some(ServerSideEncryption::from_static(ServerSideEncryption::AES256)), + ssekms_key_id: None, + ssekms_context: None, + sse_customer_algorithm: None, + sse_customer_key: None, + sse_customer_key_md5: None, + content_size: 1024, + }) + .await + .expect_err("SSE-S3 should fail closed with an invalid local master key"); - assert!(err.message.contains("valid base64")); - assert_eq!( - err.code, - S3ErrorCode::InvalidRequest, - "invalid local SSE master key is a configuration error, not a 500 (rustfs#4844)" - ); - }, - ) - .await; + assert!(err.message.contains(expected_message)); + assert_eq!( + err.code, + S3ErrorCode::InvalidRequest, + "invalid local SSE master key is a configuration error, not a 500 (rustfs#4844)" + ); + }, + ) + .await; + } reset_sse_dek_provider(); } @@ -4036,7 +4318,7 @@ mod tests { } #[tokio::test] - async fn test_kms_sse_dek_provider_uses_latest_reconfigured_service() { + async fn test_kms_sse_dek_provider_keeps_one_service_snapshot_per_request() { use rustfs_kms::config::KmsConfig; use rustfs_kms::types::{CreateKeyRequest, KeyUsage}; use tempfile::TempDir; @@ -4046,7 +4328,11 @@ mod tests { let first_dir = TempDir::new().expect("first temp dir"); manager - .reconfigure(KmsConfig::local(first_dir.path().to_path_buf()).with_insecure_development_defaults()) + .reconfigure( + KmsConfig::local(first_dir.path().to_path_buf()) + .with_default_key("first-key".to_string()) + .with_insecure_development_defaults(), + ) .await .expect("first KMS reconfigure should succeed"); manager @@ -4067,15 +4353,21 @@ mod tests { let provider = KmsSseDekProvider::new_with_service_manager(manager.clone()) .await .expect("provider should initialize"); + assert_eq!(provider.default_kms_key_id().as_deref(), Some("first-key")); let context = ObjectEncryptionContext::new("bucket".to_string(), "object".to_string()); + let first_default_key = provider.default_kms_key_id().expect("first default key should exist"); provider - .generate_sse_dek(&context, "first-key") + .generate_sse_dek(&context, &first_default_key) .await .expect("provider should use the initial service"); let second_dir = TempDir::new().expect("second temp dir"); manager - .reconfigure(KmsConfig::local(second_dir.path().to_path_buf()).with_insecure_development_defaults()) + .reconfigure( + KmsConfig::local(second_dir.path().to_path_buf()) + .with_default_key("second-key".to_string()) + .with_insecure_development_defaults(), + ) .await .expect("second KMS reconfigure should succeed"); manager @@ -4093,10 +4385,23 @@ mod tests { .await .expect("second key should be created"); + assert_eq!(provider.default_kms_key_id().as_deref(), Some("first-key")); + let captured_default_key = provider + .default_kms_key_id() + .expect("captured default key should remain available"); provider + .generate_sse_dek(&context, &captured_default_key) + .await + .expect("existing provider should retain its default key and service snapshot"); + + let reconfigured_provider = KmsSseDekProvider::new_with_service_manager(manager.clone()) + .await + .expect("new provider should capture the reconfigured service"); + assert_eq!(reconfigured_provider.default_kms_key_id().as_deref(), Some("second-key")); + reconfigured_provider .generate_sse_dek(&context, "second-key") .await - .expect("provider should resolve the latest reconfigured service"); + .expect("new provider should use the reconfigured service"); manager.stop().await.expect("kms service should stop cleanly"); } diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 6df459857..8ad82439d 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -502,6 +502,10 @@ pub(crate) mod ecstore_set_disk { pub(crate) use rustfs_ecstore::api::set_disk::{DEFAULT_READ_BUFFER_SIZE, get_lock_acquire_timeout, is_valid_storage_class}; } +pub(crate) mod ecstore_sse { + pub(crate) use rustfs_ecstore::api::sse::{ManagedDekProvider, ManagedSseScheme, managed_dek_provider}; +} + pub(crate) mod ecstore_storage { #[cfg(test)] pub(crate) use rustfs_ecstore::api::storage::init_local_disks;