From 2aaac85160f8fb718a06b257307f41feca551ce0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Wed, 29 Jul 2026 11:32:58 +0800 Subject: [PATCH] refactor(sse): decouple encryption from ecstore --- Cargo.lock | 3 - crates/ecstore/Cargo.toml | 3 - crates/ecstore/src/api/mod.rs | 10 +- crates/ecstore/src/client/utils.rs | 16 +- crates/ecstore/src/object_api/encryption.rs | 80 + crates/ecstore/src/object_api/mod.rs | 5 + crates/ecstore/src/object_api/readers.rs | 1646 ++--------------- crates/ecstore/src/object_api/types.rs | 63 +- crates/ecstore/src/runtime/instance.rs | 17 + crates/ecstore/src/runtime/sources.rs | 5 - crates/ecstore/src/set_disk/metadata.rs | 4 +- crates/ecstore/src/set_disk/mod.rs | 13 +- crates/ecstore/src/set_disk/ops/object.rs | 73 +- crates/ecstore/tests/README.md | 4 +- .../rio-v2/tests/minio_fixture_lab/README.md | 2 +- crates/utils/src/http/header_compat.rs | 76 +- .../ecstore-validation-suite-design.md | 2 +- .../src/storage}/minio_generated_read_test.rs | 20 +- rustfs/src/storage/mod.rs | 2 + rustfs/src/storage/sse.rs | 283 ++- rustfs/src/storage/storage_api.rs | 39 +- scripts/run_ecstore_validation_suite.sh | 4 +- 22 files changed, 736 insertions(+), 1634 deletions(-) create mode 100644 crates/ecstore/src/object_api/encryption.rs rename {crates/ecstore/tests => rustfs/src/storage}/minio_generated_read_test.rs (96%) diff --git a/Cargo.lock b/Cargo.lock index e769854d8..6a248e795 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9106,7 +9106,6 @@ dependencies = [ name = "rustfs-ecstore" version = "1.0.0-beta.11" dependencies = [ - "aes-gcm", "arc-swap", "async-channel", "async-recursion", @@ -9122,7 +9121,6 @@ dependencies = [ "byteorder", "bytes", "bytesize", - "chacha20poly1305", "chrono", "criterion", "enumset", @@ -9175,7 +9173,6 @@ dependencies = [ "rustfs-erasure-codec", "rustfs-filemeta", "rustfs-io-metrics", - "rustfs-kms", "rustfs-lifecycle", "rustfs-lock", "rustfs-madmin", diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index c35289f72..a4cad1097 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -57,7 +57,6 @@ rustfs-policy.workspace = true rustfs-protos.workspace = true rustfs-replication.workspace = true rustfs-lifecycle.workspace = true -rustfs-kms.workspace = true rustfs-s3-types = { workspace = true } rustfs-data-usage.workspace = true rustfs-object-capacity.workspace = true @@ -124,8 +123,6 @@ libc.workspace = true rustix = { workspace = true, features = ["process", "fs"] } rustfs-madmin.workspace = true reqwest = { workspace = true } -aes-gcm = { workspace = true, features = ["rand_core"] } -chacha20poly1305.workspace = true aws-sdk-s3 = { workspace = true, default-features = false, features = ["sigv4a", "default-https-client", "rt-tokio"] } urlencoding = { workspace = true } smallvec = { workspace = true, features = ["serde"] } diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index fc042f97f..cea6a2887 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -381,10 +381,12 @@ pub mod notification { pub mod object { pub use crate::object_api::{ - BLOCK_SIZE_V2, ERASURE_ALGORITHM, GetObjectBodyCacheHook, GetObjectBodyCacheHookLookup, GetObjectBodySource, - GetObjectReader, ObjectInfo, ObjectMutationHook, ObjectOptions, PutObjReader, RangedDecompressReader, StreamConsumer, - get_object_body_cache_plaintext_len, lookup_get_object_body_cache_hook, register_get_object_body_cache_hook, - register_object_mutation_hook, unregister_get_object_body_cache_hook, unregister_object_mutation_hook, + BLOCK_SIZE_V2, ERASURE_ALGORITHM, EncryptionResolutionError, EncryptionResolutionErrorKind, GetObjectBodyCacheHook, + GetObjectBodyCacheHookLookup, GetObjectBodySource, GetObjectReader, ObjectEncryptionResolver, ObjectInfo, + ObjectMutationHook, ObjectOptions, PutObjReader, RangedDecompressReader, ReadEncryptionMaterial, ReadEncryptionMode, + ReadEncryptionRequest, StreamConsumer, get_object_body_cache_plaintext_len, lookup_get_object_body_cache_hook, + register_get_object_body_cache_hook, register_object_mutation_hook, unregister_get_object_body_cache_hook, + unregister_object_mutation_hook, }; pub use crate::store::PreparedGetObjectReader; } diff --git a/crates/ecstore/src/client/utils.rs b/crates/ecstore/src/client/utils.rs index 5d6c975be..3a05e7cf3 100644 --- a/crates/ecstore/src/client/utils.rs +++ b/crates/ecstore/src/client/utils.rs @@ -46,16 +46,6 @@ lazy_static! { m.insert("x-amz-replication-status".to_string(), true); m }; - static ref SSE_HEADERS: HashMap = { - let mut m = HashMap::new(); - m.insert("x-amz-server-side-encryption".to_string(), true); - m.insert("x-amz-server-side-encryption-aws-kms-key-id".to_string(), true); - m.insert("x-amz-server-side-encryption-context".to_string(), true); - m.insert("x-amz-server-side-encryption-customer-algorithm".to_string(), true); - m.insert("x-amz-server-side-encryption-customer-key".to_string(), true); - m.insert("x-amz-server-side-encryption-customer-key-md5".to_string(), true); - m - }; } pub fn is_standard_query_value(qs_key: &str) -> bool { @@ -70,16 +60,12 @@ pub fn is_standard_header(header_key: &str) -> bool { *SUPPORTED_HEADERS.get(&header_key.to_lowercase()).unwrap_or(&false) } -pub fn is_sse_header(header_key: &str) -> bool { - *SSE_HEADERS.get(&header_key.to_lowercase()).unwrap_or(&false) -} - pub fn is_amz_header(header_key: &str) -> bool { let key = header_key.to_lowercase(); key.starts_with("x-amz-meta-") || key.starts_with("x-amz-grant-") || key == "x-amz-acl" - || is_sse_header(header_key) + || rustfs_utils::http::is_sse_header(header_key) || key.starts_with("x-amz-checksum-") } diff --git a/crates/ecstore/src/object_api/encryption.rs b/crates/ecstore/src/object_api/encryption.rs new file mode 100644 index 000000000..07413cdd9 --- /dev/null +++ b/crates/ecstore/src/object_api/encryption.rs @@ -0,0 +1,80 @@ +// 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 async_trait::async_trait; +use http::{HeaderMap, HeaderValue}; +use std::collections::HashMap; +use std::error::Error; +use std::fmt::{Display, Formatter}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ReadEncryptionMode { + Direct { base_nonce: [u8; 12] }, + Object, +} + +pub struct ReadEncryptionMaterial { + pub key_bytes: [u8; 32], + pub mode: ReadEncryptionMode, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum EncryptionResolutionErrorKind { + InvalidRequest, + InvalidMetadata, + ServiceUnavailable, + DecryptionFailed, +} + +#[derive(Debug)] +pub struct EncryptionResolutionError { + kind: EncryptionResolutionErrorKind, + message: String, +} + +impl EncryptionResolutionError { + pub fn new(kind: EncryptionResolutionErrorKind, message: impl Into) -> Self { + Self { + kind, + message: message.into(), + } + } + + pub fn kind(&self) -> EncryptionResolutionErrorKind { + self.kind + } +} + +impl Display for EncryptionResolutionError { + fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result { + formatter.write_str(&self.message) + } +} + +impl Error for EncryptionResolutionError {} + +pub struct ReadEncryptionRequest<'a> { + pub bucket: &'a str, + pub object: &'a str, + pub metadata: &'a HashMap, + pub headers: &'a HeaderMap, +} + +#[async_trait] +pub trait ObjectEncryptionResolver: Send + Sync { + async fn resolve_read_material( + &self, + request: ReadEncryptionRequest<'_>, + ) -> Result, EncryptionResolutionError>; +} diff --git a/crates/ecstore/src/object_api/mod.rs b/crates/ecstore/src/object_api/mod.rs index 5c57b69b2..28d567e75 100644 --- a/crates/ecstore/src/object_api/mod.rs +++ b/crates/ecstore/src/object_api/mod.rs @@ -84,6 +84,7 @@ pub(crate) fn legacy_encrypted_range_seek_enabled() -> bool { } mod body_cache_hook; +mod encryption; mod hook_slot; mod object_mutation_hook; mod readers; @@ -98,6 +99,10 @@ pub use body_cache_hook::{ pub(crate) use body_cache_hook::{ get_object_body_cache_hook, get_object_body_cache_hook_suppressed, without_get_object_body_cache_hook, }; +pub use encryption::{ + EncryptionResolutionError, EncryptionResolutionErrorKind, ObjectEncryptionResolver, ReadEncryptionMaterial, + ReadEncryptionMode, ReadEncryptionRequest, +}; pub(crate) use object_mutation_hook::notify_object_mutation; pub use object_mutation_hook::{ObjectMutationHook, register_object_mutation_hook, unregister_object_mutation_hook}; pub use readers::*; diff --git a/crates/ecstore/src/object_api/readers.rs b/crates/ecstore/src/object_api/readers.rs index 1330cc10b..ff82e2a4b 100644 --- a/crates/ecstore/src/object_api/readers.rs +++ b/crates/ecstore/src/object_api/readers.rs @@ -13,110 +13,13 @@ // limitations under the License. use super::*; -#[cfg(feature = "rio-v2")] -use aes_gcm::aead::Payload; -use aes_gcm::{ - Aes256Gcm, Key, Nonce, - aead::{Aead, KeyInit}, -}; -use base64::{Engine, engine::general_purpose::STANDARD as BASE64_STANDARD}; -#[cfg(feature = "rio-v2")] -use chacha20poly1305::ChaCha20Poly1305; -#[cfg(feature = "rio-v2")] -use hmac::{Hmac, Mac}; -use md5::{Digest, Md5}; -use rustfs_kms::{KmsUnavailableError, is_data_key_envelope, types::ObjectEncryptionContext}; -use rustfs_utils::http::{SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER}; -use rustfs_utils::path::path_join_buf; -use serde::Deserialize; -#[cfg(feature = "rio-v2")] -use sha2::Sha256; -use std::collections::HashMap; -use std::env; 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 LOCAL_SSE_DEK_FORMAT_VERSION: u8 = 1; #[cfg(feature = "rio-v2")] const DARE_PAYLOAD_SIZE: i64 = 64 * 1024; #[cfg(feature = "rio-v2")] const DARE_PACKAGE_SIZE: i64 = DARE_PAYLOAD_SIZE + 32; -const MINIO_INTERNAL_ENCRYPTION_IV_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-Iv"; -#[cfg(feature = "rio-v2")] -const MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-Seal-Algorithm"; -#[cfg(feature = "rio-v2")] -const MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-S3-Sealed-Key"; -#[cfg(feature = "rio-v2")] -const MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-Kms-Sealed-Key"; -#[cfg(feature = "rio-v2")] -const MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER: &str = "X-Minio-Internal-Server-Side-Encryption-S3-Kms-Key-Id"; -#[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"; -#[cfg(feature = "rio-v2")] -const DARE_VERSION_20: u8 = 0x20; -#[cfg(feature = "rio-v2")] -const DARE_CIPHER_AES_256_GCM: u8 = 0x00; -#[cfg(feature = "rio-v2")] -const DARE_CIPHER_CHACHA20_POLY1305: u8 = 0x01; -#[cfg(feature = "rio-v2")] -const DARE_HEADER_SIZE: usize = 16; -#[cfg(feature = "rio-v2")] -const DARE_TAG_SIZE: usize = 16; -#[cfg(feature = "rio-v2")] -const SEALED_KEY_IV_SIZE: usize = 32; -#[cfg(feature = "rio-v2")] -const SEALED_KEY_SIZE: usize = DARE_HEADER_SIZE + 32 + DARE_TAG_SIZE; -#[cfg(feature = "rio-v2")] -const MINIO_SECRET_KEY_RANDOM_SIZE: usize = 28; -#[cfg(feature = "rio-v2")] -const MINIO_SECRET_KEY_IV_SIZE: usize = 16; -#[cfg(feature = "rio-v2")] -const MINIO_SECRET_KEY_NONCE_SIZE: usize = 12; - -#[cfg(feature = "rio-v2")] -type HmacSha256 = Hmac; - -fn canonical_kms_bucket_path(bucket: &str, object: &str) -> String { - path_join_buf(&[bucket, object]) -} - -fn build_object_encryption_context( - bucket: &str, - object: &str, - provided_context: Option<&HashMap>, -) -> ObjectEncryptionContext { - let mut context = provided_context.cloned().unwrap_or_default(); - context - .entry(bucket.to_string()) - .or_insert_with(|| canonical_kms_bucket_path(bucket, object)); - - let mut object_context = ObjectEncryptionContext::new(bucket.to_string(), object.to_string()); - for (ctx_key, ctx_value) in context { - object_context = object_context.with_encryption_context(ctx_key, ctx_value); - } - object_context -} - -#[cfg(feature = "rio-v2")] -fn is_legacy_rustfs_managed_metadata(metadata: &HashMap) -> bool { - metadata_get(metadata, INTERNAL_ENCRYPTION_KEY_HEADER).is_some() - && metadata_get(metadata, INTERNAL_ENCRYPTION_IV_HEADER).is_some() - && metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER).is_none() - && metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER).is_none() -} fn part_plaintext_size(part: &ObjectPartInfo) -> i64 { if part.actual_size > 0 { @@ -544,21 +447,6 @@ impl GetObjectReader { } } -#[derive(Debug, Clone, Copy)] -struct EncryptionMaterial { - key_bytes: [u8; 32], - base_nonce: [u8; 12], - key_kind: EncryptionKeyKind, - reader_backend: crate::io_support::rio::ReadEncryptionBackend, -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -enum EncryptionKeyKind { - Direct, - Object, -} - -#[derive(Debug, Clone)] enum ReadTransform { Plain { visible_offset: usize, @@ -572,7 +460,7 @@ enum ReadTransform { total_plaintext_size: usize, }, Encrypted { - material: EncryptionMaterial, + material: ReadEncryptionMaterial, is_multipart: bool, part_numbers: Vec, sequence_number: u32, @@ -584,7 +472,6 @@ enum ReadTransform { }, } -#[derive(Debug, Clone)] struct ReadPlan { storage_offset: usize, storage_length: i64, @@ -593,7 +480,18 @@ struct ReadPlan { } impl ReadPlan { + #[cfg(test)] async fn build(rs: Option, oi: &ObjectInfo, opts: &ObjectOptions, h: &HeaderMap) -> Result { + Self::build_with_resolver(rs, oi, opts, h, Some(&tests::TEST_RESOLVER)).await + } + + async fn build_with_resolver( + rs: Option, + oi: &ObjectInfo, + opts: &ObjectOptions, + h: &HeaderMap, + resolver: Option<&dyn ObjectEncryptionResolver>, + ) -> Result { let mut rs = rs; if let Some(part_number) = opts.part_number && rs.is_none() @@ -657,9 +555,20 @@ impl ReadPlan { } if is_encrypted { - let material = resolve_encryption_material(oi, h).await?; + let resolver = resolver.ok_or_else(|| Error::other("object encryption resolver is unavailable"))?; + let resolved = resolver + .resolve_read_material(ReadEncryptionRequest { + bucket: &oi.bucket, + object: &oi.name, + metadata: &oi.user_defined, + headers: h, + }) + .await + .map_err(Error::other)? + .ok_or_else(|| Error::other("encrypted object metadata is incomplete"))?; + let material = resolved; #[cfg(feature = "rio-v2")] - let encryption_backend = material.reader_backend; + let uses_legacy_encryption = matches!(material.mode, ReadEncryptionMode::Direct { .. }); let is_multipart = is_multipart_encrypted_object(&oi.parts, oi.etag.as_deref()); let recorded_plaintext_size = oi.encryption_original_size()?; let plaintext_size = encrypted_plaintext_size(oi, is_multipart, is_compressed, recorded_plaintext_size)?; @@ -678,7 +587,7 @@ impl ReadPlan { let (requested_offset, requested_length) = rs.get_offset_length(plaintext_size)?; #[cfg(feature = "rio-v2")] { - if encryption_backend == crate::io_support::rio::ReadEncryptionBackend::Legacy { + if uses_legacy_encryption { legacy_encrypted_range_plan( oi, is_multipart, @@ -869,32 +778,32 @@ impl ReadPlan { #[cfg(not(feature = "rio-v2"))] let _ = sequence_number; let decrypted_reader: Box = if is_multipart { - match material.key_kind { - EncryptionKeyKind::Object => crate::io_support::rio::decrypt_multipart_reader_with_object_key( + match material.mode { + ReadEncryptionMode::Object => crate::io_support::rio::decrypt_multipart_reader_with_object_key( reader, material.key_bytes, part_numbers, sequence_number, ), - EncryptionKeyKind::Direct => crate::io_support::rio::decrypt_multipart_reader( + ReadEncryptionMode::Direct { base_nonce } => crate::io_support::rio::decrypt_multipart_reader( reader, material.key_bytes, - material.base_nonce, + base_nonce, part_numbers, - material.reader_backend, + crate::io_support::rio::ReadEncryptionBackend::Legacy, sequence_number, ), } } else { - match material.key_kind { - EncryptionKeyKind::Object => { + match material.mode { + ReadEncryptionMode::Object => { crate::io_support::rio::decrypt_reader_with_object_key(reader, material.key_bytes, sequence_number) } - EncryptionKeyKind::Direct => crate::io_support::rio::decrypt_reader( + ReadEncryptionMode::Direct { base_nonce } => crate::io_support::rio::decrypt_reader( reader, material.key_bytes, - material.base_nonce, - material.reader_backend, + base_nonce, + crate::io_support::rio::ReadEncryptionBackend::Legacy, sequence_number, ), } @@ -962,14 +871,28 @@ impl ReadPlan { } impl GetObjectReader { - pub async fn new( + #[cfg(test)] + pub(crate) async fn new( reader: Box, rs: Option, oi: &ObjectInfo, opts: &ObjectOptions, h: &HeaderMap, ) -> Result<(Self, usize, i64)> { - ReadPlan::build(rs, oi, opts, h).await?.into_reader(reader, oi) + Self::new_with_resolver(reader, rs, oi, opts, h, Some(&tests::TEST_RESOLVER)).await + } + + pub async fn new_with_resolver( + reader: Box, + rs: Option, + oi: &ObjectInfo, + opts: &ObjectOptions, + h: &HeaderMap, + resolver: Option<&dyn ObjectEncryptionResolver>, + ) -> Result<(Self, usize, i64)> { + ReadPlan::build_with_resolver(rs, oi, opts, h, resolver) + .await? + .into_reader(reader, oi) } pub async fn read_all(&mut self) -> Result> { let mut data = Vec::new(); @@ -1327,599 +1250,80 @@ fn multipart_part_numbers(parts: &[ObjectPartInfo]) -> Vec { parts.iter().map(|part| part.number).collect() } -fn metadata_get<'a>(metadata: &'a HashMap, key: &str) -> Option<&'a str> { - metadata.get(key).map(String::as_str).or_else(|| { - metadata - .iter() - .find_map(|(candidate, value)| candidate.eq_ignore_ascii_case(key).then_some(value.as_str())) - }) -} - -#[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) -} - -#[cfg(feature = "rio-v2")] -fn decrypt_sealed_object_key_payload(sealing_key: [u8; 32], header: &[u8], sealed_key: &[u8]) -> Result> { - let nonce = &header[4..16]; - let ciphertext = &sealed_key[DARE_HEADER_SIZE..]; - let aad = &header[..4]; - match header[1] { - DARE_CIPHER_AES_256_GCM => { - let cipher = Aes256Gcm::new_from_slice(&sealing_key) - .map_err(|err| Error::other(format!("invalid AES-GCM sealing key: {err}")))?; - let nonce = Nonce::try_from(nonce).map_err(|_| Error::other("invalid sealed object-key package nonce"))?; - cipher.decrypt(&nonce, Payload { msg: ciphertext, aad }) - } - DARE_CIPHER_CHACHA20_POLY1305 => { - let cipher = ChaCha20Poly1305::new_from_slice(&sealing_key) - .map_err(|err| Error::other(format!("invalid ChaCha20-Poly1305 sealing key: {err}")))?; - let nonce = - chacha20poly1305::Nonce::try_from(nonce).map_err(|_| Error::other("invalid sealed object-key package nonce"))?; - cipher.decrypt(&nonce, Payload { msg: ciphertext, aad }) - } - _ => return Err(Error::other("unsupported sealed object-key DARE header")), - } - .map_err(|err| Error::other(format!("failed to unseal object key: {err}"))) -} - -async fn resolve_encryption_material(oi: &ObjectInfo, headers: &HeaderMap) -> Result { - if metadata_get(&oi.user_defined, SSEC_ALGORITHM_HEADER).is_some() { - return resolve_ssec_material(oi, headers); - } - - if contains_managed_encryption_metadata(&oi.user_defined) { - return resolve_managed_material(&oi.bucket, &oi.name, &oi.user_defined).await; - } - - Err(Error::other("encrypted object metadata is incomplete")) -} - -fn contains_managed_encryption_metadata(metadata: &HashMap) -> bool { - if metadata_get(metadata, INTERNAL_ENCRYPTION_KEY_HEADER).is_some() { - return true; - } - - #[cfg(feature = "rio-v2")] - { - metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER).is_some() - || metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER).is_some() - || metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_DATA_KEY_HEADER).is_some() - } - - #[cfg(not(feature = "rio-v2"))] - { - false - } -} - -#[cfg(feature = "rio-v2")] -fn canonical_sse_path(bucket: &str, object: &str) -> String { - let bucket = bucket.trim_matches('/'); - let object = object.trim_matches('/'); - if object.is_empty() { - bucket.to_string() - } else if bucket.is_empty() { - object.to_string() - } else { - format!("{bucket}/{object}") - } -} - -#[cfg(feature = "rio-v2")] -fn managed_sse_domain(metadata: &HashMap) -> &'static str { - if metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER).is_some() - || metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER).is_some() - || matches!(metadata_get(metadata, "x-amz-server-side-encryption"), Some("aws:kms")) - { - "SSE-KMS" - } else if metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER).is_some() { - "SSE-C" - } else { - "SSE-S3" - } -} - -#[cfg(feature = "rio-v2")] -fn derive_sealing_key( - external_key: [u8; 32], - iv: [u8; SEALED_KEY_IV_SIZE], - domain: &str, - bucket: &str, - object: &str, -) -> [u8; 32] { - let mut mac = HmacSha256::new_from_slice(&external_key).expect("32-byte HMAC key"); - mac.update(&iv); - mac.update(domain.as_bytes()); - mac.update(MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.as_bytes()); - mac.update(canonical_sse_path(bucket, object).as_bytes()); - - let mut sealing_key = [0u8; 32]; - sealing_key.copy_from_slice(mac.finalize().into_bytes().as_slice()); - sealing_key -} - -#[cfg(feature = "rio-v2")] -fn try_decode_minio_sealed_key(bytes: &str) -> Result> { - let decoded = BASE64_STANDARD - .decode(bytes) - .map_err(|e| Error::other(format!("failed to decode sealed object key: {e}")))?; - match decoded.as_slice().try_into() { - Ok(sealed_key) => Ok(Some(sealed_key)), - Err(_) => Ok(None), - } -} - -#[cfg(feature = "rio-v2")] -fn try_decode_minio_sealing_iv(bytes: &str) -> Result> { - let decoded = BASE64_STANDARD - .decode(bytes) - .map_err(|e| Error::other(format!("failed to decode sealing IV: {e}")))?; - match decoded.as_slice().try_into() { - Ok(iv) => Ok(Some(iv)), - Err(_) => Ok(None), - } -} - -#[cfg(feature = "rio-v2")] -fn try_unseal_minio_object_key( - metadata: &HashMap, - bucket: &str, - object: &str, - external_key: [u8; 32], -) -> Result> { - let Some(algorithm) = metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER) else { - return Ok(None); - }; - if algorithm != MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM { - return Ok(None); - } - - let Some(iv_b64) = metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_IV_HEADER) else { - return Ok(None); - }; - let Some(iv) = try_decode_minio_sealing_iv(iv_b64)? else { - return Ok(None); - }; - - let sealed_key_b64 = metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER) - .or_else(|| metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER)) - .or_else(|| metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER)); - let Some(sealed_key_b64) = sealed_key_b64 else { - return Ok(None); - }; - let Some(sealed_key) = try_decode_minio_sealed_key(sealed_key_b64)? else { - return Ok(None); - }; - let header = &sealed_key[..DARE_HEADER_SIZE]; - if header[0] != DARE_VERSION_20 || !is_supported_sealed_object_key_cipher(header[1]) { - return Err(Error::other("unsupported sealed object-key DARE header")); - } - if u16::from_le_bytes([header[2], header[3]]) != 31 || header[4] & 0x80 == 0 { - return Err(Error::other("invalid sealed object-key payload header")); - } - - let sealing_key = derive_sealing_key(external_key, iv, managed_sse_domain(metadata), bucket, object); - let plaintext = decrypt_sealed_object_key_payload(sealing_key, header, &sealed_key)?; - let object_key: [u8; 32] = plaintext - .as_slice() - .try_into() - .map_err(|_| Error::other("sealed object key must decrypt to 32 bytes"))?; - Ok(Some(object_key)) -} - -fn resolve_ssec_material(oi: &ObjectInfo, headers: &HeaderMap) -> Result { - let algorithm = headers - .get(SSEC_ALGORITHM_HEADER) - .ok_or_else(|| Error::other("missing SSE-C algorithm header"))? - .to_str() - .map_err(|_| Error::other("invalid SSE-C algorithm header"))?; - if algorithm != DEFAULT_SSE_ALGORITHM { - return Err(Error::other(format!("unsupported SSE-C algorithm {algorithm}"))); - } - - let key_b64 = headers - .get(SSEC_KEY_HEADER) - .ok_or_else(|| Error::other("missing SSE-C key header"))? - .to_str() - .map_err(|_| Error::other("invalid SSE-C key header"))?; - let key_md5 = headers - .get(SSEC_KEY_MD5_HEADER) - .ok_or_else(|| Error::other("missing SSE-C key md5 header"))? - .to_str() - .map_err(|_| Error::other("invalid SSE-C key md5 header"))?; - - let key_bytes_vec = BASE64_STANDARD - .decode(key_b64) - .map_err(|_| Error::other("failed to decode SSE-C key"))?; - let key_bytes: [u8; 32] = key_bytes_vec - .try_into() - .map_err(|_| Error::other("SSE-C key must be 32 bytes"))?; - - let expected_md5 = BASE64_STANDARD.encode(md5_bytes(key_bytes)); - if expected_md5 != key_md5 { - return Err(Error::other("SSE-C key MD5 mismatch")); - } - - let stored_md5 = - metadata_get(&oi.user_defined, SSEC_KEY_MD5_HEADER).ok_or_else(|| Error::other("missing stored SSE-C key md5"))?; - if stored_md5 != expected_md5 { - return Err(Error::other("SSE-C key does not match object metadata")); - } - - #[cfg(feature = "rio-v2")] - if let Some(object_key) = try_unseal_minio_object_key(&oi.user_defined, &oi.bucket, &oi.name, key_bytes)? { - return Ok(EncryptionMaterial { - key_bytes: object_key, - base_nonce: [0u8; 12], - key_kind: EncryptionKeyKind::Object, - reader_backend: crate::io_support::rio::ReadEncryptionBackend::V2, - }); - } - - Ok(EncryptionMaterial { - key_bytes, - base_nonce: read_stored_ssec_nonce(&oi.user_defined, &oi.bucket, &oi.name), - key_kind: EncryptionKeyKind::Direct, - reader_backend: crate::io_support::rio::ReadEncryptionBackend::Legacy, - }) -} - -/// Resolve the SSE-C Direct base nonce for decryption. -/// -/// Since #4576 the encrypt side uses a fresh random nonce per encryption and -/// persists it under `x-rustfs-encryption-iv` (plus the MinIO interop key); -/// this reader-side resolver must read that stored value back or every SSE-C -/// GET fails its first AEAD block. Legacy objects written before random -/// nonces were persisted carry no stored IV and were encrypted with the -/// deterministic `(bucket, key)` nonce, so fall back to recomputing it. Must -/// stay in lockstep with `read_stored_ssec_nonce` in rustfs/src/storage/sse.rs -/// (the API-layer twin of this resolver). -fn read_stored_ssec_nonce(metadata: &HashMap, bucket: &str, key: &str) -> [u8; 12] { - metadata_get(metadata, INTERNAL_ENCRYPTION_IV_HEADER) - .or_else(|| metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_IV_HEADER)) - .and_then(|encoded| BASE64_STANDARD.decode(encoded).ok()) - .and_then(|bytes| <[u8; 12]>::try_from(bytes.as_slice()).ok()) - .unwrap_or_else(|| generate_ssec_nonce(bucket, key)) -} - -async fn resolve_managed_material(bucket: &str, object: &str, metadata: &HashMap) -> Result { - let normalized_metadata = normalize_managed_metadata(metadata); - let encrypted_dek = metadata_get(&normalized_metadata, INTERNAL_ENCRYPTION_KEY_HEADER) - .ok_or_else(|| Error::other("missing managed encrypted DEK"))?; - let encrypted_dek = BASE64_STANDARD - .decode(encrypted_dek) - .map_err(|e| Error::other(format!("failed to decode managed encrypted DEK: {e}")))?; - - 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()?; - #[cfg(not(feature = "rio-v2"))] - let kms_context: Option> = None; - let object_context = build_object_encryption_context(bucket, object, kms_context.as_ref()); - - // Persisted wrapping format is the read-side source of truth. The - // advertised SSE scheme and current KMS availability are write policy - // and runtime state, neither of which identifies the historical provider. - let decrypted_key = if is_data_key_envelope(&encrypted_dek) { - let service = crate::runtime::sources::object_encryption_service() - .await - .ok_or_else(|| Error::other(KmsUnavailableError))?; - #[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)? - }; - - #[cfg(feature = "rio-v2")] - if let Some(object_key) = try_unseal_minio_object_key(&normalized_metadata, bucket, object, decrypted_key)? { - return Ok(EncryptionMaterial { - key_bytes: object_key, - base_nonce: [0u8; 12], - key_kind: EncryptionKeyKind::Object, - reader_backend: crate::io_support::rio::ReadEncryptionBackend::V2, - }); - } - - let iv_b64 = metadata_get(&normalized_metadata, INTERNAL_ENCRYPTION_IV_HEADER) - .ok_or_else(|| Error::other("missing managed encryption IV"))?; - let iv = BASE64_STANDARD - .decode(iv_b64) - .map_err(|e| Error::other(format!("failed to decode managed encryption IV: {e}")))?; - let base_nonce: [u8; 12] = iv - .as_slice() - .try_into() - .map_err(|_| Error::other("managed encryption IV must be 12 bytes"))?; - - Ok(EncryptionMaterial { - key_bytes: decrypted_key, - base_nonce, - key_kind: EncryptionKeyKind::Direct, - reader_backend: crate::io_support::rio::ReadEncryptionBackend::Legacy, - }) -} - -fn normalize_managed_metadata(metadata: &HashMap) -> HashMap { - #[cfg(feature = "rio-v2")] - { - let mut normalized = metadata.clone(); - if metadata_get(&normalized, INTERNAL_ENCRYPTION_KEY_HEADER).is_none() - && let Some(value) = metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_DATA_KEY_HEADER) - .or_else(|| metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER)) - .or_else(|| metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER)) - { - normalized.insert(INTERNAL_ENCRYPTION_KEY_HEADER.to_string(), value.to_string()); - } - - if metadata_get(&normalized, INTERNAL_ENCRYPTION_IV_HEADER).is_none() - && let Some(value) = metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_IV_HEADER) - { - normalized.insert(INTERNAL_ENCRYPTION_IV_HEADER.to_string(), value.to_string()); - } - - if metadata_get(&normalized, INTERNAL_ENCRYPTION_KEY_ID_HEADER).is_none() - && let Some(value) = metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER) - { - normalized.insert(INTERNAL_ENCRYPTION_KEY_ID_HEADER.to_string(), value.to_string()); - } - - if metadata_get(&normalized, INTERNAL_ENCRYPTION_CONTEXT_HEADER).is_none() - && let Some(value) = metadata_get(metadata, MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER) - && let Ok(decoded) = BASE64_STANDARD.decode(value) - && let Ok(context) = serde_json::from_slice::>(&decoded) - && let Ok(encoded) = serde_json::to_string(&context) - { - normalized.insert(INTERNAL_ENCRYPTION_CONTEXT_HEADER.to_string(), encoded); - } - - normalized - } - - #[cfg(not(feature = "rio-v2"))] - { - metadata.clone() - } -} - -fn decrypt_local_sse_dek(encrypted_dek: &[u8], _kms_key_id: &str, object_context: &ObjectEncryptionContext) -> Result<[u8; 32]> { - if let Ok(plaintext) = decrypt_rustfs_local_sse_dek(encrypted_dek) { - return Ok(plaintext); - } - - #[cfg(feature = "rio-v2")] - { - decrypt_minio_secret_key_dek(encrypted_dek, object_context) - } - - #[cfg(not(feature = "rio-v2"))] - { - let _ = object_context; - Err(Error::other("invalid managed DEK format")) - } -} - -fn decrypt_rustfs_local_sse_dek(encrypted_dek: &[u8]) -> Result<[u8; 32]> { - let encrypted_dek = std::str::from_utf8(encrypted_dek).map_err(|_| Error::other("managed DEK is not valid UTF-8"))?; - #[derive(Deserialize)] - #[serde(deny_unknown_fields)] - struct LocalSseDekEnvelope<'a> { - version: u8, - nonce: &'a str, - ciphertext: &'a str, - } - - let (nonce, ciphertext) = match serde_json::from_str::>(encrypted_dek) { - Ok(envelope) => { - if envelope.version != LOCAL_SSE_DEK_FORMAT_VERSION { - return Err(Error::other(format!("unsupported managed DEK format version: {}", envelope.version))); - } - (envelope.nonce, envelope.ciphertext) - } - Err(_) => { - // DEPRECATED: read-only compatibility for persisted colon-delimited DEKs. - // RUSTFS_COMPAT_TODO(sse-local-dek-json-v1): Remove after all supported upgrades have rewritten legacy DEKs. - let Some((nonce, ciphertext)) = encrypted_dek.split_once(':') else { - return Err(Error::other("invalid managed DEK format")); - }; - if ciphertext.contains(':') { - return Err(Error::other("invalid managed DEK format")); - } - (nonce, ciphertext) - } - }; - - let nonce_vec = BASE64_STANDARD - .decode(nonce) - .map_err(|_| Error::other("invalid managed DEK nonce"))?; - let ciphertext = BASE64_STANDARD - .decode(ciphertext) - .map_err(|_| Error::other("invalid managed DEK ciphertext"))?; - - let nonce_array: [u8; 12] = nonce_vec - .as_slice() - .try_into() - .map_err(|_| Error::other("invalid managed DEK nonce length"))?; - - let key = Key::::from(local_sse_master_key()?); - let cipher = Aes256Gcm::new(&key); - let plaintext = cipher - .decrypt(&Nonce::from(nonce_array), ciphertext.as_slice()) - .map_err(|e| Error::other(format!("failed to decrypt managed DEK: {e}")))?; - - plaintext - .as_slice() - .try_into() - .map_err(|_| Error::other("managed DEK has invalid plaintext length")) -} - -#[cfg(feature = "rio-v2")] -#[derive(Deserialize)] -struct MinioLegacyCiphertext { - #[serde(rename = "aead")] - algorithm: String, - iv: Vec, - nonce: Vec, - bytes: Vec, -} - -#[cfg(feature = "rio-v2")] -fn decrypt_minio_secret_key_dek(encrypted_dek: &[u8], object_context: &ObjectEncryptionContext) -> Result<[u8; 32]> { - let key = local_sse_master_key()?; - let (ciphertext, iv, nonce) = parse_minio_secret_key_ciphertext(encrypted_dek)?; - let associated_data = marshal_minio_kms_context(&object_context.encryption_context); - - let mut mac = HmacSha256::new_from_slice(&key).map_err(|err| Error::other(format!("invalid local SSE master key: {err}")))?; - mac.update(&iv); - let sealing_key = mac.finalize().into_bytes(); - let cipher = Aes256Gcm::new_from_slice(sealing_key.as_slice()) - .map_err(|err| Error::other(format!("invalid MinIO sealing key: {err}")))?; - let nonce = Nonce::try_from(&nonce[..]).map_err(|_| Error::other("invalid MinIO managed DEK nonce"))?; - let plaintext = cipher - .decrypt( - &nonce, - aes_gcm::aead::Payload { - msg: &ciphertext, - aad: &associated_data, - }, - ) - .map_err(|err| Error::other(format!("failed to decrypt MinIO managed DEK: {err}")))?; - - plaintext - .as_slice() - .try_into() - .map_err(|_| Error::other("MinIO managed DEK has invalid plaintext length")) -} - -#[cfg(feature = "rio-v2")] -fn parse_minio_secret_key_ciphertext( - encrypted_dek: &[u8], -) -> Result<(Vec, [u8; MINIO_SECRET_KEY_IV_SIZE], [u8; MINIO_SECRET_KEY_NONCE_SIZE])> { - if encrypted_dek.first() == Some(&b'{') && encrypted_dek.last() == Some(&b'}') { - let legacy: MinioLegacyCiphertext = serde_json::from_slice(encrypted_dek) - .map_err(|err| Error::other(format!("failed to parse MinIO legacy managed DEK: {err}")))?; - if legacy.algorithm != "AES-256-GCM-HMAC-SHA-256" { - return Err(Error::other(format!( - "unsupported MinIO legacy managed DEK algorithm {}", - legacy.algorithm - ))); - } - let iv = legacy - .iv - .as_slice() - .try_into() - .map_err(|_| Error::other("invalid MinIO legacy managed DEK IV length"))?; - let nonce = legacy - .nonce - .as_slice() - .try_into() - .map_err(|_| Error::other("invalid MinIO legacy managed DEK nonce length"))?; - return Ok((legacy.bytes, iv, nonce)); - } - - if encrypted_dek.len() <= MINIO_SECRET_KEY_RANDOM_SIZE { - return Err(Error::other("invalid MinIO managed DEK length")); - } - - let split_at = encrypted_dek.len() - MINIO_SECRET_KEY_RANDOM_SIZE; - let (ciphertext, random) = encrypted_dek.split_at(split_at); - let iv = random[..MINIO_SECRET_KEY_IV_SIZE] - .try_into() - .map_err(|_| Error::other("invalid MinIO managed DEK IV length"))?; - let nonce = random[MINIO_SECRET_KEY_IV_SIZE..] - .try_into() - .map_err(|_| Error::other("invalid MinIO managed DEK nonce length"))?; - Ok((ciphertext.to_vec(), iv, nonce)) -} - -#[cfg(feature = "rio-v2")] -fn marshal_minio_kms_context(context: &HashMap) -> Vec { - let mut entries: Vec<_> = context.iter().collect(); - entries.sort_by_key(|(left, _)| *left); - - let mut json = String::from("{"); - for (index, (key, value)) in entries.into_iter().enumerate() { - if index > 0 { - json.push(','); - } - json.push_str(&serde_json::to_string(key).expect("string key serializes")); - json.push(':'); - json.push_str(&serde_json::to_string(value).expect("string value serializes")); - } - json.push('}'); - json.into_bytes() -} - -fn local_sse_master_key() -> Result<[u8; 32]> { - 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")? { - return Ok(key); - } - - Ok([0u8; 32]) -} - -fn decode_master_key_env(name: &str) -> Result> { - let Ok(value) = env::var(name) else { - return Ok(None); - }; - - let value = value.trim(); - if value.is_empty() { - return Ok(None); - } - - let decoded = BASE64_STANDARD - .decode(value) - .map_err(|e| Error::other(format!("{name} is not valid base64: {e}")))?; - let key = - <[u8; 32]>::try_from(decoded.as_slice()).map_err(|_| Error::other(format!("{name} must decode to exactly 32 bytes")))?; - - Ok(Some(key)) -} - -fn generate_ssec_nonce(bucket: &str, key: &str) -> [u8; 12] { - let digest = md5_bytes(format!("{bucket}-{key}").as_bytes()); - let mut nonce = [0u8; 12]; - nonce.copy_from_slice(&digest[..12]); - nonce -} - -fn md5_bytes(data: impl AsRef<[u8]>) -> [u8; 16] { - let digest = Md5::digest(data.as_ref()); - let mut out = [0u8; 16]; - out.copy_from_slice(&digest); - out -} - #[cfg(test)] mod tests { use super::*; use base64::Engine; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use md5::{Digest, Md5}; + use rustfs_utils::http::{SSEC_ALGORITHM_HEADER, SSEC_KEY_MD5_HEADER}; + use std::collections::HashMap; use std::io::Cursor; use temp_env::async_with_vars; use tokio::io::AsyncReadExt; + const TEST_DIRECT_KEY_HEADER: &str = "x-rustfs-test-direct-key"; + const TEST_OBJECT_KEY_HEADER: &str = "x-rustfs-test-object-key"; + const TEST_NONCE_HEADER: &str = "x-rustfs-test-nonce"; + + pub(super) static TEST_RESOLVER: TestObjectEncryptionResolver = TestObjectEncryptionResolver; + + pub(super) struct TestObjectEncryptionResolver; + + #[async_trait::async_trait] + impl ObjectEncryptionResolver for TestObjectEncryptionResolver { + async fn resolve_read_material( + &self, + request: ReadEncryptionRequest<'_>, + ) -> std::result::Result, EncryptionResolutionError> { + if let Some(encoded) = request.metadata.get(TEST_OBJECT_KEY_HEADER) { + let decoded = BASE64_STANDARD.decode(encoded).map_err(|_| { + EncryptionResolutionError::new(EncryptionResolutionErrorKind::InvalidMetadata, "invalid test object key") + })?; + let key_bytes = decoded.try_into().map_err(|_| { + EncryptionResolutionError::new( + EncryptionResolutionErrorKind::InvalidMetadata, + "invalid test object key length", + ) + })?; + return Ok(Some(ReadEncryptionMaterial { + key_bytes, + mode: ReadEncryptionMode::Object, + })); + } + + let encoded = request + .headers + .get(TEST_DIRECT_KEY_HEADER) + .ok_or_else(|| { + EncryptionResolutionError::new(EncryptionResolutionErrorKind::InvalidRequest, "missing test direct key") + })? + .to_str() + .map_err(|_| { + EncryptionResolutionError::new(EncryptionResolutionErrorKind::InvalidRequest, "invalid test encryption key") + })?; + let decoded = BASE64_STANDARD.decode(encoded).map_err(|_| { + EncryptionResolutionError::new(EncryptionResolutionErrorKind::InvalidRequest, "invalid test encryption key") + })?; + let key_bytes = decoded.try_into().map_err(|_| { + EncryptionResolutionError::new( + EncryptionResolutionErrorKind::InvalidRequest, + "invalid test encryption key length", + ) + })?; + let base_nonce = request + .metadata + .get(TEST_NONCE_HEADER) + .and_then(|encoded| BASE64_STANDARD.decode(encoded).ok()) + .and_then(|bytes| bytes.try_into().ok()) + .unwrap_or_else(|| fixture_nonce(request.bucket, request.object)); + Ok(Some(ReadEncryptionMaterial { + key_bytes, + mode: ReadEncryptionMode::Direct { base_nonce }, + })) + } + } + fn md5_bytes(data: impl AsRef<[u8]>) -> [u8; 16] { let digest = Md5::digest(data.as_ref()); let mut bytes = [0u8; 16]; @@ -1927,6 +1331,22 @@ mod tests { bytes } + fn fixture_nonce(bucket: &str, object: &str) -> [u8; 12] { + let digest = md5_bytes(format!("{bucket}-{object}")); + let mut nonce = [0; 12]; + nonce.copy_from_slice(&digest[..12]); + nonce + } + + fn ssec_headers_from_key(key_bytes: [u8; 32]) -> HeaderMap { + let mut headers = HeaderMap::new(); + headers.insert( + TEST_DIRECT_KEY_HEADER, + HeaderValue::from_str(&BASE64_STANDARD.encode(key_bytes)).expect("test key header is valid"), + ); + headers + } + #[tokio::test] async fn cache_body_uses_plaintext_length_for_compressed_metadata() { let mut metadata = HashMap::new(); @@ -1961,104 +1381,6 @@ mod tests { assert_eq!(restored, body); } - /// Regression for the #4576 fallout: the encrypt side persists a random - /// SSE-C nonce, and this reader-side resolver must read it back — falling - /// back to the deterministic legacy nonce only when no IV was stored. - /// Reverting the stored-nonce lookup breaks the first two cases. - #[test] - fn read_stored_ssec_nonce_prefers_persisted_iv_and_falls_back_for_legacy() { - let stored = [7u8; 12]; - let deterministic = generate_ssec_nonce("bucket", "object"); - assert_ne!(stored, deterministic, "test nonce must differ from the deterministic value"); - - let mut metadata = HashMap::new(); - metadata.insert(INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(stored)); - assert_eq!(read_stored_ssec_nonce(&metadata, "bucket", "object"), stored); - - // MinIO interop key only, in non-canonical casing: the lookup is - // case-insensitive like every other internal-metadata read here. - let mut metadata = HashMap::new(); - metadata.insert(MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_ascii_lowercase(), BASE64_STANDARD.encode(stored)); - assert_eq!(read_stored_ssec_nonce(&metadata, "bucket", "object"), stored); - - // Legacy object: no stored IV → deterministic fallback. - assert_eq!(read_stored_ssec_nonce(&HashMap::new(), "bucket", "object"), deterministic); - - // Corrupt values (bad base64 / wrong length) also fall back instead of erroring. - let mut metadata = HashMap::new(); - metadata.insert(INTERNAL_ENCRYPTION_IV_HEADER.to_string(), "not-base64!!".to_string()); - assert_eq!(read_stored_ssec_nonce(&metadata, "bucket", "object"), deterministic); - let mut metadata = HashMap::new(); - metadata.insert(INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode([1u8; 8])); - assert_eq!(read_stored_ssec_nonce(&metadata, "bucket", "object"), deterministic); - } - - fn ssec_headers_from_key(key_bytes: [u8; 32]) -> HeaderMap { - let mut headers = HeaderMap::new(); - headers.insert(SSEC_ALGORITHM_HEADER, HeaderValue::from_static("AES256")); - headers.insert( - SSEC_KEY_HEADER, - HeaderValue::from_str(&BASE64_STANDARD.encode(key_bytes)).expect("valid base64 header"), - ); - headers.insert( - SSEC_KEY_MD5_HEADER, - HeaderValue::from_str(&BASE64_STANDARD.encode(md5_bytes(key_bytes))).expect("valid md5 header"), - ); - headers - } - - #[cfg(feature = "rio-v2")] - #[test] - fn test_legacy_managed_metadata_excludes_sealed_keys() { - let legacy_metadata = HashMap::from([ - (INTERNAL_ENCRYPTION_KEY_HEADER.to_string(), "encrypted-dek".to_string()), - (INTERNAL_ENCRYPTION_IV_HEADER.to_string(), "nonce".to_string()), - ]); - assert!(is_legacy_rustfs_managed_metadata(&legacy_metadata)); - - let sealed_metadata = HashMap::from([ - (INTERNAL_ENCRYPTION_KEY_HEADER.to_string(), "encrypted-dek".to_string()), - (INTERNAL_ENCRYPTION_IV_HEADER.to_string(), "nonce".to_string()), - (MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER.to_string(), "sealed-key".to_string()), - ]); - - assert!(!is_legacy_rustfs_managed_metadata(&sealed_metadata)); - } - - #[cfg(feature = "rio-v2")] - fn seal_ssec_object_key_for_test( - bucket: &str, - object: &str, - customer_key: [u8; 32], - object_key: [u8; 32], - ) -> ([u8; 32], Vec) { - let iv = [0x23u8; SEALED_KEY_IV_SIZE]; - let sealing_key = derive_sealing_key(customer_key, iv, "SSE-C", bucket, object); - let cipher = Aes256Gcm::new_from_slice(&sealing_key).expect("valid sealing key"); - - let mut header = [0u8; DARE_HEADER_SIZE]; - header[0] = DARE_VERSION_20; - header[1] = DARE_CIPHER_AES_256_GCM; - header[2..4].copy_from_slice(&31u16.to_le_bytes()); - header[4] = 0x80; - header[5..16].copy_from_slice(&[0x45u8; 11]); - - let nonce = Nonce::try_from(&header[4..16]).expect("valid nonce"); - let mut sealed = header.to_vec(); - sealed.extend_from_slice( - &cipher - .encrypt( - &nonce, - aes_gcm::aead::Payload { - msg: &object_key, - aad: &header[..4], - }, - ) - .expect("seal object key"), - ); - (iv, sealed) - } - #[tokio::test] async fn test_ranged_decompress_reader() { // Create test data @@ -2342,261 +1664,6 @@ mod tests { assert_eq!(actual, b"fghijkl"); } - fn encrypt_managed_dek_for_test(dek: [u8; 32], master_key: [u8; 32]) -> String { - let key = Key::::from(master_key); - let cipher = Aes256Gcm::new(&key); - let nonce = Nonce::from([0u8; 12]); - let ciphertext = cipher.encrypt(&nonce, dek.as_slice()).expect("encrypt managed dek"); - serde_json::json!({ - "version": LOCAL_SSE_DEK_FORMAT_VERSION, - "nonce": BASE64_STANDARD.encode(nonce), - "ciphertext": BASE64_STANDARD.encode(ciphertext), - }) - .to_string() - } - - fn encrypt_legacy_managed_dek_for_test(dek: [u8; 32], master_key: [u8; 32]) -> String { - let key = Key::::from(master_key); - let cipher = Aes256Gcm::new(&key); - let nonce = Nonce::from([0u8; 12]); - let ciphertext = cipher.encrypt(&nonce, dek.as_slice()).expect("encrypt legacy managed dek"); - format!("{}:{}", BASE64_STANDARD.encode(nonce), BASE64_STANDARD.encode(ciphertext)) - } - - #[test] - fn decrypt_rustfs_local_sse_dek_rejects_unknown_json_version() { - let envelope = serde_json::json!({ - "version": LOCAL_SSE_DEK_FORMAT_VERSION + 1, - "nonce": BASE64_STANDARD.encode([0u8; 12]), - "ciphertext": BASE64_STANDARD.encode([0u8; 48]), - }) - .to_string(); - - let error = - decrypt_rustfs_local_sse_dek(envelope.as_bytes()).expect_err("unknown local SSE DEK versions must fail closed"); - assert!(error.to_string().contains("unsupported managed DEK format version")); - } - - #[cfg(feature = "rio-v2")] - fn seal_managed_s3_object_key_for_test( - bucket: &str, - object: &str, - data_key: [u8; 32], - object_key: [u8; 32], - ) -> ([u8; 32], Vec) { - seal_managed_s3_object_key_for_test_with_cipher(bucket, object, data_key, object_key, DARE_CIPHER_AES_256_GCM) - } - - #[cfg(feature = "rio-v2")] - fn seal_managed_s3_object_key_for_test_with_cipher( - bucket: &str, - object: &str, - data_key: [u8; 32], - object_key: [u8; 32], - cipher_id: u8, - ) -> ([u8; 32], Vec) { - let iv = [0x24u8; SEALED_KEY_IV_SIZE]; - let sealing_key = derive_sealing_key(data_key, iv, "SSE-S3", bucket, object); - - let mut header = [0u8; DARE_HEADER_SIZE]; - header[0] = DARE_VERSION_20; - header[1] = cipher_id; - header[2..4].copy_from_slice(&31u16.to_le_bytes()); - header[4] = 0x80; - header[5..16].copy_from_slice(&[0x46u8; 11]); - - let ciphertext = match cipher_id { - DARE_CIPHER_AES_256_GCM => { - let cipher = Aes256Gcm::new_from_slice(&sealing_key).expect("valid sealing key"); - let nonce = Nonce::try_from(&header[4..16]).expect("valid nonce"); - cipher - .encrypt( - &nonce, - Payload { - msg: &object_key, - aad: &header[..4], - }, - ) - .expect("seal managed object key") - } - DARE_CIPHER_CHACHA20_POLY1305 => { - let cipher = ChaCha20Poly1305::new_from_slice(&sealing_key).expect("valid sealing key"); - let nonce = chacha20poly1305::Nonce::try_from(&header[4..16]).expect("valid nonce"); - cipher - .encrypt( - &nonce, - Payload { - msg: &object_key, - aad: &header[..4], - }, - ) - .expect("seal managed object key") - } - _ => panic!("unsupported test cipher"), - }; - let mut sealed = header.to_vec(); - sealed.extend_from_slice(&ciphertext); - (iv, sealed) - } - - #[cfg(feature = "rio-v2")] - #[test] - fn test_supported_sealed_object_key_cipher_accepts_current_minio_fixture_value() { - assert!(is_supported_sealed_object_key_cipher(DARE_CIPHER_AES_256_GCM)); - assert!(is_supported_sealed_object_key_cipher(DARE_CIPHER_CHACHA20_POLY1305)); - assert!(!is_supported_sealed_object_key_cipher(0x02)); - } - - #[tokio::test] - async fn resolve_managed_material_accepts_case_insensitive_metadata_keys() { - async_with_vars([("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([0u8; 32])))], async { - let data_key = [0x24; 32]; - let base_nonce = [0x14; 12]; - let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0u8; 32]); - let metadata = HashMap::from([ - ("X-Rustfs-Encryption-Key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), - ("X-Rustfs-Encryption-IV".to_string(), BASE64_STANDARD.encode(base_nonce)), - ]); - - let material = resolve_managed_material("", "", &metadata) - .await - .expect("managed material should resolve mixed-case metadata keys"); - - assert_eq!(material.key_bytes, data_key); - assert_eq!(material.base_nonce, base_nonce); - }) - .await; - } - - #[tokio::test] - async fn resolve_managed_material_selects_provider_from_persisted_dek() { - use rustfs_kms::KmsConfig; - use tempfile::TempDir; - - let key_dir = TempDir::new().expect("create KMS key directory"); - let manager = rustfs_kms::init_global_kms_service_manager(); - manager - .reconfigure(KmsConfig::local(key_dir.path().to_path_buf()).with_insecure_development_defaults()) - .await - .expect("start test KMS service"); - - async_with_vars([("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([7u8; 32])))], async { - let data_key = [0x24; 32]; - let base_nonce = [0x14; 12]; - let encrypted_dek = encrypt_legacy_managed_dek_for_test(data_key, [7u8; 32]); - let metadata = HashMap::from([ - ( - INTERNAL_ENCRYPTION_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(encrypted_dek.as_bytes()), - ), - (INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(base_nonce)), - (INTERNAL_ENCRYPTION_KEY_ID_HEADER.to_string(), "legacy-local-key".to_string()), - ]); - - let material = resolve_managed_material("bucket", "object", &metadata) - .await - .expect("legacy local DEK should not be routed to the running KMS"); - assert_eq!(material.key_bytes, data_key); - assert_eq!(material.base_nonce, base_nonce); - }) - .await; - - manager.stop().await.expect("stop test KMS service"); - - let kms_envelope = br#"{ - "key_id": "test-key-id", - "master_key_id": "master-key-id", - "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" - }"#; - let metadata = HashMap::from([ - (INTERNAL_ENCRYPTION_KEY_HEADER.to_string(), BASE64_STANDARD.encode(kms_envelope)), - (INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode([0x14; 12])), - (INTERNAL_ENCRYPTION_KEY_ID_HEADER.to_string(), "test-key-id".to_string()), - ]); - let error = match resolve_managed_material("bucket", "object", &metadata).await { - Ok(_) => panic!("KMS envelope must not fall back to the local provider"), - Err(error) => error, - }; - let Error::Io(io_error) = error else { - panic!("KMS absence should retain its typed source"); - }; - assert!( - io_error - .get_ref() - .and_then(|source| source.downcast_ref::()) - .is_some() - ); - } - - #[cfg(feature = "rio-v2")] - #[tokio::test] - async fn resolve_managed_material_accepts_chacha20_poly1305_header_variant() { - async_with_vars([("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([0u8; 32])))], async { - let data_key = [0x24; 32]; - let object_key = [0x33; 32]; - let (iv, sealed_key) = seal_managed_s3_object_key_for_test_with_cipher( - "bucket", - "object", - data_key, - object_key, - DARE_CIPHER_CHACHA20_POLY1305, - ); - - let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0u8; 32]); - let metadata = HashMap::from([ - ( - MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(sealed_key), - ), - (MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(iv)), - ( - MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), - MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.to_string(), - ), - ( - MINIO_INTERNAL_ENCRYPTION_KMS_DATA_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(encrypted_dek.as_bytes()), - ), - (MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER.to_string(), "default".to_string()), - ]); - - let material = resolve_managed_material("bucket", "object", &metadata) - .await - .expect("managed material should accept current MinIO header variant"); - assert_eq!(material.key_kind, EncryptionKeyKind::Object); - assert_eq!(material.key_bytes, object_key); - }) - .await; - } - - #[tokio::test] - async fn resolve_encryption_material_accepts_case_insensitive_metadata_keys() { - async_with_vars([("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([0u8; 32])))], async { - let data_key = [0x24; 32]; - let base_nonce = [0x14; 12]; - let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0u8; 32]); - let metadata = HashMap::from([ - ("X-Rustfs-Encryption-Key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), - ("X-Rustfs-Encryption-IV".to_string(), BASE64_STANDARD.encode(base_nonce)), - ]); - let object_info = ObjectInfo { - user_defined: Arc::new(metadata), - ..Default::default() - }; - let material = resolve_encryption_material(&object_info, &HeaderMap::new()) - .await - .expect("resolve_encryption_material should accept mixed-case managed metadata"); - - assert_eq!(material.key_bytes, data_key); - assert_eq!(material.base_nonce, base_nonce); - }) - .await; - } - #[tokio::test] async fn test_get_object_reader_rejects_ssec_read_without_headers() { let object_info = ObjectInfo { @@ -2791,310 +1858,6 @@ mod tests { )); } - #[tokio::test] - async fn test_get_object_reader_allows_encrypted_full_object_passthrough() { - async_with_vars([("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([0u8; 32])))], async { - let plaintext = b"managed-full-object".to_vec(); - let data_key = [0x21; 32]; - let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0u8; 32]); - let bucket = "bucket"; - let object = "managed-full-object"; - - let mut encrypted = Vec::new(); - #[cfg(feature = "rio-v2")] - let user_defined = { - let object_key = [0x41; 32]; - let (sealing_iv, sealed_key) = seal_managed_s3_object_key_for_test(bucket, object, data_key, object_key); - crate::io_support::rio::EncryptReader::new_with_object_key(Cursor::new(plaintext.clone()), object_key) - .read_to_end(&mut encrypted) - .await - .expect("encrypt managed object"); - HashMap::from([ - ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), - ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), - ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ( - MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), - MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.to_string(), - ), - (MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(sealing_iv)), - ( - MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(sealed_key), - ), - ]) - }; - #[cfg(not(feature = "rio-v2"))] - let user_defined = { - let base_nonce = [0x11; 12]; - crate::io_support::rio::EncryptReader::new(Cursor::new(plaintext.clone()), data_key, base_nonce) - .read_to_end(&mut encrypted) - .await - .expect("encrypt managed object"); - HashMap::from([ - ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), - ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), - ("x-rustfs-encryption-iv".to_string(), BASE64_STANDARD.encode(base_nonce)), - ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ]) - }; - - let object_info = ObjectInfo { - bucket: bucket.to_string(), - name: object.to_string(), - size: encrypted.len() as i64, - user_defined: Arc::new(user_defined), - ..Default::default() - }; - - let (mut reader, offset, length) = GetObjectReader::new( - Box::new(Cursor::new(encrypted.clone())), - None, - &object_info, - &ObjectOptions::default(), - &HeaderMap::new(), - ) - .await - .expect("managed encrypted full-object reads should decrypt inside ecstore"); - - let mut actual = Vec::new(); - reader.read_to_end(&mut actual).await.expect("read managed plaintext"); - - assert_eq!(offset, 0); - assert_eq!(length, object_info.size); - assert_eq!(reader.object_info.size, plaintext.len() as i64); - assert_eq!(actual, plaintext); - }) - .await; - } - - #[tokio::test] - async fn test_get_object_reader_decrypts_managed_sse_range_on_plaintext_semantics() { - async_with_vars([("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([0u8; 32])))], async { - let plaintext = b"0123456789abcdefghijklmnopqrstuvwxyz".to_vec(); - let data_key = [0x23; 32]; - let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0u8; 32]); - let bucket = "bucket"; - let object = "managed-range-object"; - - let mut encrypted = Vec::new(); - #[cfg(feature = "rio-v2")] - let user_defined = { - let object_key = [0x43; 32]; - let (sealing_iv, sealed_key) = seal_managed_s3_object_key_for_test(bucket, object, data_key, object_key); - crate::io_support::rio::EncryptReader::new_with_object_key(Cursor::new(plaintext.clone()), object_key) - .read_to_end(&mut encrypted) - .await - .expect("encrypt managed ranged object"); - HashMap::from([ - ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), - ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), - ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ( - MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), - MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.to_string(), - ), - (MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(sealing_iv)), - ( - MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(sealed_key), - ), - ]) - }; - #[cfg(not(feature = "rio-v2"))] - let user_defined = { - let base_nonce = [0x13; 12]; - crate::io_support::rio::EncryptReader::new(Cursor::new(plaintext.clone()), data_key, base_nonce) - .read_to_end(&mut encrypted) - .await - .expect("encrypt managed ranged object"); - HashMap::from([ - ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), - ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), - ("x-rustfs-encryption-iv".to_string(), BASE64_STANDARD.encode(base_nonce)), - ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ]) - }; - - let object_info = ObjectInfo { - bucket: bucket.to_string(), - name: object.to_string(), - size: encrypted.len() as i64, - user_defined: Arc::new(user_defined), - ..Default::default() - }; - let range = HTTPRangeSpec { - is_suffix_length: false, - start: 5, - end: 11, - }; - - let (mut reader, offset, length) = GetObjectReader::new( - Box::new(Cursor::new(encrypted.clone())), - Some(range), - &object_info, - &ObjectOptions::default(), - &HeaderMap::new(), - ) - .await - .expect("managed encrypted range reads should decrypt inside ecstore"); - - let mut actual = Vec::new(); - reader.read_to_end(&mut actual).await.expect("read managed ranged plaintext"); - - assert_eq!(offset, 0); - assert_eq!(length, encrypted.len() as i64); - assert_eq!(reader.object_info.size, 7); - assert_eq!(actual, b"56789ab"); - }) - .await; - } - - #[tokio::test] - async fn test_get_object_reader_uses_local_managed_fallback_with_explicit_sse_s3_key() { - async_with_vars( - [ - ("__RUSTFS_SSE_SIMPLE_CMK", None::), - ("RUSTFS_SSE_S3_MASTER_KEY", Some(BASE64_STANDARD.encode([0u8; 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 bucket = "bucket"; - let object = "managed-local-fallback"; - - let mut encrypted = Vec::new(); - #[cfg(feature = "rio-v2")] - let user_defined = { - let object_key = [0x42; 32]; - let (sealing_iv, sealed_key) = seal_managed_s3_object_key_for_test(bucket, object, data_key, object_key); - crate::io_support::rio::EncryptReader::new_with_object_key(Cursor::new(plaintext.clone()), object_key) - .read_to_end(&mut encrypted) - .await - .expect("encrypt managed object with local fallback key"); - HashMap::from([ - ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), - ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), - ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ( - MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), - MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.to_string(), - ), - (MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(sealing_iv)), - ( - MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(sealed_key), - ), - ]) - }; - #[cfg(not(feature = "rio-v2"))] - let user_defined = { - let base_nonce = [0x12; 12]; - crate::io_support::rio::EncryptReader::new(Cursor::new(plaintext.clone()), data_key, base_nonce) - .read_to_end(&mut encrypted) - .await - .expect("encrypt managed object with local fallback key"); - HashMap::from([ - ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), - ("x-rustfs-encryption-key".to_string(), BASE64_STANDARD.encode(encrypted_dek.as_bytes())), - ("x-rustfs-encryption-iv".to_string(), BASE64_STANDARD.encode(base_nonce)), - ("x-rustfs-encryption-original-size".to_string(), plaintext.len().to_string()), - ]) - }; - - let object_info = ObjectInfo { - bucket: bucket.to_string(), - name: object.to_string(), - size: encrypted.len() as i64, - user_defined: Arc::new(user_defined), - ..Default::default() - }; - - let (mut reader, _, _) = GetObjectReader::new( - Box::new(Cursor::new(encrypted)), - None, - &object_info, - &ObjectOptions::default(), - &HeaderMap::new(), - ) - .await - .expect("managed encrypted reads should use the configured local SSE-S3 key"); - - let mut actual = Vec::new(); - reader.read_to_end(&mut actual).await.expect("read managed plaintext"); - - assert_eq!(reader.object_info.size, plaintext.len() as i64); - assert_eq!(actual, plaintext); - }, - ) - .await; - } - - #[cfg(feature = "rio-v2")] - #[tokio::test] - async fn test_get_object_reader_accepts_minio_only_managed_metadata() { - async_with_vars([("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([0u8; 32])))], async { - let plaintext = b"managed-minio-metadata".to_vec(); - let data_key = [0x23; 32]; - let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0u8; 32]); - let bucket = "bucket"; - let object = "managed-minio-metadata"; - let object_key = [0x44; 32]; - let (sealing_iv, sealed_key) = seal_managed_s3_object_key_for_test(bucket, object, data_key, object_key); - - let mut encrypted = Vec::new(); - crate::io_support::rio::EncryptReader::new_with_object_key(Cursor::new(plaintext.clone()), object_key) - .read_to_end(&mut encrypted) - .await - .expect("encrypt managed object"); - - let object_info = ObjectInfo { - bucket: bucket.to_string(), - name: object.to_string(), - size: encrypted.len() as i64, - user_defined: Arc::new(HashMap::from([ - ("x-amz-server-side-encryption".to_string(), "AES256".to_string()), - ( - MINIO_INTERNAL_ENCRYPTION_KMS_DATA_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(encrypted_dek.as_bytes()), - ), - ( - MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(sealed_key), - ), - (MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(sealing_iv)), - ( - MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), - MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.to_string(), - ), - (MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER.to_string(), "default".to_string()), - ("x-minio-internal-actual-size".to_string(), plaintext.len().to_string()), - ])), - ..Default::default() - }; - - let (mut reader, offset, length) = GetObjectReader::new( - Box::new(Cursor::new(encrypted.clone())), - None, - &object_info, - &ObjectOptions::default(), - &HeaderMap::new(), - ) - .await - .expect("managed encrypted reads should accept MinIO-style metadata"); - - let mut actual = Vec::new(); - reader.read_to_end(&mut actual).await.expect("read managed plaintext"); - - assert_eq!(offset, 0); - assert_eq!(length, object_info.size); - assert_eq!(reader.object_info.size, plaintext.len() as i64); - assert_eq!(actual, plaintext); - }) - .await; - } - #[tokio::test] async fn test_get_object_reader_compressed_range_returns_physical_offset_from_index() { let mut index = Index::new(); @@ -3384,8 +2147,6 @@ mod tests { let object_key = [0x67; 32]; let bucket = "bucket"; let object = "sealed-object"; - let (sealing_iv, sealed_key) = seal_ssec_object_key_for_test(bucket, object, customer_key, object_key); - let mut encrypted = Vec::new(); crate::io_support::rio::EncryptReader::new_with_object_key(Cursor::new(plaintext.clone()), object_key) .read_to_end(&mut encrypted) @@ -3397,6 +2158,7 @@ mod tests { name: object.to_string(), size: encrypted.len() as i64, user_defined: Arc::new(HashMap::from([ + (TEST_OBJECT_KEY_HEADER.to_string(), BASE64_STANDARD.encode(object_key)), ("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string()), ( "x-amz-server-side-encryption-customer-key-md5".to_string(), @@ -3406,15 +2168,6 @@ mod tests { "x-amz-server-side-encryption-customer-original-size".to_string(), plaintext.len().to_string(), ), - ( - MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), - MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.to_string(), - ), - (MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(sealing_iv)), - ( - MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(sealed_key), - ), ])), ..Default::default() }; @@ -3623,10 +2376,7 @@ mod tests { "x-amz-server-side-encryption-customer-original-size".to_string(), total_plaintext.to_string(), ), - ( - INTERNAL_ENCRYPTION_IV_HEADER.to_string(), - BASE64_STANDARD.encode(LEGACY_FIXTURE_BASE_NONCE), - ), + (TEST_NONCE_HEADER.to_string(), BASE64_STANDARD.encode(LEGACY_FIXTURE_BASE_NONCE)), ]) } @@ -4005,65 +2755,6 @@ mod tests { .await; } - #[tokio::test] - async fn test_legacy_managed_multipart_range_seek_byte_exact() { - async_with_vars( - [ - ("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([0u8; 32]))), - (ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true".to_string())), - ], - async { - let data_key = [0x74; 32]; - let encrypted_dek = encrypt_managed_dek_for_test(data_key, [0u8; 32]); - let total_plaintext: usize = 20_000 + 9_000 + 5_000; - let metadata = HashMap::from([ - ( - INTERNAL_ENCRYPTION_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(encrypted_dek.as_bytes()), - ), - (INTERNAL_ENCRYPTION_KEY_ID_HEADER.to_string(), "default".to_string()), - ( - INTERNAL_ENCRYPTION_IV_HEADER.to_string(), - BASE64_STANDARD.encode(LEGACY_FIXTURE_BASE_NONCE), - ), - (INTERNAL_ENCRYPTION_ORIGINAL_SIZE_HEADER.to_string(), total_plaintext.to_string()), - ]); - let fixture = - build_legacy_multipart_fixture("bucket", "managed-multipart", data_key, &[20_000, 9_000, 5_000], metadata) - .await; - let headers = HeaderMap::new(); - let opts = ObjectOptions::default(); - - for (rs, expected_offset, expected_length, label) in [ - ( - range(33_900, 33_999), - fixture.physical_part_start(2), - fixture.part_physical_sizes[2] as i64, - "managed tail range", - ), - ( - range(19_990, 20_010), - 0, - (fixture.part_physical_sizes[0] + fixture.part_physical_sizes[1]) as i64, - "managed boundary straddle", - ), - ] { - let (start, len) = rs.get_offset_length(total_plaintext as i64).expect("valid managed range"); - let expected_body = - &fixture.plaintext[start..start + usize::try_from(len).expect("valid managed range length fits usize")]; - - let (body, offset, length, reported_size) = read_via_seek_window(&fixture, Some(rs), &opts, &headers).await; - - assert_eq!(offset, expected_offset, "{label}: physical offset"); - assert_eq!(length, expected_length, "{label}: physical length"); - assert_eq!(reported_size, len, "{label}: reported plaintext size"); - assert_eq!(body, expected_body, "{label}: body bytes"); - } - }, - ) - .await; - } - #[tokio::test] async fn test_legacy_single_part_multipart_object_keeps_full_read_shape() { async_with_vars([(ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true"))], async { @@ -4299,50 +2990,6 @@ mod tests { .await; } - #[tokio::test] - async fn test_legacy_ssec_multipart_range_rejects_wrong_or_missing_key() { - async_with_vars([(ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true"))], async { - let key_bytes = [0x79; 32]; - let fixture = build_legacy_ssec_multipart_fixture(key_bytes, &[20_000, 9_000, 5_000]).await; - let rs = range(33_900, 33_999); - - let missing = match GetObjectReader::new( - Box::new(Cursor::new(fixture.ciphertext.clone())), - Some(rs.clone()), - &fixture.object_info, - &ObjectOptions::default(), - &HeaderMap::new(), - ) - .await - { - Ok(_) => panic!("missing SSE-C key must fail before any body is produced"), - Err(err) => err, - }; - assert!( - missing.to_string().contains("SSE-C"), - "missing-key failure must come from SSE-C validation: {missing}" - ); - - let wrong = match GetObjectReader::new( - Box::new(Cursor::new(fixture.ciphertext.clone())), - Some(rs), - &fixture.object_info, - &ObjectOptions::default(), - &ssec_headers_from_key([0x00; 32]), - ) - .await - { - Ok(_) => panic!("wrong SSE-C key must fail before any body is produced"), - Err(err) => err, - }; - assert!( - wrong.to_string().contains("SSE-C key does not match object metadata"), - "wrong-key failure must come from the stored key check: {wrong}" - ); - }) - .await; - } - #[tokio::test] async fn test_legacy_ssec_multipart_seek_tamper_fails_hard_with_no_plaintext() { async_with_vars([(ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true"))], async { @@ -4394,8 +3041,6 @@ mod tests { let object_key = [0x68; 32]; let bucket = "bucket"; let object = "large-range-object"; - let (sealing_iv, sealed_key) = seal_ssec_object_key_for_test(bucket, object, customer_key, object_key); - let mut encrypted = Vec::new(); crate::io_support::rio::EncryptReader::new_with_object_key(Cursor::new(plaintext.clone()), object_key) .read_to_end(&mut encrypted) @@ -4407,6 +3052,7 @@ mod tests { name: object.to_string(), size: encrypted.len() as i64, user_defined: Arc::new(HashMap::from([ + (TEST_OBJECT_KEY_HEADER.to_string(), BASE64_STANDARD.encode(object_key)), ("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string()), ( "x-amz-server-side-encryption-customer-key-md5".to_string(), @@ -4416,15 +3062,6 @@ mod tests { "x-amz-server-side-encryption-customer-original-size".to_string(), plaintext.len().to_string(), ), - ( - MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), - MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.to_string(), - ), - (MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(sealing_iv)), - ( - MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(sealed_key), - ), ])), ..Default::default() }; @@ -4536,7 +3173,6 @@ mod tests { let object_key = [0x74; 32]; let bucket = "bucket"; let object = "compressed-large-object"; - let (sealing_iv, sealed_key) = seal_ssec_object_key_for_test(bucket, object, customer_key, object_key); let mut compressor = crate::io_support::rio::CompressReader::with_encrypted_padding( Cursor::new(plaintext.clone()), CompressionAlgorithm::default(), @@ -4629,6 +3265,7 @@ mod tests { ..Default::default() }]), user_defined: Arc::new(HashMap::from([ + (TEST_OBJECT_KEY_HEADER.to_string(), BASE64_STANDARD.encode(object_key)), ("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string()), ( "x-amz-server-side-encryption-customer-key-md5".to_string(), @@ -4638,15 +3275,6 @@ mod tests { "x-amz-server-side-encryption-customer-original-size".to_string(), plaintext.len().to_string(), ), - ( - MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER.to_string(), - MINIO_INTERNAL_ENCRYPTION_SEAL_ALGORITHM.to_string(), - ), - (MINIO_INTERNAL_ENCRYPTION_IV_HEADER.to_string(), BASE64_STANDARD.encode(sealing_iv)), - ( - MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER.to_string(), - BASE64_STANDARD.encode(sealed_key), - ), ( "x-minio-internal-compression".to_string(), crate::io_support::rio::compression_metadata_value(CompressionAlgorithm::default()), @@ -4680,7 +3308,7 @@ mod tests { assert_eq!(plaintext_offset as i64, range.start - uncomp_off); assert_eq!(plaintext_length, 64); } - other => panic!("expected encrypted read plan, got {other:?}"), + _ => panic!("expected encrypted read plan"), } let (mut reader, offset, length) = GetObjectReader::new( diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index c9a0e560b..afe4dc7f2 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -271,29 +271,9 @@ impl ObjectInfo { } pub fn is_encrypted(&self) -> bool { - // Corresponding to the logic in rustfs/src/sse.rs/encryption_material_to_metadata function - use rustfs_utils::http::{SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER}; - - self.user_defined.keys().any(|key| { - let lower = key.to_ascii_lowercase(); - lower.starts_with("x-minio-encryption-") - || lower.starts_with("x-minio-internal-server-side-encryption-") - || matches!( - lower.as_str(), - "x-minio-internal-encrypted-multipart" - | "x-rustfs-encryption-key" - | "x-rustfs-encryption-algorithm" - | "x-rustfs-encryption-iv" - | "x-rustfs-encryption-key-id" - | "x-rustfs-encryption-context" - | "x-rustfs-encryption-tag" - | "x-amz-server-side-encryption-aws-kms-key-id" - | SSEC_ALGORITHM_HEADER - | SSEC_KEY_HEADER - | SSEC_KEY_MD5_HEADER - | "x-amz-server-side-encryption" - ) - }) + self.user_defined + .keys() + .any(|key| rustfs_utils::http::is_object_encryption_marker(key)) } /// Maximum inline size for non-versioned objects (128 KiB). @@ -337,26 +317,7 @@ impl ObjectInfo { } pub fn encryption_original_size(&self) -> std::io::Result> { - let actual_size = rustfs_utils::http::get_str(&self.user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE); - if let Some(size_str) = self - .user_defined - .get("x-rustfs-encryption-original-size") - .map(String::as_str) - .or_else(|| { - self.user_defined - .get("x-amz-server-side-encryption-customer-original-size") - .map(String::as_str) - }) - .or(actual_size.as_deref()) - && !size_str.is_empty() - { - let size = size_str - .parse::() - .map_err(|e| std::io::Error::other(format!("Failed to parse encryption original size: {e}")))?; - return Ok(Some(size)); - } - - Ok(None) + rustfs_utils::http::get_object_encryption_original_size(&self.user_defined) } pub fn decrypted_size(&self) -> std::io::Result { @@ -386,9 +347,6 @@ impl ObjectInfo { return Ok(actual_size); } - // Check if object is encrypted - // Managed SSE stores original size in x-rustfs-encryption-original-size metadata - // SSE-C stores original size in x-amz-server-side-encryption-customer-original-size if let Some(size) = self.encryption_original_size()? { return Ok(size); } @@ -878,6 +836,19 @@ mod tests { assert!(!object.is_inline_fast_path_eligible(), "transitioned objects must fall back"); } + #[test] + fn minio_internal_encryption_metadata_is_not_treated_as_plaintext() { + let object = ObjectInfo { + user_defined: Arc::new(HashMap::from([( + "X-Minio-Internal-Server-Side-Encryption-Sealed-Key".to_string(), + "sealed".to_string(), + )])), + ..Default::default() + }; + + assert!(object.is_encrypted()); + } + #[test] fn versions_after_marker_handles_null_version_marker() { let first_version = Uuid::parse_str("11111111-2222-3333-4444-555555555555").unwrap(); diff --git a/crates/ecstore/src/runtime/instance.rs b/crates/ecstore/src/runtime/instance.rs index 2c8ab72c6..ed0bda1dd 100644 --- a/crates/ecstore/src/runtime/instance.rs +++ b/crates/ecstore/src/runtime/instance.rs @@ -46,6 +46,7 @@ use crate::bucket::metadata_sys::BucketMetadataSys; use crate::bucket::replication::{DynReplicationPool, ReplicationStats}; use crate::disk::DiskStore; use crate::layout::endpoints::{EndpointServerPools, SetupType}; +use crate::object_api::ObjectEncryptionResolver; use crate::services::event_notification::EventNotifier; use crate::services::tier::tier::TierConfigMgr; use rustfs_lock::{GlobalLockManager, get_global_lock_manager}; @@ -159,6 +160,8 @@ pub struct InstanceContext { /// workers (scanner/heal/tier/lifecycle) without touching another instance. /// Replaces the process-global cancel-token static. background_cancel_token: OnceLock, + /// Resolves object-encryption material at the application boundary. + object_encryption_resolver: OnceLock>, tier_delete_journal_recovery_stores: std::sync::Mutex>, transition_transaction_recovery_stores: std::sync::Mutex>, #[cfg(test)] @@ -197,6 +200,7 @@ impl InstanceContext { local_disk_set_drives: Arc::new(RwLock::new(Vec::new())), bucket_metadata_sys: std::sync::Mutex::new(None), background_cancel_token: OnceLock::new(), + object_encryption_resolver: OnceLock::new(), tier_delete_journal_recovery_stores: std::sync::Mutex::new(HashSet::new()), transition_transaction_recovery_stores: std::sync::Mutex::new(HashSet::new()), #[cfg(test)] @@ -209,6 +213,19 @@ impl InstanceContext { self.lock_manager.clone() } + /// Install the application-owned object-encryption resolver once. + pub fn set_object_encryption_resolver( + &self, + resolver: Arc, + ) -> Result<(), Arc> { + self.object_encryption_resolver.set(resolver) + } + + /// Return the configured object-encryption resolver, if startup installed one. + pub fn object_encryption_resolver(&self) -> Option<&dyn ObjectEncryptionResolver> { + self.object_encryption_resolver.get().map(Arc::as_ref) + } + /// Set this instance's S3 region. /// /// Write-once: panics on a second write, preserving the startup fail-fast diff --git a/crates/ecstore/src/runtime/sources.rs b/crates/ecstore/src/runtime/sources.rs index 78ac568fb..8695ab397 100644 --- a/crates/ecstore/src/runtime/sources.rs +++ b/crates/ecstore/src/runtime/sources.rs @@ -46,7 +46,6 @@ use crate::{ use rustfs_concurrency::WorkloadAdmissionSnapshotProvider; use rustfs_config::server_config::{Config, get_global_server_config, set_global_server_config}; use rustfs_io_metrics::internode_metrics::global_internode_metrics; -use rustfs_kms::{ObjectEncryptionService, get_global_encryption_service}; use rustfs_lock::client::LockClient; use s3s::dto::BucketLifecycleConfiguration; use s3s::region::Region; @@ -105,10 +104,6 @@ pub(crate) fn record_erasure_write_quorum_failure(stage: &'static str, dominant_ global_internode_metrics().record_erasure_write_quorum_failure(stage, dominant_error); } -pub(crate) async fn object_encryption_service() -> Option> { - get_global_encryption_service().await -} - pub fn object_store_handle() -> Option> { resolve_object_store_handle() } diff --git a/crates/ecstore/src/set_disk/metadata.rs b/crates/ecstore/src/set_disk/metadata.rs index 136a6cf7e..a00a83ee5 100644 --- a/crates/ecstore/src/set_disk/metadata.rs +++ b/crates/ecstore/src/set_disk/metadata.rs @@ -505,9 +505,7 @@ impl SetDisks { } fn file_info_has_encryption_metadata(meta: &FileInfo) -> bool { - meta.metadata - .keys() - .any(|name| http::is_encryption_metadata_key(name) || http::is_sse_header(name)) + meta.metadata.keys().any(|name| http::is_object_encryption_marker(name)) } fn starts_with_ignore_ascii_case(value: &str, prefix: &str) -> bool { diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index bdd66c945..4f61df4b5 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -143,15 +143,17 @@ use rustfs_object_capacity::capacity_scope::{ CapacityScope, CapacityScopeDisk, current_dirty_generation, record_capacity_scope, record_global_dirty_scope, }; use rustfs_s3_types::EventName; +#[cfg(test)] +use rustfs_utils::http::SSEC_ALGORITHM_HEADER; use rustfs_utils::http::headers::AMZ_OBJECT_TAGGING; use rustfs_utils::http::headers::AMZ_STORAGE_CLASS; use rustfs_utils::http::headers::{ CACHE_CONTROL, CONTENT_DISPOSITION, CONTENT_ENCODING, CONTENT_LANGUAGE, CONTENT_TYPE, EXPIRES, HeaderExt as _, }; use rustfs_utils::http::{ - SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER, SUFFIX_ACTUAL_OBJECT_SIZE_CAP, SUFFIX_ACTUAL_SIZE, - SUFFIX_COMPRESSION, SUFFIX_COMPRESSION_SIZE, SUFFIX_REPLICATION_SSEC_CRC, SUFFIX_RESTORE_OPERATION_ID, contains_key_str, - get_header_map, get_str, insert_str, is_encryption_metadata_key, remove_header_map, + SUFFIX_ACTUAL_OBJECT_SIZE_CAP, SUFFIX_ACTUAL_SIZE, SUFFIX_COMPRESSION, SUFFIX_COMPRESSION_SIZE, SUFFIX_REPLICATION_SSEC_CRC, + SUFFIX_RESTORE_OPERATION_ID, contains_key_str, get_header_map, get_str, insert_str, is_object_encryption_marker, + remove_header_map, }; use rustfs_utils::{ HashAlgorithm, @@ -407,10 +409,7 @@ pub(crate) fn strip_internal_multipart_metadata(metadata: &mut HashMap) -> bool { - metadata.keys().any(|key| is_encryption_metadata_key(key)) - || metadata.contains_key(SSEC_ALGORITHM_HEADER) - || metadata.contains_key(SSEC_KEY_HEADER) - || metadata.contains_key(SSEC_KEY_MD5_HEADER) + metadata.keys().any(|key| is_object_encryption_marker(key)) } /// Per-set memoized capacity dirty scope. diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index d64718777..185d7c742 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -39,6 +39,7 @@ use crate::object_api::{GetObjectBodySource, get_object_body_cache_hook_suppress use crate::services::tier::tier::{TierConfigMgr, TierOperationLease}; use crate::store::ECStore; use futures::FutureExt as _; +use http::HeaderValue; use std::future::Future; fn erasure_from_file_info(fi: &FileInfo, uses_legacy: bool) -> Result { @@ -46,6 +47,17 @@ fn erasure_from_file_info(fi: &FileInfo, uses_legacy: bool) -> Result, + range: Option, + object_info: &ObjectInfo, + opts: &ObjectOptions, + headers: &HeaderMap, +) -> Result<(GetObjectReader, usize, i64)> { + GetObjectReader::new_with_resolver(reader, range, object_info, opts, headers, ctx.object_encryption_resolver()).await +} + /// Length of the full plaintext body when — and only when — this read's output /// is exactly the object's complete plaintext, so the app-layer body cache may /// serve it in place of the erasure read. @@ -704,7 +716,8 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { size_bucket, ); record_get_object_reader_path_observation(GET_OBJECT_PATH_CODEC_STREAMING, object_class, size_bucket); - let (mut reader, _offset, _length) = GetObjectReader::new(stream, range, &object_info, opts, &h).await?; + let (mut reader, _offset, _length) = + get_object_reader_with_context(&self.ctx, stream, range, &object_info, opts, &h).await?; // Carry the hook probe result so the app layer skips its // now-redundant lookup on the streaming miss path (ODC-16). reader.body_source = body_source; @@ -742,7 +755,8 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { let (rd, wd) = tokio::io::duplex(duplex_buffer_size); debug!(bucket, object, duplex_buffer_size, "Created duplex pipe for object data transfer"); - let (mut reader, offset, length) = GetObjectReader::new(Box::new(rd), range, &object_info, opts, &h).await?; + let (mut reader, offset, length) = + get_object_reader_with_context(&self.ctx, Box::new(rd), range, &object_info, opts, &h).await?; // Carry the hook probe result so the app layer skips its now-redundant // lookup on the streaming miss path (ODC-16). reader.body_source = body_source; @@ -4275,6 +4289,61 @@ mod erasure_construction_tests { } } +#[cfg(test)] +mod object_encryption_resolver_wiring_tests { + use super::*; + use crate::object_api::{EncryptionResolutionError, ObjectEncryptionResolver, ReadEncryptionMaterial, ReadEncryptionRequest}; + use std::io::Cursor; + use std::sync::atomic::{AtomicUsize, Ordering}; + + struct CountingResolver { + calls: AtomicUsize, + } + + #[async_trait::async_trait] + impl ObjectEncryptionResolver for CountingResolver { + async fn resolve_read_material( + &self, + _request: ReadEncryptionRequest<'_>, + ) -> std::result::Result, EncryptionResolutionError> { + self.calls.fetch_add(1, Ordering::Relaxed); + Ok(None) + } + } + + #[tokio::test] + async fn get_object_reader_forwards_instance_resolver() { + let resolver = Arc::new(CountingResolver { + calls: AtomicUsize::new(0), + }); + let ctx = InstanceContext::new(); + assert!( + ctx.set_object_encryption_resolver(resolver.clone()).is_ok(), + "fresh context should accept resolver" + ); + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: 1, + user_defined: Arc::new(HashMap::from([("x-amz-server-side-encryption".to_string(), "AES256".to_string())])), + ..Default::default() + }; + + let result = get_object_reader_with_context( + &ctx, + Box::new(Cursor::new(Vec::::new())), + None, + &object_info, + &ObjectOptions::default(), + &HeaderMap::new(), + ) + .await; + + assert!(result.is_err(), "resolver returning no material must fail closed"); + assert_eq!(resolver.calls.load(Ordering::Relaxed), 1); + } +} + #[cfg(test)] pub(in crate::set_disk::ops) mod hermetic_set_disks_support { //! Shared hermetic `SetDisks` construction for the ops tests below: the diff --git a/crates/ecstore/tests/README.md b/crates/ecstore/tests/README.md index 0833a467d..1514ba7f3 100644 --- a/crates/ecstore/tests/README.md +++ b/crates/ecstore/tests/README.md @@ -2,7 +2,7 @@ ## MinIO-generated encrypted fixtures -`minio_generated_read_test.rs` validates the `bitrot -> GetObjectReader` path against raw MinIO backend data captured by +`rustfs/src/storage/minio_generated_read_test.rs` validates the `bitrot -> GetObjectReader` path against raw MinIO backend data captured by `.\rustfs\scripts\minio_fixture_lab\lab.py`. It currently covers multipart fixtures for: @@ -20,5 +20,5 @@ Example: ```powershell $env:RUSTFS_MINIO_FIXTURE_ROOT = '.\rustfs\tmp\minio-fixture-lab-local-key' $env:RUSTFS_MINIO_STATIC_KMS_KEY_B64 = '' -cargo +1.97.1 test -p rustfs-ecstore --features rio-v2 --test minio_generated_read_test -- --ignored +cargo +1.97.1 test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored ``` diff --git a/crates/rio-v2/tests/minio_fixture_lab/README.md b/crates/rio-v2/tests/minio_fixture_lab/README.md index f20abe205..81ae85ef6 100644 --- a/crates/rio-v2/tests/minio_fixture_lab/README.md +++ b/crates/rio-v2/tests/minio_fixture_lab/README.md @@ -137,7 +137,7 @@ tests read): ./capture_via_docker.sh RUSTFS_MINIO_STATIC_KMS_KEY_B64=IyqsU3kMFloCNup4BsZtf/rmfHVcTgznO2F25CkEH1g= \ - cargo test -p rustfs-ecstore --features rio-v2 --test minio_generated_read_test -- --ignored + cargo test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored ``` This is exactly what the nightly `minio-interop` GitHub Actions workflow runs diff --git a/crates/utils/src/http/header_compat.rs b/crates/utils/src/http/header_compat.rs index 6dda300e6..d1bd38b82 100644 --- a/crates/utils/src/http/header_compat.rs +++ b/crates/utils/src/http/header_compat.rs @@ -25,6 +25,11 @@ 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-"; +const MINIO_INTERNAL_ENCRYPTION_PREFIX: &str = "x-minio-internal-server-side-encryption-"; +const MINIO_INTERNAL_ENCRYPTED_MULTIPART: &str = "x-minio-internal-encrypted-multipart"; +const RUSTFS_ENCRYPTION_ORIGINAL_SIZE: &str = "x-rustfs-encryption-original-size"; +const MINIO_ENCRYPTION_ORIGINAL_SIZE: &str = "x-minio-encryption-original-size"; +const SSEC_ORIGINAL_SIZE: &str = "x-amz-server-side-encryption-customer-original-size"; // Suffix constants (part after x-rustfs- or x-minio-). Use with get_header/insert_header. pub const SUFFIX_FORCE_DELETE: &str = "force-delete"; @@ -40,11 +45,49 @@ pub const SUFFIX_SOURCE_REPLICATION_REQUEST: &str = "source-replication-request" pub const SUFFIX_SOURCE_REPLICATION_CHECK: &str = "source-replication-check"; pub const SUFFIX_REPLICATION_SSEC_CRC: &str = "replication-ssec-crc"; -/// Returns true if the key is an internal encryption metadata key (x-rustfs-encryption-* or -/// x-minio-encryption-*). Case-insensitive for metadata filtering. +/// Returns true if the key is object-encryption metadata understood by RustFS or MinIO. +/// Case-insensitive for metadata filtering. pub fn is_encryption_metadata_key(key: &str) -> bool { let lower = key.to_lowercase(); - lower.starts_with(RUSTFS_ENCRYPTION_PREFIX) || lower.starts_with(MINIO_ENCRYPTION_PREFIX) + lower.starts_with(RUSTFS_ENCRYPTION_PREFIX) + || lower.starts_with(MINIO_ENCRYPTION_PREFIX) + || lower.starts_with(MINIO_INTERNAL_ENCRYPTION_PREFIX) + || lower == MINIO_INTERNAL_ENCRYPTED_MULTIPART +} + +/// Returns true when a metadata key proves that object data is encrypted. +/// +/// Original-size metadata alone is not proof: older plaintext objects can +/// retain that compatibility field after metadata migration. +pub fn is_object_encryption_marker(key: &str) -> bool { + (is_encryption_metadata_key(key) + && !key.eq_ignore_ascii_case(RUSTFS_ENCRYPTION_ORIGINAL_SIZE) + && !key.eq_ignore_ascii_case(MINIO_ENCRYPTION_ORIGINAL_SIZE)) + || super::is_sse_header(key) +} + +/// Reads the logical object size recorded by encryption metadata. +pub fn get_object_encryption_original_size(metadata: &std::collections::HashMap) -> std::io::Result> { + let actual_size = super::get_str(metadata, super::SUFFIX_ACTUAL_SIZE); + let size = get_case_insensitive(metadata, RUSTFS_ENCRYPTION_ORIGINAL_SIZE) + .or_else(|| get_case_insensitive(metadata, SSEC_ORIGINAL_SIZE)) + .or(actual_size.as_deref()); + + let Some(size) = size.filter(|size| !size.is_empty()) else { + return Ok(None); + }; + size.parse::() + .map(Some) + .map_err(|error| std::io::Error::other(format!("Failed to parse encryption original size: {error}"))) +} + +fn get_case_insensitive<'a>(metadata: &'a std::collections::HashMap, key: &str) -> Option<&'a str> { + metadata.get(key).map(String::as_str).or_else(|| { + metadata + .iter() + .find(|(candidate, _)| candidate.eq_ignore_ascii_case(key)) + .map(|(_, value)| value.as_str()) + }) } fn rustfs_key(suffix: &str) -> String { @@ -106,10 +149,37 @@ mod tests { assert!(is_encryption_metadata_key("x-rustfs-encryption-iv")); assert!(is_encryption_metadata_key("X-Rustfs-Encryption-Key")); assert!(is_encryption_metadata_key("x-minio-encryption-iv")); + assert!(is_encryption_metadata_key("X-Minio-Internal-Server-Side-Encryption-Sealed-Key")); + assert!(is_encryption_metadata_key("X-Minio-Internal-Encrypted-Multipart")); assert!(!is_encryption_metadata_key("x-amz-meta-custom")); assert!(!is_encryption_metadata_key("x-rustfs-internal-healing")); } + #[test] + fn object_encryption_marker_excludes_size_only_metadata() { + assert!(!is_object_encryption_marker(RUSTFS_ENCRYPTION_ORIGINAL_SIZE)); + assert!(is_object_encryption_marker("X-Minio-Internal-Server-Side-Encryption-Sealed-Key")); + assert!(is_object_encryption_marker("x-amz-server-side-encryption")); + } + + #[test] + fn object_encryption_original_size_is_case_insensitive() { + let metadata = std::collections::HashMap::from([( + "X-Amz-Server-Side-Encryption-Customer-Original-Size".to_string(), + "42".to_string(), + )]); + assert_eq!(get_object_encryption_original_size(&metadata).expect("valid size"), Some(42)); + } + + #[test] + fn object_encryption_original_size_prefers_rustfs_metadata() { + let metadata = std::collections::HashMap::from([ + (SSEC_ORIGINAL_SIZE.to_string(), "21".to_string()), + (RUSTFS_ENCRYPTION_ORIGINAL_SIZE.to_string(), "42".to_string()), + ]); + assert_eq!(get_object_encryption_original_size(&metadata).expect("valid size"), Some(42)); + } + #[test] fn test_get_header() { let mut headers = HeaderMap::new(); diff --git a/docs/testing/ecstore-validation-suite-design.md b/docs/testing/ecstore-validation-suite-design.md index 4eb3b41d5..62e95f340 100644 --- a/docs/testing/ecstore-validation-suite-design.md +++ b/docs/testing/ecstore-validation-suite-design.md @@ -358,7 +358,7 @@ Fixture-backed tests should run when the fixture path is present: ```bash cargo test -p rustfs-ecstore --test legacy_bitrot_read_test -- --nocapture -cargo test -p rustfs-ecstore --features rio-v2 --test minio_generated_read_test -- --ignored --nocapture +cargo test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored --nocapture ``` ## Multi-Expert Adversarial Review Summary diff --git a/crates/ecstore/tests/minio_generated_read_test.rs b/rustfs/src/storage/minio_generated_read_test.rs similarity index 96% rename from crates/ecstore/tests/minio_generated_read_test.rs rename to rustfs/src/storage/minio_generated_read_test.rs index 875aa9be9..36285c565 100644 --- a/crates/ecstore/tests/minio_generated_read_test.rs +++ b/rustfs/src/storage/minio_generated_read_test.rs @@ -4,14 +4,13 @@ use std::fs; use std::io::Cursor; use std::path::{Path, PathBuf}; -mod storage_api; - +use super::sse::SseObjectEncryptionResolver; +use super::storage_api::ecstore_test_support::{ + DiskAPI as _, DiskOption, Endpoint, Erasure, GetObjectReader, ObjectInfo, ObjectOptions, create_bitrot_reader, new_disk, +}; use rustfs_filemeta::{FileInfo, FileInfoOpts, get_file_info}; use serde::Deserialize; use sha2::{Digest, Sha256}; -use storage_api::minio_generated_read::{ - DiskAPI as _, DiskOption, Endpoint, Erasure, GetObjectReader, ObjectInfo, ObjectOptions, create_bitrot_reader, new_disk, -}; use temp_env::async_with_vars; use tokio::io::{AsyncReadExt, AsyncWrite}; @@ -49,7 +48,7 @@ impl AsyncWrite for VecAsyncWriter { fn fixture_root() -> PathBuf { std::env::var_os("RUSTFS_MINIO_FIXTURE_ROOT") .map(PathBuf::from) - .unwrap_or_else(|| PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../rio-v2/tests/fixtures/minio-generated")) + .unwrap_or_else(|| PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../crates/rio-v2/tests/fixtures/minio-generated")) } fn case_dir(case_id: &str) -> PathBuf { @@ -137,12 +136,14 @@ async fn read_fixture_plaintext(encrypted: Vec, object_info: ObjectInfo, kms ("RUSTFS_SSE_S3_MASTER_KEY", None::), ], async move { - let (mut reader, offset, length) = GetObjectReader::new( + let resolver = SseObjectEncryptionResolver; + let (mut reader, offset, length) = GetObjectReader::new_with_resolver( Box::new(Cursor::new(encrypted)), None, &object_info, &ObjectOptions::default(), &http::HeaderMap::new(), + Some(&resolver), ) .await .map_err(|err| format!("construct GetObjectReader from MinIO raw fixture: {err:?}"))?; @@ -219,11 +220,12 @@ async fn encrypted_fixture_bytes(case_dir: &Path, manifest: &ManifestRecord, fil readers.push(reader); } - let erasure = Erasure::new( + let erasure = Erasure::try_new( file_info.erasure.data_blocks, file_info.erasure.parity_blocks, file_info.erasure.block_size, - ); + ) + .expect("fixture erasure geometry"); let mut writer = VecAsyncWriter::default(); let (written, err) = erasure.decode(&mut writer, readers, 0, part.size, part.size).await; if let Some(err) = err { diff --git a/rustfs/src/storage/mod.rs b/rustfs/src/storage/mod.rs index 8b68f3aeb..30ac7ecc1 100644 --- a/rustfs/src/storage/mod.rs +++ b/rustfs/src/storage/mod.rs @@ -36,6 +36,8 @@ mod ecfs_extend; mod ecfs_test; pub(crate) mod head_prefix; #[cfg(test)] +mod minio_generated_read_test; +#[cfg(test)] mod multi_factor_scheduler_integration_test; pub(crate) mod runtime_sources; #[cfg(test)] diff --git a/rustfs/src/storage/sse.rs b/rustfs/src/storage/sse.rs index 36aa2116b..5d3c1dc17 100644 --- a/rustfs/src/storage/sse.rs +++ b/rustfs/src/storage/sse.rs @@ -70,6 +70,10 @@ //! ``` use super::StorageError; +use super::storage_api::ecstore_object::{ + EncryptionResolutionError, EncryptionResolutionErrorKind, ObjectEncryptionResolver, ReadEncryptionMaterial, + ReadEncryptionMode, ReadEncryptionRequest, +}; use crate::storage::storage_api::runtime_sources_consumer::runtime_sources; #[cfg(feature = "rio-v2")] use aes_gcm::aead::Payload; @@ -144,6 +148,7 @@ use rustfs_utils::http::headers::{ }; use rustfs_utils::path::path_join_buf; use s3s::dto::{SSECustomerAlgorithm, SSECustomerKey, SSECustomerKeyMD5, SSEKMSKeyId}; +use std::borrow::Cow; // ============================================================================ // High-Level SSE Configuration @@ -641,6 +646,23 @@ pub(crate) fn validate_sse_headers_for_read(metadata: &HashMap, } pub(crate) fn map_get_object_reader_error(err: StorageError) -> ApiError { + if let StorageError::Io(io_error) = &err + && let Some(resolution_error) = io_error + .get_ref() + .and_then(|source| source.downcast_ref::()) + { + let code = match resolution_error.kind() { + EncryptionResolutionErrorKind::InvalidRequest => S3ErrorCode::InvalidRequest, + EncryptionResolutionErrorKind::ServiceUnavailable => S3ErrorCode::ServiceUnavailable, + _ => S3ErrorCode::InternalError, + }; + return ApiError { + code, + message: resolution_error.to_string(), + source: Some(Box::new(err)), + }; + } + if let Some(message) = map_ssec_get_object_reader_error_message(&err) { return ApiError { code: S3ErrorCode::InvalidRequest, @@ -763,6 +785,104 @@ pub enum EncryptionKeyKind { Object, } +pub(crate) struct SseObjectEncryptionResolver; + +#[async_trait] +impl ObjectEncryptionResolver for SseObjectEncryptionResolver { + async fn resolve_read_material( + &self, + request: ReadEncryptionRequest<'_>, + ) -> Result, EncryptionResolutionError> { + let metadata = normalize_encryption_metadata_case(request.metadata)?; + let (_, customer_key, customer_key_md5) = + extract_ssec_params_from_headers(request.headers).map_err(map_encryption_resolution_error)?; + let material = sse_decryption(DecryptionRequest { + bucket: request.bucket, + key: request.object, + metadata: &metadata, + sse_customer_key: customer_key.as_ref(), + sse_customer_key_md5: customer_key_md5.as_ref(), + }) + .await + .map_err(map_encryption_resolution_error)?; + + Ok(material.map(|material| ReadEncryptionMaterial { + key_bytes: material.key_bytes, + mode: match material.key_kind { + EncryptionKeyKind::Direct => ReadEncryptionMode::Direct { + base_nonce: material.base_nonce, + }, + EncryptionKeyKind::Object => ReadEncryptionMode::Object, + }, + })) + } +} + +fn normalize_encryption_metadata_case( + metadata: &HashMap, +) -> Result>, EncryptionResolutionError> { + const CANONICAL_KEYS: &[&str] = &[ + "x-amz-server-side-encryption", + "x-amz-server-side-encryption-aws-kms-key-id", + "x-amz-server-side-encryption-customer-algorithm", + "x-amz-server-side-encryption-customer-key-md5", + SSEC_ORIGINAL_SIZE_HEADER, + INTERNAL_ENCRYPTION_KEY_ID_HEADER, + INTERNAL_ENCRYPTION_KEY_HEADER, + INTERNAL_ENCRYPTION_ALGORITHM_HEADER, + INTERNAL_ENCRYPTION_IV_HEADER, + "x-rustfs-encryption-context", + "x-rustfs-encryption-tag", + INTERNAL_ENCRYPTION_ORIGINAL_SIZE_HEADER, + MINIO_INTERNAL_ENCRYPTION_MULTIPART_HEADER, + MINIO_INTERNAL_ENCRYPTION_IV_HEADER, + MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER, + MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER, + MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER, + MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER, + MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER, + MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER, + ]; + + let needs_normalization = metadata.keys().any(|key| { + CANONICAL_KEYS + .iter() + .any(|canonical| key != canonical && key.eq_ignore_ascii_case(canonical)) + }); + if !needs_normalization { + return Ok(Cow::Borrowed(metadata)); + } + + let mut normalized = metadata.clone(); + for canonical in CANONICAL_KEYS { + let mut matching_values = metadata + .iter() + .filter_map(|(key, value)| key.eq_ignore_ascii_case(canonical).then_some(value)); + let Some(value) = matching_values.next() else { + continue; + }; + if matching_values.any(|candidate| candidate != value) { + return Err(EncryptionResolutionError::new( + EncryptionResolutionErrorKind::InvalidMetadata, + format!("conflicting object encryption metadata for {canonical}"), + )); + } + if !normalized.contains_key(*canonical) { + normalized.insert((*canonical).to_string(), value.clone()); + } + } + Ok(Cow::Owned(normalized)) +} + +fn map_encryption_resolution_error(error: ApiError) -> EncryptionResolutionError { + let kind = match error.code { + S3ErrorCode::InvalidArgument | S3ErrorCode::InvalidRequest => EncryptionResolutionErrorKind::InvalidRequest, + S3ErrorCode::ServiceUnavailable => EncryptionResolutionErrorKind::ServiceUnavailable, + _ => EncryptionResolutionErrorKind::DecryptionFailed, + }; + EncryptionResolutionError::new(kind, error.message) +} + #[derive(Debug, Clone)] pub struct ManagedSealedKey { #[cfg(feature = "rio-v2")] @@ -2631,19 +2751,20 @@ fn ssec_invalid_request(message: &str) -> ApiError { mod tests { use super::{ ApiError, DataKey, DecryptionRequest, EncryptionKeyKind, EncryptionMaterial, EncryptionRequest, - INTERNAL_ENCRYPTION_ALGORITHM_HEADER, INTERNAL_ENCRYPTION_IV_HEADER, INTERNAL_ENCRYPTION_KEY_HEADER, - INTERNAL_ENCRYPTION_KEY_ID_HEADER, KmsSseDekProvider, KmsUnavailableError, MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER, - MINIO_INTERNAL_ENCRYPTION_IV_HEADER, 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, apply_managed_decryption_material, - apply_managed_encryption_material, 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, - kms_operation_error, 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, + EncryptionResolutionErrorKind, INTERNAL_ENCRYPTION_ALGORITHM_HEADER, INTERNAL_ENCRYPTION_IV_HEADER, + INTERNAL_ENCRYPTION_KEY_HEADER, INTERNAL_ENCRYPTION_KEY_ID_HEADER, KmsSseDekProvider, KmsUnavailableError, + MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER, MINIO_INTERNAL_ENCRYPTION_IV_HEADER, + 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, + ObjectEncryptionResolver, PrepareEncryptionRequest, ReadEncryptionMode, ReadEncryptionRequest, SSEC_ORIGINAL_SIZE_HEADER, + SSEType, SseDekProvider, SseObjectEncryptionResolver, SsecParams, StorageError, TestSseDekProvider, + apply_managed_decryption_material, apply_managed_encryption_material, 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, kms_operation_error, 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, }; #[cfg(feature = "rio-v2")] use super::{ @@ -2703,6 +2824,97 @@ mod tests { SSE_TEST_LOCK.get_or_init(|| Mutex::new(())).lock().await } + #[tokio::test] + async fn object_encryption_resolver_returns_ssec_read_material() { + let key = [0x31; 32]; + let key_b64 = BASE64_STANDARD.encode(key); + let key_md5 = BASE64_STANDARD.encode(md5::compute(key).0); + let nonce = [0x42; 12]; + let metadata = HashMap::from([ + ("X-Amz-Server-Side-Encryption-Customer-Algorithm".to_string(), "AES256".to_string()), + ("X-Amz-Server-Side-Encryption-Customer-Key-Md5".to_string(), key_md5.clone()), + ("X-Rustfs-Encryption-Iv".to_string(), BASE64_STANDARD.encode(nonce)), + ]); + let mut headers = HeaderMap::new(); + headers.insert("x-amz-server-side-encryption-customer-algorithm", HeaderValue::from_static("AES256")); + headers.insert( + "x-amz-server-side-encryption-customer-key", + HeaderValue::from_str(&key_b64).expect("base64 key is a valid header"), + ); + headers.insert( + "x-amz-server-side-encryption-customer-key-md5", + HeaderValue::from_str(&key_md5).expect("base64 MD5 is a valid header"), + ); + + let material = SseObjectEncryptionResolver + .resolve_read_material(ReadEncryptionRequest { + bucket: "bucket", + object: "object", + metadata: &metadata, + headers: &headers, + }) + .await + .expect("SSE-C material should resolve") + .expect("SSE-C metadata should produce material"); + + assert_eq!(material.key_bytes, key); + assert_eq!(material.mode, ReadEncryptionMode::Direct { base_nonce: nonce }); + } + + #[tokio::test] + async fn object_encryption_resolver_classifies_missing_ssec_key_as_invalid_request() { + let metadata = HashMap::from([("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string())]); + let result = SseObjectEncryptionResolver + .resolve_read_material(ReadEncryptionRequest { + bucket: "bucket", + object: "object", + metadata: &metadata, + headers: &HeaderMap::new(), + }) + .await; + let error = match result { + Err(error) => error, + Ok(_) => panic!("missing SSE-C key must fail closed"), + }; + + assert_eq!(error.kind(), EncryptionResolutionErrorKind::InvalidRequest); + } + + #[tokio::test] + async fn object_encryption_resolver_rejects_conflicting_metadata_case_variants() { + let metadata = HashMap::from([ + ("x-rustfs-encryption-key".to_string(), "first".to_string()), + ("X-Rustfs-Encryption-Key".to_string(), "second".to_string()), + ]); + let result = SseObjectEncryptionResolver + .resolve_read_material(ReadEncryptionRequest { + bucket: "bucket", + object: "object", + metadata: &metadata, + headers: &HeaderMap::new(), + }) + .await; + let error = match result { + Err(error) => error, + Ok(_) => panic!("conflicting metadata aliases must fail closed"), + }; + + assert_eq!(error.kind(), EncryptionResolutionErrorKind::InvalidMetadata); + } + + #[test] + fn normalize_encryption_metadata_case_accepts_lowercase_minio_internal_keys() { + let lowercase_key = MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER.to_ascii_lowercase(); + let metadata = HashMap::from([(lowercase_key, "sealed-key".to_string())]); + + let normalized = super::normalize_encryption_metadata_case(&metadata).expect("metadata aliases should normalize"); + + assert_eq!( + normalized.get(MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER), + Some(&"sealed-key".to_string()) + ); + } + struct UnavailableSseDekProvider; #[async_trait::async_trait] @@ -3725,6 +3937,19 @@ mod tests { assert_eq!(decrypted.key_kind, EncryptionKeyKind::Object); assert_eq!(decrypted.key_bytes, material.key_bytes); + + let resolved = SseObjectEncryptionResolver + .resolve_read_material(ReadEncryptionRequest { + bucket: "bucket", + object: "object", + metadata: &metadata, + headers: &HeaderMap::new(), + }) + .await + .expect("managed resolver") + .expect("managed material"); + assert_eq!(resolved.mode, ReadEncryptionMode::Object); + assert_eq!(resolved.key_bytes, material.key_bytes); }, ) .await; @@ -3785,6 +4010,29 @@ mod tests { assert_eq!(decrypted.key_kind, EncryptionKeyKind::Object); assert_eq!(decrypted.key_bytes, material.key_bytes); + + let mut headers = HeaderMap::new(); + headers.insert("x-amz-server-side-encryption-customer-algorithm", HeaderValue::from_static("AES256")); + headers.insert( + "x-amz-server-side-encryption-customer-key", + HeaderValue::from_str(&customer_key).expect("customer key header"), + ); + headers.insert( + "x-amz-server-side-encryption-customer-key-md5", + HeaderValue::from_str(&customer_key_md5).expect("customer key MD5 header"), + ); + let resolved = SseObjectEncryptionResolver + .resolve_read_material(ReadEncryptionRequest { + bucket: "bucket", + object: "object", + metadata: &metadata, + headers: &headers, + }) + .await + .expect("SSE-C resolver") + .expect("SSE-C material"); + assert_eq!(resolved.mode, ReadEncryptionMode::Object); + assert_eq!(resolved.key_bytes, material.key_bytes); } #[cfg(feature = "rio-v2")] @@ -4715,6 +4963,15 @@ mod tests { ); } + #[test] + fn test_map_get_object_reader_error_preserves_typed_service_unavailable() { + let resolution_error = + super::EncryptionResolutionError::new(EncryptionResolutionErrorKind::ServiceUnavailable, "KMS unavailable"); + let err = map_get_object_reader_error(StorageError::other(resolution_error)); + assert_eq!(err.code, S3ErrorCode::ServiceUnavailable); + assert_eq!(err.message, "KMS unavailable"); + } + #[test] fn test_map_get_object_reader_error_leaves_non_ssec_errors_unchanged() { let err = map_get_object_reader_error(StorageError::other("plain io failure")); diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 1735ce2b4..bb6aac980 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -510,12 +510,21 @@ pub(crate) mod ecstore_object { #[cfg(test)] pub(crate) use rustfs_ecstore::api::object::GetObjectBodySource; pub(crate) use rustfs_ecstore::api::object::{ - GetObjectBodyCacheHook, GetObjectBodyCacheHookLookup, ObjectMutationHook, get_object_body_cache_plaintext_len, - lookup_get_object_body_cache_hook, register_get_object_body_cache_hook, register_object_mutation_hook, - unregister_get_object_body_cache_hook, unregister_object_mutation_hook, + EncryptionResolutionError, EncryptionResolutionErrorKind, GetObjectBodyCacheHook, GetObjectBodyCacheHookLookup, + ObjectEncryptionResolver, ObjectMutationHook, ReadEncryptionMaterial, ReadEncryptionMode, ReadEncryptionRequest, + get_object_body_cache_plaintext_len, lookup_get_object_body_cache_hook, register_get_object_body_cache_hook, + register_object_mutation_hook, unregister_get_object_body_cache_hook, unregister_object_mutation_hook, }; } +#[cfg(test)] +pub(crate) mod ecstore_test_support { + pub(crate) use rustfs_ecstore::api::bitrot::create_bitrot_reader; + pub(crate) use rustfs_ecstore::api::disk::{DiskAPI, DiskOption, endpoint::Endpoint, new_disk}; + pub(crate) use rustfs_ecstore::api::erasure::Erasure; + pub(crate) use rustfs_ecstore::api::object::{GetObjectReader, ObjectInfo, ObjectOptions}; +} + 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}; } @@ -945,13 +954,21 @@ pub(crate) async fn init_local_disks(endpoint_pools: EndpointServerPools) -> Res /// The process-level bootstrap instance context that single-instance startup /// threads through the storage foundation (Phase 5 follow-up, backlog#1052). pub(crate) fn bootstrap_instance_ctx() -> Arc { - ecstore_runtime::bootstrap_ctx() + let context = ecstore_runtime::bootstrap_ctx(); + configure_object_encryption_resolver(&context); + context } /// Construct a fresh per-server instance context (backlog#1052 S5): a second /// embedded server owns its own erasure/region/endpoint/deployment id cells. pub(crate) fn new_instance_ctx() -> Arc { - Arc::new(InstanceContext::new()) + let context = Arc::new(InstanceContext::new()); + configure_object_encryption_resolver(&context); + context +} + +fn configure_object_encryption_resolver(context: &InstanceContext) { + let _ = context.set_object_encryption_resolver(Arc::new(super::sse::SseObjectEncryptionResolver)); } pub(crate) fn init_lock_clients(endpoint_pools: EndpointServerPools) { @@ -1713,7 +1730,7 @@ pub(crate) async fn init_compression_total_memory_from_backend(store: Arc&2 exit 1