fix: restore SSE baseline on latest main (#1951)

This commit is contained in:
安正超
2026-02-25 14:19:04 +08:00
committed by GitHub
parent 62b51b5649
commit 672c255567
6 changed files with 300 additions and 610 deletions
+141 -113
View File
@@ -18,7 +18,7 @@
use crate::app::context::{AppContext, get_global_app_context};
use crate::error::ApiError;
use crate::storage::concurrency::get_concurrency_manager;
use crate::storage::ecfs::{ManagedEncryptionMaterial, RUSTFS_OWNER};
use crate::storage::ecfs::RUSTFS_OWNER;
use crate::storage::entity;
use crate::storage::helper::OperationHelper;
use crate::storage::options::{
@@ -30,7 +30,6 @@ use futures::StreamExt;
use rustfs_ecstore::StorageAPI;
use rustfs_ecstore::bucket::quota::checker::QuotaChecker;
use rustfs_ecstore::bucket::{
metadata_sys,
quota::QuotaOperation,
replication::{get_must_replicate_options, must_replicate, schedule_replication},
};
@@ -41,7 +40,7 @@ use rustfs_ecstore::new_object_layer_fn;
use rustfs_ecstore::set_disk::{MAX_PARTS_COUNT, is_valid_storage_class};
use rustfs_ecstore::store_api::{CompletePart, MultipartUploadResult, ObjectIO, ObjectOptions, PutObjReader};
use rustfs_filemeta::{ReplicationStatusType, ReplicationType};
use rustfs_rio::{CompressReader, DecryptReader, EncryptReader, HashReader, Reader, WarpReader};
use rustfs_rio::{CompressReader, HashReader, Reader, WarpReader};
use rustfs_targets::EventName;
use rustfs_utils::CompressionAlgorithm;
use rustfs_utils::http::{
@@ -51,9 +50,10 @@ use rustfs_utils::http::{
use s3s::dto::*;
use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
use std::collections::HashMap;
use std::str::FromStr;
use std::sync::Arc;
use tokio_util::io::StreamReader;
use tracing::{debug, info, instrument, warn};
use tracing::{info, instrument, warn};
pub type MultipartUsecaseResult<T> = Result<T, ApiError>;
@@ -518,72 +518,26 @@ impl DefaultMultipartUsecase {
metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), tags);
}
// TDD: Get bucket SSE configuration for multipart upload
let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok();
debug!("TDD: Got bucket SSE config for multipart: {:?}", bucket_sse_config);
let encryption_request = PrepareEncryptionRequest {
bucket: &bucket,
key: &key,
server_side_encryption,
ssekms_key_id,
sse_customer_algorithm: sse_customer_algorithm.clone(),
sse_customer_key_md5: sse_customer_key_md5.clone(),
};
// TDD: Determine effective encryption (request parameters override bucket defaults)
let original_sse = server_side_encryption.clone();
let effective_sse = server_side_encryption.or_else(|| {
bucket_sse_config.as_ref().and_then(|(config, _timestamp)| {
debug!("TDD: Processing bucket SSE config for multipart: {:?}", config);
config.rules.first().and_then(|rule| {
debug!("TDD: Processing SSE rule for multipart: {:?}", rule);
rule.apply_server_side_encryption_by_default.as_ref().map(|sse| {
debug!("TDD: Found SSE default for multipart: {:?}", sse);
match sse.sse_algorithm.as_str() {
"AES256" => ServerSideEncryption::from_static(ServerSideEncryption::AES256),
"aws:kms" => ServerSideEncryption::from_static(ServerSideEncryption::AWS_KMS),
_ => ServerSideEncryption::from_static(ServerSideEncryption::AES256),
}
})
})
})
});
debug!("TDD: effective_sse for multipart={:?} (original={:?})", effective_sse, original_sse);
let (effective_sse, effective_kms_key_id) = match sse_prepare_encryption(encryption_request).await? {
Some(material) => {
let server_side_encryption = Some(material.server_side_encryption.clone());
let ssekms_key_id = material.kms_key_id.clone();
let _original_kms_key_id = ssekms_key_id.clone();
let mut effective_kms_key_id = ssekms_key_id.or_else(|| {
bucket_sse_config.as_ref().and_then(|(config, _timestamp)| {
config.rules.first().and_then(|rule| {
rule.apply_server_side_encryption_by_default
.as_ref()
.and_then(|sse| sse.kms_master_key_id.clone())
})
})
});
metadata.extend(material.metadata);
// Store effective SSE information in metadata for multipart upload
if let Some(sse_alg) = &sse_customer_algorithm {
metadata.insert(
"x-amz-server-side-encryption-customer-algorithm".to_string(),
sse_alg.as_str().to_string(),
);
}
if let Some(sse_md5) = &sse_customer_key_md5 {
metadata.insert("x-amz-server-side-encryption-customer-key-md5".to_string(), sse_md5.clone());
}
if let Some(sse) = &effective_sse {
if is_managed_sse(sse) {
let material = create_managed_encryption_material(&bucket, &key, sse, effective_kms_key_id.clone(), 0).await?;
let ManagedEncryptionMaterial {
data_key: _,
headers,
kms_key_id: kms_key_used,
} = material;
metadata.extend(headers.into_iter());
effective_kms_key_id = Some(kms_key_used.clone());
} else {
metadata.insert("x-amz-server-side-encryption".to_string(), sse.as_str().to_string());
(server_side_encryption, ssekms_key_id)
}
}
if let Some(kms_key_id) = &effective_kms_key_id {
metadata.insert("x-amz-server-side-encryption-aws-kms-key-id".to_string(), kms_key_id.clone());
}
None => (None, None),
};
if is_compressible(&req.headers, &key) {
metadata.insert(
@@ -646,9 +600,9 @@ impl DefaultMultipartUsecase {
upload_id,
part_number,
content_length,
sse_customer_algorithm: _sse_customer_algorithm,
sse_customer_key: _sse_customer_key,
sse_customer_key_md5: _sse_customer_key_md5,
sse_customer_algorithm,
sse_customer_key,
sse_customer_key_md5,
// content_md5,
..
} = input;
@@ -691,37 +645,11 @@ impl DefaultMultipartUsecase {
};
let opts = ObjectOptions::default();
let fi = store
let mut fi = store
.get_multipart_info(&bucket, &key, &upload_id, &opts)
.await
.map_err(ApiError::from)?;
// Check if managed encryption will be applied
let will_apply_managed_encryption = decrypt_managed_encryption_key(&bucket, &key, &fi.user_defined)
.await?
.is_some();
// If managed encryption will be applied, and we have Content-Length, buffer the entire body
// This is necessary because encryption changes the data size, which causes Content-Length mismatches
if will_apply_managed_encryption && size.is_some() {
let mut total = 0i64;
let mut buffer = bytes::BytesMut::new();
while let Some(chunk) = body_stream.next().await {
let chunk = chunk.map_err(|e| ApiError::from(StorageError::other(e.to_string())))?;
total += chunk.len() as i64;
buffer.extend_from_slice(&chunk);
}
if total <= 0 {
return Err(s3_error!(UnexpectedContent));
}
size = Some(total);
let combined = buffer.freeze();
let stream = futures::stream::once(async move { Ok::<Bytes, std::io::Error>(combined) });
body_stream = StreamingBlob::wrap(stream);
}
let mut size = size.ok_or_else(|| s3_error!(UnexpectedContent))?;
// Apply adaptive buffer sizing based on part size for optimal streaming performance.
@@ -772,12 +700,51 @@ impl DefaultMultipartUsecase {
return Err(ApiError::from(StorageError::other(format!("add_checksum error={err:?}"))).into());
}
if let Some((key_bytes, base_nonce, _)) = decrypt_managed_encryption_key(&bucket, &key, &fi.user_defined).await? {
let part_nonce = derive_part_nonce(base_nonce, part_id);
let encrypt_reader = EncryptReader::new(reader, key_bytes, part_nonce);
reader = HashReader::new(Box::new(encrypt_reader), HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false)
.map_err(ApiError::from)?;
}
let server_side_encryption = fi
.user_defined
.get("x-amz-server-side-encryption")
.map(|s| {
ServerSideEncryption::from_str(s)
.map_err(|e| ApiError::from(StorageError::other(format!("Invalid server-side encryption: {e}"))))
})
.transpose()?;
let ssekms_key_id = fi
.user_defined
.get("x-amz-server-side-encryption-aws-kms-key-id")
.map(|s| s.to_string());
let part_key = fi.user_defined.get("x-rustfs-encryption-key").cloned();
let part_nonce = fi.user_defined.get("x-rustfs-encryption-iv").cloned();
let encryption_request = EncryptionRequest {
bucket: &bucket,
key: &key,
server_side_encryption,
ssekms_key_id,
sse_customer_algorithm: sse_customer_algorithm.clone(),
sse_customer_key,
sse_customer_key_md5: sse_customer_key_md5.clone(),
content_size: actual_size,
part_number: Some(part_id),
part_key,
part_nonce,
};
encryption_request.check_upload_part_customer_key_md5(&fi.user_defined, sse_customer_key_md5.clone())?;
let (requested_sse, requested_kms_key_id) = match sse_encryption(encryption_request).await? {
Some(material) => {
let requested_sse = Some(material.server_side_encryption.clone());
let requested_kms_key_id = material.kms_key_id.clone();
let encrypted_reader = material.wrap_reader(reader);
reader = HashReader::new(encrypted_reader, HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false)
.map_err(ApiError::from)?;
fi.user_defined.extend(material.metadata);
(requested_sse, requested_kms_key_id)
}
None => (None, None),
};
let mut reader = PutObjReader::new(reader);
@@ -820,6 +787,10 @@ impl DefaultMultipartUsecase {
}
let output = UploadPartOutput {
server_side_encryption: requested_sse,
ssekms_key_id: requested_kms_key_id,
sse_customer_algorithm,
sse_customer_key_md5,
checksum_crc32,
checksum_crc32c,
checksum_sha1,
@@ -988,6 +959,11 @@ impl DefaultMultipartUsecase {
upload_id,
copy_source_if_match,
copy_source_if_none_match,
sse_customer_algorithm,
sse_customer_key,
sse_customer_key_md5,
copy_source_sse_customer_key,
copy_source_sse_customer_key_md5,
..
} = req.input;
@@ -1012,7 +988,7 @@ impl DefaultMultipartUsecase {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
};
let mp_info = store
let mut mp_info = store
.get_multipart_info(&bucket, &key, &upload_id, &ObjectOptions::default())
.await
.map_err(ApiError::from)?;
@@ -1090,11 +1066,19 @@ impl DefaultMultipartUsecase {
let mut reader: Box<dyn Reader> = Box::new(WarpReader::new(src_stream));
if let Some((key_bytes, nonce, original_size_opt)) =
decrypt_managed_encryption_key(&src_bucket, &src_key, &src_info.user_defined).await?
{
reader = Box::new(DecryptReader::new(reader, key_bytes, nonce));
if let Some(original) = original_size_opt {
let src_decryption_request = DecryptionRequest {
bucket: &src_bucket,
key: &src_key,
metadata: &src_info.user_defined,
sse_customer_key: copy_source_sse_customer_key.as_ref(),
sse_customer_key_md5: copy_source_sse_customer_key_md5.as_ref(),
part_number: None,
parts: &src_info.parts,
};
if let Some(material) = sse_decryption(src_decryption_request).await? {
reader = material.wrap_single_reader(reader);
if let Some(original) = material.original_size {
src_info.actual_size = original;
}
}
@@ -1110,11 +1094,51 @@ impl DefaultMultipartUsecase {
let mut reader = HashReader::new(reader, size, actual_size, None, None, false).map_err(ApiError::from)?;
if let Some((key_bytes, base_nonce, _)) = decrypt_managed_encryption_key(&bucket, &key, &mp_info.user_defined).await? {
let part_nonce = derive_part_nonce(base_nonce, part_id);
let encrypt_reader = EncryptReader::new(reader, key_bytes, part_nonce);
reader = HashReader::new(Box::new(encrypt_reader), -1, actual_size, None, None, false).map_err(ApiError::from)?;
}
let server_side_encryption = mp_info
.user_defined
.get("x-amz-server-side-encryption")
.map(|s| {
ServerSideEncryption::from_str(s)
.map_err(|e| ApiError::from(StorageError::other(format!("Invalid server-side encryption: {e}"))))
})
.transpose()?;
let ssekms_key_id = mp_info
.user_defined
.get("x-amz-server-side-encryption-aws-kms-key-id")
.map(|s| s.to_string());
let part_key = mp_info.user_defined.get("x-rustfs-encryption-key").cloned();
let part_nonce = mp_info.user_defined.get("x-rustfs-encryption-iv").cloned();
let encryption_request = EncryptionRequest {
bucket: &bucket,
key: &key,
server_side_encryption,
ssekms_key_id,
sse_customer_algorithm: sse_customer_algorithm.clone(),
sse_customer_key,
sse_customer_key_md5: sse_customer_key_md5.clone(),
content_size: actual_size,
part_number: Some(part_id),
part_key,
part_nonce,
};
encryption_request.check_upload_part_customer_key_md5(&mp_info.user_defined, sse_customer_key_md5.clone())?;
let (requested_sse, requested_kms_key_id) = match sse_encryption(encryption_request).await? {
Some(material) => {
let requested_sse = Some(material.server_side_encryption.clone());
let requested_kms_key_id = material.kms_key_id.clone();
let encrypted_reader = material.wrap_reader(reader);
reader = HashReader::new(encrypted_reader, HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false)
.map_err(ApiError::from)?;
mp_info.user_defined.extend(material.metadata);
(requested_sse, requested_kms_key_id)
}
None => (None, None),
};
let mut reader = PutObjReader::new(reader);
@@ -1137,6 +1161,10 @@ impl DefaultMultipartUsecase {
let output = UploadPartCopyOutput {
copy_part_result: Some(copy_part_result),
copy_source_version_id: src_version_id,
server_side_encryption: requested_sse,
ssekms_key_id: requested_kms_key_id,
sse_customer_algorithm,
sse_customer_key_md5,
..Default::default()
};
+97 -312
View File
@@ -31,7 +31,6 @@ use crate::storage::options::{
};
use crate::storage::s3_api::{restore, select};
use crate::storage::*;
use base64::{Engine, engine::general_purpose::STANDARD as BASE64_STANDARD};
use bytes::Bytes;
use datafusion::arrow::{
csv::WriterBuilder as CsvWriterBuilder, json::WriterBuilder as JsonWriterBuilder, json::writer::JsonArray,
@@ -74,7 +73,7 @@ use rustfs_filemeta::{
};
use rustfs_notify::EventArgsBuilder;
use rustfs_policy::policy::action::{Action, S3Action};
use rustfs_rio::{CompressReader, DecryptReader, EncryptReader, EtagReader, HardLimitReader, HashReader, Reader, WarpReader};
use rustfs_rio::{CompressReader, EtagReader, HashReader, Reader, WarpReader};
use rustfs_s3select_api::{
object_store::bytes_stream,
query::{Context, Query},
@@ -326,7 +325,7 @@ impl DefaultObjectUsecase {
// TDD: Determine effective encryption configuration (request overrides bucket default)
let original_sse = server_side_encryption.clone();
let effective_sse = server_side_encryption.or_else(|| {
let mut effective_sse = server_side_encryption.or_else(|| {
bucket_sse_config.as_ref().and_then(|(config, _timestamp)| {
debug!("TDD: Processing bucket SSE config: {:?}", config);
config.rules.first().and_then(|rule| {
@@ -392,24 +391,6 @@ impl DefaultObjectUsecase {
metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), tags.to_string());
}
// TDD: Store effective SSE information in metadata for GET responses
if let Some(sse_alg) = &sse_customer_algorithm {
metadata.insert(
"x-amz-server-side-encryption-customer-algorithm".to_string(),
sse_alg.as_str().to_string(),
);
}
if let Some(sse_md5) = &sse_customer_key_md5 {
metadata.insert("x-amz-server-side-encryption-customer-key-md5".to_string(), sse_md5.clone());
}
if let Some(sse) = &effective_sse {
metadata.insert("x-amz-server-side-encryption".to_string(), sse.as_str().to_string());
}
if let Some(kms_key_id) = &effective_kms_key_id {
metadata.insert("x-amz-server-side-encryption-aws-kms-key-id".to_string(), kms_key_id.clone());
}
let mut opts: ObjectOptions = put_opts(&bucket, &key, version_id.clone(), &req.headers, metadata.clone())
.await
.map_err(ApiError::from)?;
@@ -463,75 +444,38 @@ impl DefaultObjectUsecase {
opts.want_checksum = reader.checksum();
}
// Apply SSE-C encryption if customer provided key
if let (Some(_), Some(sse_key), Some(sse_key_md5_provided)) =
(&sse_customer_algorithm, &sse_customer_key, &sse_customer_key_md5)
{
// Decode the base64 key
let key_bytes = BASE64_STANDARD
.decode(sse_key)
.map_err(|e| ApiError::from(StorageError::other(format!("Invalid SSE-C key: {e}"))))?;
// Apply encryption using unified SSE API.
let encryption_request = EncryptionRequest {
bucket: &bucket,
key: &key,
server_side_encryption: effective_sse.clone(),
ssekms_key_id: effective_kms_key_id.clone(),
sse_customer_algorithm: sse_customer_algorithm.clone(),
sse_customer_key,
sse_customer_key_md5: sse_customer_key_md5.clone(),
content_size: actual_size,
part_number: None,
part_key: None,
part_nonce: None,
};
// Verify key length (should be 32 bytes for AES-256)
if key_bytes.len() != 32 {
return Err(ApiError::from(StorageError::other("SSE-C key must be 32 bytes")).into());
}
if let Some(material) = sse_encryption(encryption_request).await? {
effective_sse = Some(material.server_side_encryption.clone());
effective_kms_key_id = material.kms_key_id.clone();
// Convert Vec<u8> to [u8; 32]
let mut key_array = [0u8; 32];
key_array.copy_from_slice(&key_bytes[..32]);
let encrypted_reader = material.wrap_reader(reader);
reader = HashReader::new(encrypted_reader, HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false)
.map_err(ApiError::from)?;
// Verify MD5 hash of the key matches what the client claims
let computed_md5 = BASE64_STANDARD.encode(md5::compute(&key_bytes).0);
if computed_md5 != *sse_key_md5_provided {
return Err(ApiError::from(StorageError::other("SSE-C key MD5 mismatch")).into());
}
// Store original size for later retrieval during decryption
let original_size = if size >= 0 { size } else { actual_size };
metadata.insert(
"x-amz-server-side-encryption-customer-original-size".to_string(),
original_size.to_string(),
);
// Generate a deterministic nonce from object key for consistency
let mut nonce = [0u8; 12];
let nonce_source = format!("{bucket}-{key}");
let nonce_hash = md5::compute(nonce_source.as_bytes());
nonce.copy_from_slice(&nonce_hash.0[..12]);
// Apply encryption
let encrypt_reader = EncryptReader::new(reader, key_array, nonce);
reader = HashReader::new(Box::new(encrypt_reader), -1, actual_size, None, None, false).map_err(ApiError::from)?;
}
// Apply managed SSE (SSE-S3 or SSE-KMS) when requested
if sse_customer_algorithm.is_none()
&& let Some(sse_alg) = &effective_sse
&& is_managed_sse(sse_alg)
{
let material =
create_managed_encryption_material(&bucket, &key, sse_alg, effective_kms_key_id.clone(), actual_size).await?;
let ManagedEncryptionMaterial {
data_key,
headers,
kms_key_id: kms_key_used,
} = material;
let key_bytes = data_key.plaintext_key;
let nonce = data_key.nonce;
metadata.extend(headers);
effective_kms_key_id = Some(kms_key_used.clone());
let encrypt_reader = EncryptReader::new(reader, key_bytes, nonce);
reader = HashReader::new(Box::new(encrypt_reader), -1, actual_size, None, None, false).map_err(ApiError::from)?;
let encryption_metadata = material.metadata;
metadata.extend(encryption_metadata.clone());
opts.user_defined.extend(encryption_metadata);
}
let mut reader = PutObjReader::new(reader);
let mt2 = metadata.clone();
opts.user_defined.extend(metadata);
let repoptions =
get_must_replicate_options(&mt2, "".to_string(), ReplicationStatusType::Empty, ReplicationType::Object, opts.clone());
@@ -1332,151 +1276,45 @@ impl DefaultObjectUsecase {
None
};
// Apply SSE-C decryption if customer provided key and object was encrypted with SSE-C
let mut final_stream = reader.stream;
let stored_sse_algorithm = info.user_defined.get("x-amz-server-side-encryption-customer-algorithm");
let stored_sse_key_md5 = info.user_defined.get("x-amz-server-side-encryption-customer-key-md5");
let mut managed_encryption_applied = false;
let mut managed_original_size: Option<i64> = None;
let mut response_content_length = content_length;
debug!(
"GET object metadata check: stored_sse_algorithm={:?}, stored_sse_key_md5={:?}, provided_sse_key={:?}",
stored_sse_algorithm,
stored_sse_key_md5,
"GET object metadata check: parts={}, provided_sse_key={:?}",
info.parts.len(),
req.input.sse_customer_key.is_some()
);
if stored_sse_algorithm.is_some() {
// Object was encrypted with SSE-C, so customer must provide matching key
if let (Some(sse_key), Some(sse_key_md5_provided)) = (&req.input.sse_customer_key, &req.input.sse_customer_key_md5) {
// For true multipart objects (more than 1 part), SSE-C decryption is currently not fully implemented
// Each part needs to be decrypted individually, which requires storage layer changes
// Note: Single part objects also have info.parts.len() == 1, but they are not true multipart uploads
if info.parts.len() > 1 {
warn!(
"SSE-C multipart object detected with {} parts. Currently, multipart SSE-C upload parts are not encrypted during upload_part, so no decryption is needed during GET.",
info.parts.len()
);
// Verify that the provided key MD5 matches the stored MD5 for security
if let Some(stored_md5) = stored_sse_key_md5 {
debug!("SSE-C MD5 comparison: provided='{}', stored='{}'", sse_key_md5_provided, stored_md5);
if sse_key_md5_provided != stored_md5 {
error!("SSE-C key MD5 mismatch: provided='{}', stored='{}'", sse_key_md5_provided, stored_md5);
return Err(
ApiError::from(StorageError::other("SSE-C key does not match object encryption key")).into()
);
}
} else {
return Err(ApiError::from(StorageError::other(
"Object encrypted with SSE-C but stored key MD5 not found",
))
.into());
}
// Since upload_part currently doesn't encrypt the data (SSE-C code is commented out),
// we don't need to decrypt it either. Just return the data as-is.
// TODO: Implement proper multipart SSE-C encryption/decryption
} else {
// Verify that the provided key MD5 matches the stored MD5
if let Some(stored_md5) = stored_sse_key_md5 {
debug!("SSE-C MD5 comparison: provided='{}', stored='{}'", sse_key_md5_provided, stored_md5);
if sse_key_md5_provided != stored_md5 {
error!("SSE-C key MD5 mismatch: provided='{}', stored='{}'", sse_key_md5_provided, stored_md5);
return Err(
ApiError::from(StorageError::other("SSE-C key does not match object encryption key")).into()
);
}
} else {
return Err(ApiError::from(StorageError::other(
"Object encrypted with SSE-C but stored key MD5 not found",
))
.into());
}
// Decode the base64 key
let key_bytes = BASE64_STANDARD
.decode(sse_key)
.map_err(|e| ApiError::from(StorageError::other(format!("Invalid SSE-C key: {e}"))))?;
// Verify key length (should be 32 bytes for AES-256)
if key_bytes.len() != 32 {
return Err(ApiError::from(StorageError::other("SSE-C key must be 32 bytes")).into());
}
// Convert Vec<u8> to [u8; 32]
let mut key_array = [0u8; 32];
key_array.copy_from_slice(&key_bytes[..32]);
// Verify MD5 hash of the key matches what the client claims
let computed_md5 = BASE64_STANDARD.encode(md5::compute(&key_bytes).0);
if computed_md5 != *sse_key_md5_provided {
return Err(ApiError::from(StorageError::other("SSE-C key MD5 mismatch")).into());
}
// Generate the same deterministic nonce from object key
let mut nonce = [0u8; 12];
let nonce_source = format!("{bucket}-{key}");
let nonce_hash = md5::compute(nonce_source.as_bytes());
nonce.copy_from_slice(&nonce_hash.0[..12]);
// Apply decryption
// We need to wrap the stream in a Reader first since DecryptReader expects a Reader
let warp_reader = WarpReader::new(final_stream);
let decrypt_reader = DecryptReader::new(warp_reader, key_array, nonce);
final_stream = Box::new(decrypt_reader);
}
} else {
return Err(
ApiError::from(StorageError::other("Object encrypted with SSE-C but no customer key provided")).into(),
);
}
}
if stored_sse_algorithm.is_none()
&& let Some((key_bytes, nonce, original_size)) =
decrypt_managed_encryption_key(&bucket, &key, &info.user_defined).await?
{
if info.parts.len() > 1 {
let (reader, plain_size) = decrypt_multipart_managed_stream(final_stream, &info.parts, key_bytes, nonce)
.await
.map_err(ApiError::from)?;
final_stream = reader;
managed_original_size = Some(plain_size);
} else {
let warp_reader = WarpReader::new(final_stream);
let decrypt_reader = DecryptReader::new(warp_reader, key_bytes, nonce);
final_stream = Box::new(decrypt_reader);
managed_original_size = original_size;
}
managed_encryption_applied = true;
}
// For SSE-C encrypted objects, use the original size instead of encrypted size
let response_content_length = if stored_sse_algorithm.is_some() {
if let Some(original_size_str) = info.user_defined.get("x-amz-server-side-encryption-customer-original-size") {
let original_size = original_size_str.parse::<i64>().unwrap_or(content_length);
info!(
"SSE-C decryption: using original size {} instead of encrypted size {}",
original_size, content_length
);
original_size
} else {
debug!("SSE-C decryption: no original size found, using content_length {}", content_length);
content_length
}
} else if managed_encryption_applied {
managed_original_size.unwrap_or(content_length)
} else {
content_length
let decryption_request = DecryptionRequest {
bucket: &bucket,
key: &key,
metadata: &info.user_defined,
sse_customer_key: req.input.sse_customer_key.as_ref(),
sse_customer_key_md5: req.input.sse_customer_key_md5.as_ref(),
part_number: None,
parts: &info.parts,
};
info!("Final response_content_length: {}", response_content_length);
let (server_side_encryption, sse_customer_algorithm, sse_customer_key_md5, ssekms_key_id, encryption_applied) =
match sse_decryption(decryption_request).await? {
Some(material) => {
let server_side_encryption = Some(material.server_side_encryption.clone());
let sse_customer_algorithm = Some(material.algorithm.clone());
let sse_customer_key_md5 = material.customer_key_md5.clone();
let ssekms_key_id = material.kms_key_id.clone();
if stored_sse_algorithm.is_some() || managed_encryption_applied {
let limit_reader = HardLimitReader::new(Box::new(WarpReader::new(final_stream)), response_content_length);
final_stream = Box::new(limit_reader);
}
let (decrypted_stream, plaintext_size) = material
.wrap_reader(final_stream, content_length)
.await
.map_err(ApiError::from)?;
final_stream = decrypted_stream;
response_content_length = plaintext_size;
(server_side_encryption, sse_customer_algorithm, sse_customer_key_md5, ssekms_key_id, true)
}
None => (None, None, None, None, false),
};
// Calculate concurrency-aware buffer size for optimal performance
// This adapts based on the number of concurrent GetObject requests
@@ -1506,8 +1344,7 @@ impl DefaultObjectUsecase {
&& io_strategy.cache_writeback_enabled
&& part_number.is_none()
&& rs.is_none()
&& !managed_encryption_applied
&& stored_sse_algorithm.is_none()
&& !encryption_applied
&& response_content_length > 0
&& (response_content_length as usize) <= manager.max_object_size();
@@ -1568,11 +1405,11 @@ impl DefaultObjectUsecase {
ReaderStream::with_capacity(Box::new(mem_reader), optimal_buffer_size),
response_content_length as usize,
)))
} else if stored_sse_algorithm.is_some() || managed_encryption_applied {
// For SSE-C encrypted objects, don't use bytes_stream to limit the stream
// because DecryptReader needs to read all encrypted data to produce decrypted output
} else if encryption_applied {
// For encrypted objects (SSE-C or managed SSE), avoid bytes_stream length limiting
// because DecryptReader may need to consume the full encrypted stream.
info!(
"Managed SSE: Using unlimited stream for decryption with buffer size {}",
"Encrypted object: Using unlimited stream for decryption with buffer size {}",
optimal_buffer_size
);
Some(StreamingBlob::wrap(ReaderStream::with_capacity(final_stream, optimal_buffer_size)))
@@ -1628,21 +1465,6 @@ impl DefaultObjectUsecase {
}
};
// Extract SSE information from metadata for response
let server_side_encryption = info
.user_defined
.get("x-amz-server-side-encryption")
.map(|v| ServerSideEncryption::from(v.clone()));
let sse_customer_algorithm = info
.user_defined
.get("x-amz-server-side-encryption-customer-algorithm")
.map(|v| SSECustomerAlgorithm::from(v.clone()));
let sse_customer_key_md5 = info
.user_defined
.get("x-amz-server-side-encryption-customer-key-md5")
.cloned();
let ssekms_key_id = info.user_defined.get("x-amz-server-side-encryption-aws-kms-key-id").cloned();
let mut checksum_crc32 = None;
let mut checksum_crc32c = None;
let mut checksum_sha1 = None;
@@ -2035,6 +1857,8 @@ impl DefaultObjectUsecase {
sse_customer_algorithm,
sse_customer_key,
sse_customer_key_md5,
copy_source_sse_customer_key,
copy_source_sse_customer_key_md5,
metadata_directive,
metadata,
copy_source_if_match,
@@ -2096,7 +1920,7 @@ impl DefaultObjectUsecase {
};
let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok();
let effective_sse = requested_sse.or_else(|| {
let mut effective_sse = requested_sse.or_else(|| {
bucket_sse_config.as_ref().and_then(|(config, _)| {
config.rules.first().and_then(|rule| {
rule.apply_server_side_encryption_by_default
@@ -2158,11 +1982,19 @@ impl DefaultObjectUsecase {
let mut reader: Box<dyn Reader> = Box::new(WarpReader::new(gr.stream));
if let Some((key_bytes, nonce, original_size_opt)) =
decrypt_managed_encryption_key(&src_bucket, &src_key, &src_info.user_defined).await?
{
reader = Box::new(DecryptReader::new(reader, key_bytes, nonce));
if let Some(original) = original_size_opt {
let decryption_request = DecryptionRequest {
bucket: &src_bucket,
key: &src_key,
metadata: &src_info.user_defined,
sse_customer_key: copy_source_sse_customer_key.as_ref(),
sse_customer_key_md5: copy_source_sse_customer_key_md5.as_ref(),
part_number: None,
parts: &src_info.parts,
};
if let Some(material) = sse_decryption(decryption_request).await? {
reader = material.wrap_single_reader(reader);
if let Some(original) = material.original_size {
src_info.actual_size = original;
}
}
@@ -2225,63 +2057,29 @@ impl DefaultObjectUsecase {
let mut reader = HashReader::new(reader, length, actual_size, None, None, false).map_err(ApiError::from)?;
if let Some(ref sse_alg) = effective_sse
&& is_managed_sse(sse_alg)
{
let material =
create_managed_encryption_material(&bucket, &key, sse_alg, effective_kms_key_id.clone(), actual_size).await?;
let encryption_request = EncryptionRequest {
bucket: &bucket,
key: &key,
server_side_encryption: effective_sse.clone(),
ssekms_key_id: effective_kms_key_id.clone(),
sse_customer_algorithm: sse_customer_algorithm.clone(),
sse_customer_key,
sse_customer_key_md5: sse_customer_key_md5.clone(),
content_size: actual_size,
part_number: None,
part_key: None,
part_nonce: None,
};
let ManagedEncryptionMaterial {
data_key,
headers,
kms_key_id: kms_key_used,
} = material;
if let Some(material) = sse_encryption(encryption_request).await? {
effective_sse = Some(material.server_side_encryption.clone());
effective_kms_key_id = material.kms_key_id.clone();
let key_bytes = data_key.plaintext_key;
let nonce = data_key.nonce;
let encrypted_reader = material.wrap_reader(reader);
reader = HashReader::new(encrypted_reader, HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false)
.map_err(ApiError::from)?;
src_info.user_defined.extend(headers.into_iter());
effective_kms_key_id = Some(kms_key_used.clone());
let encrypt_reader = EncryptReader::new(reader, key_bytes, nonce);
reader = HashReader::new(Box::new(encrypt_reader), -1, actual_size, None, None, false).map_err(ApiError::from)?;
}
// Apply SSE-C encryption if customer-provided key is specified
if let (Some(sse_alg), Some(sse_key), Some(sse_md5)) = (&sse_customer_algorithm, &sse_customer_key, &sse_customer_key_md5)
&& sse_alg.as_str() == "AES256"
{
let key_bytes = BASE64_STANDARD.decode(sse_key.as_str()).map_err(|e| {
error!("Failed to decode SSE-C key: {}", e);
ApiError::from(StorageError::other("Invalid SSE-C key"))
})?;
if key_bytes.len() != 32 {
return Err(ApiError::from(StorageError::other("SSE-C key must be 32 bytes")).into());
}
let computed_md5 = BASE64_STANDARD.encode(md5::compute(&key_bytes).0);
if computed_md5 != sse_md5.as_str() {
return Err(ApiError::from(StorageError::other("SSE-C key MD5 mismatch")).into());
}
// Store original size before encryption
src_info
.user_defined
.insert("x-amz-server-side-encryption-customer-original-size".to_string(), actual_size.to_string());
let key_array: [u8; 32] = key_bytes
.try_into()
.map_err(|_| ApiError::from(StorageError::other("SSE-C key must be 32 bytes")))?;
// Generate deterministic nonce from bucket-key
let nonce_source = format!("{bucket}-{key}");
let nonce_hash = md5::compute(nonce_source.as_bytes());
let nonce: [u8; 12] = nonce_hash.0[..12]
.try_into()
.map_err(|_| ApiError::from(StorageError::other("Failed to derive SSE-C nonce")))?;
let encrypt_reader = EncryptReader::new(reader, key_array, nonce);
reader = HashReader::new(Box::new(encrypt_reader), -1, actual_size, None, None, false).map_err(ApiError::from)?;
src_info.user_defined.extend(material.metadata);
}
src_info.put_object_reader = Some(PutObjReader::new(reader));
@@ -2292,19 +2090,6 @@ impl DefaultObjectUsecase {
src_info.user_defined.insert(k, v);
}
// Store SSE-C metadata for GET responses
if let Some(ref sse_alg) = sse_customer_algorithm {
src_info.user_defined.insert(
"x-amz-server-side-encryption-customer-algorithm".to_string(),
sse_alg.as_str().to_string(),
);
}
if let Some(ref sse_md5) = sse_customer_key_md5 {
src_info
.user_defined
.insert("x-amz-server-side-encryption-customer-key-md5".to_string(), sse_md5.clone());
}
// check quota for copy operation
if let Some(metadata_sys) = rustfs_ecstore::bucket::metadata_sys::GLOBAL_BucketMetadataSys.get() {
let quota_checker = QuotaChecker::new(metadata_sys.clone());
-7
View File
@@ -34,7 +34,6 @@ use rustfs_ecstore::{
// RESERVED_METADATA_PREFIX,
},
};
use rustfs_kms::DataKey;
use rustfs_notify::notifier_global;
use rustfs_rio::{CompressReader, HashReader, Reader, WarpReader};
use rustfs_targets::EventName;
@@ -565,12 +564,6 @@ pub struct FS {
// pub store: ECStore,
}
pub(crate) struct ManagedEncryptionMaterial {
pub(crate) data_key: DataKey,
pub(crate) headers: HashMap<String, String>,
pub(crate) kms_key_id: String,
}
#[derive(Debug, Default, serde::Deserialize)]
pub(crate) struct ListObjectUnorderedQuery {
#[serde(rename = "allow-unordered")]
+2 -168
View File
@@ -17,7 +17,7 @@ use crate::config::workload_profiles::{
};
use crate::error::ApiError;
use crate::server::cors;
use crate::storage::ecfs::{InMemoryAsyncReader, ListObjectUnorderedQuery};
use crate::storage::ecfs::ListObjectUnorderedQuery;
use http::{HeaderMap, HeaderValue, StatusCode};
use metrics::counter;
use rustfs_ecstore::bucket::metadata_sys;
@@ -27,9 +27,6 @@ use rustfs_ecstore::bucket::replication::ReplicationConfigurationExt;
use rustfs_ecstore::error::StorageError;
use rustfs_ecstore::store_api::{BucketOptions, ObjectInfo, ObjectToDelete};
use rustfs_ecstore::{StorageAPI, new_object_layer_fn};
use rustfs_filemeta::ObjectPartInfo;
use rustfs_kms::{EncryptionMetadata, ObjectEncryptionContext, get_global_encryption_service};
use rustfs_rio::{DecryptReader, Reader, WarpReader};
use rustfs_targets::EventName;
use rustfs_targets::arn::{TargetID, TargetIDError};
use rustfs_utils::http::{
@@ -39,7 +36,7 @@ use rustfs_utils::http::{
use s3s::dto::{
Delimiter, LambdaFunctionConfiguration, NotificationConfigurationFilter, ObjectLockConfiguration, ObjectLockEnabled,
ObjectLockLegalHold, ObjectLockLegalHoldStatus, ObjectLockRetention, ObjectLockRetentionMode, QueueConfiguration,
ServerSideEncryption, TopicConfiguration,
TopicConfiguration,
};
use s3s::{S3Error, S3ErrorCode, S3Response, S3Result};
use serde_urlencoded::from_bytes;
@@ -49,7 +46,6 @@ use std::sync::Arc;
use time::OffsetDateTime;
use time::format_description::well_known::Rfc3339;
use time::{format_description::FormatItem, macros::format_description};
use tokio::io::AsyncRead;
use tracing::{debug, warn};
pub const RFC1123: &[FormatItem<'_>] =
@@ -192,168 +188,6 @@ pub(crate) fn get_buffer_size_opt_in(file_size: i64) -> usize {
buffer_size
}
pub(crate) async fn create_managed_encryption_material(
bucket: &str,
key: &str,
algorithm: &ServerSideEncryption,
kms_key_id: Option<String>,
original_size: i64,
) -> Result<crate::storage::ecfs::ManagedEncryptionMaterial, ApiError> {
let Some(service) = get_global_encryption_service().await else {
return Err(ApiError::from(StorageError::other("KMS encryption service is not initialized")));
};
if !is_managed_sse(algorithm) {
return Err(ApiError::from(StorageError::other(format!(
"Unsupported server-side encryption algorithm: {}",
algorithm.as_str()
))));
}
let algorithm_str = algorithm.as_str();
let mut context = ObjectEncryptionContext::new(bucket.to_string(), key.to_string());
if original_size >= 0 {
context = context.with_size(original_size as u64);
}
let mut kms_key_candidate = kms_key_id;
if kms_key_candidate.is_none() {
kms_key_candidate = service.get_default_key_id().cloned();
}
let kms_key_to_use = kms_key_candidate
.clone()
.ok_or_else(|| ApiError::from(StorageError::other("No KMS key available for managed server-side encryption")))?;
let (data_key, encrypted_data_key) = service
.create_data_key(&kms_key_candidate, &context)
.await
.map_err(|e| ApiError::from(StorageError::other(format!("Failed to create data key: {e}"))))?;
let metadata = EncryptionMetadata {
algorithm: algorithm_str.to_string(),
key_id: kms_key_to_use.clone(),
key_version: 1,
iv: data_key.nonce.to_vec(),
tag: None,
encryption_context: context.encryption_context.clone(),
encrypted_at: jiff::Zoned::now(),
original_size: if original_size >= 0 { original_size as u64 } else { 0 },
encrypted_data_key,
};
let mut headers = service.metadata_to_headers(&metadata);
headers.insert("x-rustfs-encryption-original-size".to_string(), metadata.original_size.to_string());
Ok(crate::storage::ecfs::ManagedEncryptionMaterial {
data_key,
headers,
kms_key_id: kms_key_to_use,
})
}
pub(crate) async fn decrypt_managed_encryption_key(
bucket: &str,
key: &str,
metadata: &HashMap<String, String>,
) -> Result<Option<([u8; 32], [u8; 12], Option<i64>)>, ApiError> {
if !metadata.contains_key("x-rustfs-encryption-key") {
return Ok(None);
}
let Some(service) = get_global_encryption_service().await else {
return Err(ApiError::from(StorageError::other("KMS encryption service is not initialized")));
};
let parsed = service
.headers_to_metadata(metadata)
.map_err(|e| ApiError::from(StorageError::other(format!("Failed to parse encryption metadata: {e}"))))?;
if parsed.iv.len() != 12 {
return Err(ApiError::from(StorageError::other("Invalid encryption nonce length; expected 12 bytes")));
}
let context = ObjectEncryptionContext::new(bucket.to_string(), key.to_string());
let data_key = service
.decrypt_data_key(&parsed.encrypted_data_key, &context)
.await
.map_err(|e| ApiError::from(StorageError::other(format!("Failed to decrypt data key: {e}"))))?;
let key_bytes = data_key.plaintext_key;
let mut nonce = [0u8; 12];
nonce.copy_from_slice(&parsed.iv[..12]);
let original_size = metadata
.get("x-rustfs-encryption-original-size")
.and_then(|s| s.parse::<i64>().ok());
Ok(Some((key_bytes, nonce, original_size)))
}
pub(crate) fn derive_part_nonce(base: [u8; 12], part_number: usize) -> [u8; 12] {
let mut nonce = base;
let current = u32::from_be_bytes([nonce[8], nonce[9], nonce[10], nonce[11]]);
let incremented = current.wrapping_add(part_number as u32);
nonce[8..12].copy_from_slice(&incremented.to_be_bytes());
nonce
}
pub(crate) async fn decrypt_multipart_managed_stream(
mut encrypted_stream: Box<dyn AsyncRead + Unpin + Send + Sync>,
parts: &[ObjectPartInfo],
key_bytes: [u8; 32],
base_nonce: [u8; 12],
) -> Result<(Box<dyn Reader>, i64), StorageError> {
let total_plain_capacity: usize = parts.iter().map(|part| part.actual_size.max(0) as usize).sum();
let mut plaintext = Vec::with_capacity(total_plain_capacity);
for part in parts {
if part.size == 0 {
continue;
}
let mut encrypted_part = vec![0u8; part.size];
tokio::io::AsyncReadExt::read_exact(&mut encrypted_stream, &mut encrypted_part)
.await
.map_err(|e| StorageError::other(format!("failed to read encrypted multipart segment {}: {}", part.number, e)))?;
let part_nonce = derive_part_nonce(base_nonce, part.number);
let cursor = std::io::Cursor::new(encrypted_part);
let mut decrypt_reader = DecryptReader::new(WarpReader::new(cursor), key_bytes, part_nonce);
tokio::io::AsyncReadExt::read_to_end(&mut decrypt_reader, &mut plaintext)
.await
.map_err(|e| StorageError::other(format!("failed to decrypt multipart segment {}: {}", part.number, e)))?;
}
let total_plain_size = plaintext.len() as i64;
let reader = Box::new(WarpReader::new(InMemoryAsyncReader::new(plaintext))) as Box<dyn Reader>;
Ok((reader, total_plain_size))
}
pub(crate) fn strip_managed_encryption_metadata(metadata: &mut HashMap<String, String>) {
const KEYS: [&str; 7] = [
"x-amz-server-side-encryption",
"x-amz-server-side-encryption-aws-kms-key-id",
"x-rustfs-encryption-iv",
"x-rustfs-encryption-tag",
"x-rustfs-encryption-key",
"x-rustfs-encryption-context",
"x-rustfs-encryption-original-size",
];
for key in KEYS.iter() {
metadata.remove(*key);
}
}
pub(crate) fn is_managed_sse(algorithm: &ServerSideEncryption) -> bool {
matches!(algorithm.as_str(), "AES256" | "aws:kms")
}
/// Validate object key for control characters and log special characters
///
/// This function:
+8
View File
@@ -18,8 +18,10 @@ pub mod ecfs;
pub(crate) mod entity;
pub(crate) mod helper;
pub mod options;
pub(crate) mod readers;
pub mod rpc;
pub(crate) mod s3_api;
mod sse;
pub mod tonic_service;
#[cfg(test)]
@@ -28,5 +30,11 @@ mod ecfs_extend;
#[cfg(test)]
mod ecfs_test;
pub(crate) mod head_prefix;
#[cfg(test)]
mod sse_test;
pub(crate) use ecfs_extend::*;
pub(crate) use sse::{
DecryptionRequest, EncryptionRequest, PrepareEncryptionRequest, sse_decryption, sse_encryption, sse_prepare_encryption,
strip_managed_encryption_metadata,
};
+52 -10
View File
@@ -322,7 +322,7 @@ impl EncryptionRequest<'_> {
// if customer_key_md5 is provided, check if it matches the metadata
let customer_key_md5_from_metadata = user_defined.get("x-amz-server-side-encryption-customer-key-md5");
if let Some(customer_key_md5_from_metadata) = customer_key_md5_from_metadata
&& !customer_key_md5_from_metadata.eq_ignore_ascii_case(customer_key_md5.as_str())
&& customer_key_md5_from_metadata != customer_key_md5.as_str()
{
return Err(ApiError::from(StorageError::other("Customer key MD5 mismatch")));
}
@@ -1324,15 +1324,6 @@ pub async fn get_sse_dek_provider() -> Result<Arc<dyn SseDekProvider>, ApiError>
Ok(provider)
}
// check encryption metadata
pub fn check_encryption_metadata(metadata: &HashMap<String, String>) -> bool {
if !metadata.contains_key("x-rustfs-encryption-key") && !metadata.contains_key("x-amz-server-side-encryption") {
return false;
}
true
}
/// Reset the global SSE DEK provider (for testing only)
///
/// Note: OnceLock doesn't support reset in stable Rust.
@@ -1642,6 +1633,57 @@ mod tests {
assert!(result.is_err());
}
#[test]
fn test_upload_part_customer_key_md5_comparison_is_case_sensitive() {
let mut metadata = HashMap::new();
metadata.insert(
"x-amz-server-side-encryption-customer-key-md5".to_string(),
"AbCdEfGhIjKlMnOpQrStUvWxYz0123456789+/==".to_string(),
);
let request = EncryptionRequest {
bucket: "bucket",
key: "object",
server_side_encryption: None,
ssekms_key_id: None,
sse_customer_algorithm: None,
sse_customer_key: None,
sse_customer_key_md5: None,
content_size: 1,
part_number: Some(1),
part_key: None,
part_nonce: None,
};
let mismatch = "aBcDeFgHiJkLmNoPqRsTuVwXyZ0123456789+/==".to_string();
let result = request.check_upload_part_customer_key_md5(&metadata, Some(mismatch));
assert!(result.is_err());
}
#[test]
fn test_upload_part_customer_key_md5_exact_match() {
let mut metadata = HashMap::new();
let md5 = "AbCdEfGhIjKlMnOpQrStUvWxYz0123456789+/==".to_string();
metadata.insert("x-amz-server-side-encryption-customer-key-md5".to_string(), md5.clone());
let request = EncryptionRequest {
bucket: "bucket",
key: "object",
server_side_encryption: None,
ssekms_key_id: None,
sse_customer_algorithm: None,
sse_customer_key: None,
sse_customer_key_md5: None,
content_size: 1,
part_number: Some(1),
part_key: None,
part_nonce: None,
};
let result = request.check_upload_part_customer_key_md5(&metadata, Some(md5));
assert!(result.is_ok());
}
// ============================================================================
// Integration Tests - Encrypt/Decrypt with SimpleSseDekProvider
// ============================================================================