mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-10 15:16:56 +00:00
Merge branch 'main' into feat/kms-vault-transit2
This commit is contained in:
@@ -24,7 +24,7 @@ use rustfs_audit::{audit_system, start_audit_system as start_global_audit_system
|
||||
use rustfs_config::audit::{AUDIT_MQTT_KEYS, AUDIT_MQTT_SUB_SYS, AUDIT_ROUTE_PREFIX, AUDIT_WEBHOOK_KEYS, AUDIT_WEBHOOK_SUB_SYS};
|
||||
use rustfs_config::{DEFAULT_DELIMITER, ENABLE_KEY, ENV_PREFIX, EnableState, MAX_ADMIN_REQUEST_BODY_SIZE};
|
||||
use rustfs_ecstore::config::Config;
|
||||
use rustfs_targets::check_mqtt_broker_available;
|
||||
use rustfs_targets::{TargetError, check_mqtt_broker_available_with_tls, target::mqtt::MQTTTlsConfig};
|
||||
use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::{HashMap, HashSet};
|
||||
@@ -521,9 +521,24 @@ impl Operation for AuditTargetConfig {
|
||||
.ok_or_else(|| s3_error!(InvalidArgument, "topic is required"))?;
|
||||
let username = kv_map.get(rustfs_config::MQTT_USERNAME).map(String::as_str);
|
||||
let password = kv_map.get(rustfs_config::MQTT_PASSWORD).map(String::as_str);
|
||||
check_mqtt_broker_available(endpoint, topic, username, password)
|
||||
let tls = MQTTTlsConfig::from_values(
|
||||
kv_map.get(rustfs_config::MQTT_TLS_POLICY).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_TLS_CA).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_TLS_CLIENT_CERT).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_TLS_CLIENT_KEY).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_TLS_TRUST_LEAF_AS_CA).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_WS_PATH_ALLOWLIST).map(String::as_str),
|
||||
)
|
||||
.map_err(|e| s3_error!(InvalidArgument, "invalid MQTT TLS settings: {}", e))?;
|
||||
let parsed_broker = Url::parse(endpoint).map_err(|e| s3_error!(InvalidArgument, "invalid broker URL: {}", e))?;
|
||||
rustfs_targets::target::mqtt::validate_mqtt_broker_url(&parsed_broker, &tls)
|
||||
.map_err(|e| s3_error!(InvalidArgument, "{}", e))?;
|
||||
check_mqtt_broker_available_with_tls(parsed_broker.as_str(), topic, username, password, &tls)
|
||||
.await
|
||||
.map_err(|e| s3_error!(InvalidArgument, "MQTT Broker unavailable: {}", e))?;
|
||||
.map_err(|e| match e {
|
||||
TargetError::Configuration(_) => s3_error!(InvalidArgument, "{}", e),
|
||||
_ => s3_error!(InvalidArgument, "MQTT broker check failed: {}", e),
|
||||
})?;
|
||||
|
||||
if let Some(queue_dir) = kv_map.get("queue_dir") {
|
||||
validate_queue_dir(queue_dir.as_str()).await?;
|
||||
|
||||
@@ -25,7 +25,7 @@ use rustfs_config::notify::{
|
||||
};
|
||||
use rustfs_config::{DEFAULT_DELIMITER, ENABLE_KEY, ENV_PREFIX, EnableState, MAX_ADMIN_REQUEST_BODY_SIZE};
|
||||
use rustfs_ecstore::config::Config;
|
||||
use rustfs_targets::check_mqtt_broker_available;
|
||||
use rustfs_targets::{TargetError, check_mqtt_broker_available_with_tls, target::mqtt::MQTTTlsConfig};
|
||||
use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::{HashMap, HashSet};
|
||||
@@ -450,9 +450,24 @@ impl Operation for NotificationTarget {
|
||||
.ok_or_else(|| s3_error!(InvalidArgument, "topic is required"))?;
|
||||
let username = kv_map.get(rustfs_config::MQTT_USERNAME).map(String::as_str);
|
||||
let password = kv_map.get(rustfs_config::MQTT_PASSWORD).map(String::as_str);
|
||||
check_mqtt_broker_available(endpoint, topic, username, password)
|
||||
let tls = MQTTTlsConfig::from_values(
|
||||
kv_map.get(rustfs_config::MQTT_TLS_POLICY).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_TLS_CA).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_TLS_CLIENT_CERT).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_TLS_CLIENT_KEY).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_TLS_TRUST_LEAF_AS_CA).map(String::as_str),
|
||||
kv_map.get(rustfs_config::MQTT_WS_PATH_ALLOWLIST).map(String::as_str),
|
||||
)
|
||||
.map_err(|e| s3_error!(InvalidArgument, "invalid MQTT TLS settings: {}", e))?;
|
||||
let parsed_broker = Url::parse(endpoint).map_err(|e| s3_error!(InvalidArgument, "invalid broker URL: {}", e))?;
|
||||
rustfs_targets::target::mqtt::validate_mqtt_broker_url(&parsed_broker, &tls)
|
||||
.map_err(|e| s3_error!(InvalidArgument, "{}", e))?;
|
||||
check_mqtt_broker_available_with_tls(parsed_broker.as_str(), topic, username, password, &tls)
|
||||
.await
|
||||
.map_err(|e| s3_error!(InvalidArgument, "MQTT Broker unavailable: {}", e))?;
|
||||
.map_err(|e| match e {
|
||||
TargetError::Configuration(_) => s3_error!(InvalidArgument, "{}", e),
|
||||
_ => s3_error!(InvalidArgument, "MQTT broker check failed: {}", e),
|
||||
})?;
|
||||
|
||||
if let Some(queue_dir) = kv_map.get("queue_dir") {
|
||||
validate_queue_dir(queue_dir.as_str()).await?;
|
||||
|
||||
@@ -170,18 +170,18 @@ fn encode_list_objects_v2_value(value: &str, encoding_type: Option<&EncodingType
|
||||
}
|
||||
}
|
||||
|
||||
fn build_metadata_extension_user_metadata(user_defined: &HashMap<String, String>) -> Option<MinioUserMetadata> {
|
||||
fn build_metadata_extension_user_metadata(user_defined: &HashMap<String, String>) -> Option<UserMetadataCollection> {
|
||||
let mut items = extract_user_defined_metadata(user_defined)
|
||||
.into_iter()
|
||||
.filter(|(key, _)| !key.is_empty())
|
||||
.map(|(key, value)| MinioMetadataEntry { key, value })
|
||||
.map(|(key, value)| UserMetadataEntry { key, value })
|
||||
.collect::<Vec<_>>();
|
||||
items.sort_by(|left, right| left.key.cmp(&right.key));
|
||||
|
||||
if items.is_empty() {
|
||||
None
|
||||
} else {
|
||||
Some(MinioUserMetadata { items })
|
||||
Some(UserMetadataCollection { items })
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2729,7 +2729,7 @@ mod tests {
|
||||
assert_eq!(version.internal, Some(ObjectInternalInfo { k: 4, m: 2 }));
|
||||
assert_eq!(
|
||||
version.user_metadata.as_ref().map(|metadata| metadata.items.clone()),
|
||||
Some(vec![MinioMetadataEntry {
|
||||
Some(vec![UserMetadataEntry {
|
||||
key: "project".to_string(),
|
||||
value: "alpha".to_string(),
|
||||
}])
|
||||
@@ -2745,7 +2745,7 @@ mod tests {
|
||||
assert!(marker.user_tags.is_none());
|
||||
assert_eq!(
|
||||
marker.user_metadata.as_ref().map(|metadata| metadata.items.clone()),
|
||||
Some(vec![MinioMetadataEntry {
|
||||
Some(vec![UserMetadataEntry {
|
||||
key: "marker".to_string(),
|
||||
value: "true".to_string(),
|
||||
}])
|
||||
@@ -2840,7 +2840,7 @@ mod tests {
|
||||
assert!(object.owner.is_some());
|
||||
assert_eq!(
|
||||
object.user_metadata.as_ref().map(|metadata| metadata.items.clone()),
|
||||
Some(vec![MinioMetadataEntry {
|
||||
Some(vec![UserMetadataEntry {
|
||||
key: "project".to_string(),
|
||||
value: "alpha".to_string(),
|
||||
}])
|
||||
|
||||
@@ -23,7 +23,7 @@ use crate::storage::entity;
|
||||
use crate::storage::helper::OperationHelper;
|
||||
use crate::storage::options::{
|
||||
copy_src_opts, extract_metadata, get_complete_multipart_upload_opts, get_content_sha256_with_query, get_opts,
|
||||
parse_copy_source_range, put_opts,
|
||||
parse_copy_source_range, put_opts, validate_archive_content_encoding,
|
||||
};
|
||||
use crate::storage::request_context::spawn_traced;
|
||||
use crate::storage::s3_api::multipart::build_list_parts_output;
|
||||
@@ -559,6 +559,12 @@ impl DefaultMultipartUsecase {
|
||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||
};
|
||||
|
||||
validate_archive_content_encoding(
|
||||
&key,
|
||||
req.headers.get("content-type").and_then(|value| value.to_str().ok()),
|
||||
req.headers.get("content-encoding").and_then(|value| value.to_str().ok()),
|
||||
)?;
|
||||
|
||||
let mut metadata = extract_metadata(&req.headers);
|
||||
|
||||
if let Some(tags) = tagging {
|
||||
|
||||
@@ -38,6 +38,7 @@ use crate::storage::helper::OperationHelper;
|
||||
use crate::storage::options::{
|
||||
copy_dst_opts, copy_src_opts, del_opts, extract_metadata, extract_metadata_from_mime_with_object_name,
|
||||
filter_object_metadata, get_content_sha256_with_query, get_opts, normalize_content_encoding_for_storage, put_opts,
|
||||
validate_archive_content_encoding,
|
||||
};
|
||||
use crate::storage::s3_api::multipart::parse_list_parts_params;
|
||||
use crate::storage::s3_api::{acl, restore, select};
|
||||
@@ -315,6 +316,20 @@ fn apply_put_request_metadata(
|
||||
tagging: Option<TaggingHeader>,
|
||||
storage_class: Option<StorageClass>,
|
||||
) -> S3Result<()> {
|
||||
let request_content_type = content_type.as_ref().map(ToString::to_string).or_else(|| {
|
||||
headers
|
||||
.get("content-type")
|
||||
.and_then(|value| value.to_str().ok())
|
||||
.map(ToOwned::to_owned)
|
||||
});
|
||||
let request_content_encoding = content_encoding.as_ref().map(ToString::to_string).or_else(|| {
|
||||
headers
|
||||
.get("content-encoding")
|
||||
.and_then(|value| value.to_str().ok())
|
||||
.map(ToOwned::to_owned)
|
||||
});
|
||||
validate_archive_content_encoding(object_name, request_content_type.as_deref(), request_content_encoding.as_deref())?;
|
||||
|
||||
if let Some(cache_control) = cache_control {
|
||||
metadata.insert("cache-control".to_string(), cache_control.to_string());
|
||||
}
|
||||
|
||||
@@ -104,6 +104,10 @@ impl ObjectIoCachedGetObjectSource for CachedGetObject {
|
||||
self.last_modified.as_deref()
|
||||
}
|
||||
|
||||
fn expires(&self) -> Option<&str> {
|
||||
self.expires.as_deref()
|
||||
}
|
||||
|
||||
fn cache_control(&self) -> Option<&str> {
|
||||
self.cache_control.as_deref()
|
||||
}
|
||||
|
||||
@@ -46,6 +46,8 @@ use rustfs_config::{
|
||||
DEFAULT_COMPRESS_ENABLE, DEFAULT_COMPRESS_EXTENSIONS, DEFAULT_COMPRESS_MIME_TYPES, DEFAULT_COMPRESS_MIN_SIZE,
|
||||
ENV_COMPRESS_ENABLE, ENV_COMPRESS_EXTENSIONS, ENV_COMPRESS_MIME_TYPES, ENV_COMPRESS_MIN_SIZE, EnableState,
|
||||
};
|
||||
use rustfs_ecstore::compress::{STANDARD_EXCLUDE_COMPRESS_CONTENT_TYPES, STANDARD_EXCLUDE_COMPRESS_EXTENSIONS};
|
||||
use rustfs_utils::string::{has_pattern, has_string_suffix_in_slice};
|
||||
use std::str::FromStr;
|
||||
use tower_http::compression::predicate::Predicate;
|
||||
use tracing::debug;
|
||||
@@ -196,6 +198,20 @@ impl CompressionConfig {
|
||||
|
||||
None
|
||||
}
|
||||
|
||||
pub(crate) fn is_excluded_filename(filename: &str) -> bool {
|
||||
has_string_suffix_in_slice(&filename.to_ascii_lowercase(), STANDARD_EXCLUDE_COMPRESS_EXTENSIONS)
|
||||
}
|
||||
|
||||
pub(crate) fn is_excluded_mime_type(content_type: &str) -> bool {
|
||||
let main_type = content_type
|
||||
.split(';')
|
||||
.next()
|
||||
.unwrap_or(content_type)
|
||||
.trim()
|
||||
.to_ascii_lowercase();
|
||||
!main_type.is_empty() && has_pattern(STANDARD_EXCLUDE_COMPRESS_CONTENT_TYPES, &main_type)
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for CompressionConfig {
|
||||
@@ -299,23 +315,37 @@ impl Predicate for CompressionPredicate {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Check if the response matches configured extension via Content-Disposition
|
||||
// Hard-stop archive/media/package MIME types even if the whitelist matches.
|
||||
// This includes tar, gzip, bzip2, xz, zstd, zip, rar, 7z, lzip, lzma, lzop variants,
|
||||
// plus video/*, audio/*, image/*, font/*, application/pdf, and application/wasm.
|
||||
if let Some(content_type) = response.headers().get(http::header::CONTENT_TYPE)
|
||||
&& let Ok(ct) = content_type.to_str()
|
||||
{
|
||||
if CompressionConfig::is_excluded_mime_type(ct) {
|
||||
debug!("Skipping compression for excluded Content-Type '{}'", ct);
|
||||
return false;
|
||||
}
|
||||
|
||||
if self.config.matches_mime_type(ct) {
|
||||
debug!("Compressing response: Content-Type '{}' matches configured MIME pattern", ct);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
// Hard-stop archive-like attachment downloads even if the whitelist matches.
|
||||
if let Some(content_disposition) = response.headers().get(http::header::CONTENT_DISPOSITION)
|
||||
&& let Ok(cd) = content_disposition.to_str()
|
||||
&& let Some(filename) = CompressionConfig::extract_filename_from_content_disposition(cd)
|
||||
&& self.config.matches_extension(&filename)
|
||||
{
|
||||
debug!("Compressing response: filename '{}' matches configured extension", filename);
|
||||
return true;
|
||||
}
|
||||
if CompressionConfig::is_excluded_filename(&filename) {
|
||||
debug!("Skipping compression for excluded filename '{}'", filename);
|
||||
return false;
|
||||
}
|
||||
|
||||
// Check if the response matches configured MIME type
|
||||
if let Some(content_type) = response.headers().get(http::header::CONTENT_TYPE)
|
||||
&& let Ok(ct) = content_type.to_str()
|
||||
&& self.config.matches_mime_type(ct)
|
||||
{
|
||||
debug!("Compressing response: Content-Type '{}' matches configured MIME pattern", ct);
|
||||
return true;
|
||||
if self.config.matches_extension(&filename) {
|
||||
debug!("Compressing response: filename '{}' matches configured extension", filename);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
// Default: don't compress (whitelist approach)
|
||||
@@ -405,7 +435,11 @@ use std::task::{Context, Poll};
|
||||
use tower::{Layer, Service};
|
||||
|
||||
/// Tower layer that injects `RequestPathCategory` into each response's extensions
|
||||
/// based on the incoming request URI path. Must be placed before `CompressionLayer`.
|
||||
/// based on the incoming request URI path.
|
||||
///
|
||||
/// It must be placed inside `CompressionLayer` so the category is available when
|
||||
/// the outer compression middleware evaluates its response predicate. With
|
||||
/// `tower::ServiceBuilder`, that means adding this layer after `CompressionLayer`.
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub(crate) struct PathCategoryInjectionLayer;
|
||||
|
||||
@@ -627,6 +661,42 @@ mod tests {
|
||||
assert_eq!(predicate.config.min_size, 1000);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_compression_predicate_skips_archive_mime_type_even_when_whitelisted() {
|
||||
let predicate = CompressionPredicate::new(CompressionConfig {
|
||||
enabled: true,
|
||||
extensions: vec![],
|
||||
mime_patterns: vec!["application/zip".to_string()],
|
||||
min_size: 0,
|
||||
});
|
||||
|
||||
let response = Response::builder()
|
||||
.header(http::header::CONTENT_TYPE, "application/zip")
|
||||
.header(http::header::CONTENT_LENGTH, "4096")
|
||||
.body(http_body_util::Empty::<bytes::Bytes>::new())
|
||||
.expect("response");
|
||||
|
||||
assert!(!predicate.should_compress(&response));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_compression_predicate_skips_archive_filename_even_when_whitelisted() {
|
||||
let predicate = CompressionPredicate::new(CompressionConfig {
|
||||
enabled: true,
|
||||
extensions: vec![".zip".to_string()],
|
||||
mime_patterns: vec![],
|
||||
min_size: 0,
|
||||
});
|
||||
|
||||
let response = Response::builder()
|
||||
.header(http::header::CONTENT_DISPOSITION, r#"attachment; filename="bundle.zip""#)
|
||||
.header(http::header::CONTENT_LENGTH, "4096")
|
||||
.body(http_body_util::Empty::<bytes::Bytes>::new())
|
||||
.expect("response");
|
||||
|
||||
assert!(!predicate.should_compress(&response));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_path_category_classify_s3() {
|
||||
assert_eq!(PathCategory::classify("/"), PathCategory::S3DataPlane);
|
||||
|
||||
@@ -603,8 +603,8 @@ fn process_connection(
|
||||
// 9. KeystoneAuthLayer — X-Auth-Token validation
|
||||
// 10. TraceLayer — request/response tracing + metrics
|
||||
// 11. PropagateRequestIdLayer — X-Request-ID → response
|
||||
// 12. PathCategoryInjectionLayer — injects path category for compression
|
||||
// 13. CompressionLayer — response compression (whitelist, path-aware)
|
||||
// 12. CompressionLayer — response compression (whitelist, path-aware)
|
||||
// 13. PathCategoryInjectionLayer — injects path category for compression predicate
|
||||
// 14. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes
|
||||
// 15. ConditionalCorsLayer — S3 API CORS
|
||||
// 16. RedirectLayer — console redirect (conditional)
|
||||
@@ -739,8 +739,8 @@ fn process_connection(
|
||||
.layer(PropagateRequestIdLayer::x_request_id())
|
||||
// Compress responses based on whitelist configuration
|
||||
// Only compresses when enabled and matches configured extensions/MIME types
|
||||
.layer(PathCategoryInjectionLayer)
|
||||
.layer(CompressionLayer::new().compress_when(PathAwareCompressionPredicate::new(compression_config)))
|
||||
.layer(PathCategoryInjectionLayer)
|
||||
.layer(ObjectAttributesEtagFixLayer)
|
||||
// Conditional CORS layer: only applies to S3 API requests (not Admin, not Console)
|
||||
// Admin has its own CORS handling in router.rs
|
||||
@@ -927,8 +927,16 @@ fn get_default_tcp_keepalive() -> TcpKeepalive {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::server::compress::RequestPathCategory;
|
||||
use bytes::Bytes;
|
||||
use http::HeaderMap;
|
||||
use http::Request as HttpRequest;
|
||||
use http_body_util::Empty;
|
||||
use opentelemetry::propagation::Extractor;
|
||||
use std::convert::Infallible;
|
||||
use std::future::Ready;
|
||||
use std::task::{Context, Poll};
|
||||
use tower::{Layer, Service, ServiceBuilder};
|
||||
|
||||
/// Baseline constants — reference the authoritative config defaults.
|
||||
/// If a config default changes, tests automatically follow.
|
||||
@@ -1065,4 +1073,89 @@ mod tests {
|
||||
assert_eq!(carrier.get("Content-Type"), Some("application/json"));
|
||||
assert_eq!(carrier.get("CONTENT-TYPE"), Some("application/json"));
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
struct ObserveCategoryLayer;
|
||||
|
||||
#[derive(Clone)]
|
||||
struct ObserveCategoryService<S> {
|
||||
inner: S,
|
||||
}
|
||||
|
||||
impl<S> Layer<S> for ObserveCategoryLayer {
|
||||
type Service = ObserveCategoryService<S>;
|
||||
|
||||
fn layer(&self, inner: S) -> Self::Service {
|
||||
ObserveCategoryService { inner }
|
||||
}
|
||||
}
|
||||
|
||||
impl<S, ReqBody, ResBody> Service<HttpRequest<ReqBody>> for ObserveCategoryService<S>
|
||||
where
|
||||
S: Service<HttpRequest<ReqBody>, Response = Response<ResBody>, Error = Infallible>,
|
||||
{
|
||||
type Response = Response<ResBody>;
|
||||
type Error = Infallible;
|
||||
type Future = Ready<std::result::Result<Response<ResBody>, Infallible>>;
|
||||
|
||||
fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<std::result::Result<(), Self::Error>> {
|
||||
Poll::Ready(Ok(()))
|
||||
}
|
||||
|
||||
fn call(&mut self, req: HttpRequest<ReqBody>) -> Self::Future {
|
||||
let response = futures::executor::block_on(self.inner.call(req)).expect("infallible");
|
||||
let mut response = response;
|
||||
let seen = response.extensions().get::<RequestPathCategory>().is_some();
|
||||
response
|
||||
.headers_mut()
|
||||
.insert("x-category-seen", if seen { "true" } else { "false" }.parse().expect("header"));
|
||||
std::future::ready(Ok(response))
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
struct OkService;
|
||||
|
||||
impl<ReqBody> Service<HttpRequest<ReqBody>> for OkService {
|
||||
type Response = Response<Empty<Bytes>>;
|
||||
type Error = Infallible;
|
||||
type Future = Ready<std::result::Result<Response<Empty<Bytes>>, Infallible>>;
|
||||
|
||||
fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<std::result::Result<(), Self::Error>> {
|
||||
Poll::Ready(Ok(()))
|
||||
}
|
||||
|
||||
fn call(&mut self, _req: HttpRequest<ReqBody>) -> Self::Future {
|
||||
std::future::ready(Ok(Response::new(Empty::new())))
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_service_builder_order_regression_for_response_extensions() {
|
||||
let request = HttpRequest::builder().uri("/bucket/archive.zip").body(()).expect("request");
|
||||
|
||||
let mut broken_order = ServiceBuilder::new()
|
||||
.layer(PathCategoryInjectionLayer)
|
||||
.layer(ObserveCategoryLayer)
|
||||
.service(OkService);
|
||||
|
||||
let broken_response = futures::executor::block_on(broken_order.call(request)).expect("response");
|
||||
assert_eq!(
|
||||
broken_response.headers().get("x-category-seen").and_then(|v| v.to_str().ok()),
|
||||
Some("false")
|
||||
);
|
||||
|
||||
let request = HttpRequest::builder().uri("/bucket/archive.zip").body(()).expect("request");
|
||||
|
||||
let mut fixed_order = ServiceBuilder::new()
|
||||
.layer(ObserveCategoryLayer)
|
||||
.layer(PathCategoryInjectionLayer)
|
||||
.service(OkService);
|
||||
|
||||
let fixed_response = futures::executor::block_on(fixed_order.call(request)).expect("response");
|
||||
assert_eq!(
|
||||
fixed_response.headers().get("x-category-seen").and_then(|v| v.to_str().ok()),
|
||||
Some("true")
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -33,7 +33,7 @@ use rustfs_policy::service_type::ServiceType;
|
||||
use rustfs_utils::hash::EMPTY_STRING_SHA256_HASH;
|
||||
use rustfs_utils::http::AMZ_CONTENT_SHA256;
|
||||
use rustfs_utils::path::is_dir_object;
|
||||
use s3s::{S3Result, s3_error};
|
||||
use s3s::{S3Error, S3ErrorCode, S3Result, s3_error};
|
||||
use std::collections::HashMap;
|
||||
use std::sync::LazyLock;
|
||||
use tracing::error;
|
||||
@@ -349,6 +349,73 @@ pub(crate) fn normalize_content_encoding_for_storage(value: &str) -> Option<Stri
|
||||
if normalized.is_empty() { None } else { Some(normalized) }
|
||||
}
|
||||
|
||||
const ENV_ALLOW_ARCHIVE_CONTENT_ENCODING: &str = "RUSTFS_ALLOW_ARCHIVE_CONTENT_ENCODING";
|
||||
|
||||
const ARCHIVE_CONTENT_ENCODING_BLOCKED_SUFFIXES: &[&str] = &[
|
||||
".zip",
|
||||
".tar",
|
||||
".tar.gz",
|
||||
".tgz",
|
||||
".tar.bz2",
|
||||
".tbz",
|
||||
".tbz2",
|
||||
".tar.xz",
|
||||
".txz",
|
||||
".tar.zst",
|
||||
".tar.zstd",
|
||||
".tzst",
|
||||
];
|
||||
|
||||
const ARCHIVE_CONTENT_ENCODING_BLOCKED_CONTENT_TYPES: &[&str] =
|
||||
&["application/zip", "application/x-zip-compressed", "application/x-tar"];
|
||||
|
||||
fn is_archive_object_name_for_content_encoding(object_name: &str) -> bool {
|
||||
let object_name = object_name.to_ascii_lowercase();
|
||||
ARCHIVE_CONTENT_ENCODING_BLOCKED_SUFFIXES
|
||||
.iter()
|
||||
.any(|suffix| object_name.ends_with(suffix))
|
||||
}
|
||||
|
||||
fn is_archive_content_type_for_content_encoding(content_type: &str) -> bool {
|
||||
let main_type = content_type
|
||||
.split(';')
|
||||
.next()
|
||||
.unwrap_or(content_type)
|
||||
.trim()
|
||||
.to_ascii_lowercase();
|
||||
|
||||
ARCHIVE_CONTENT_ENCODING_BLOCKED_CONTENT_TYPES
|
||||
.iter()
|
||||
.any(|candidate| main_type == *candidate)
|
||||
}
|
||||
|
||||
pub(crate) fn validate_archive_content_encoding(
|
||||
object_name: &str,
|
||||
content_type: Option<&str>,
|
||||
content_encoding: Option<&str>,
|
||||
) -> S3Result<()> {
|
||||
if rustfs_utils::get_env_bool(ENV_ALLOW_ARCHIVE_CONTENT_ENCODING, false) {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let Some(content_encoding) = content_encoding.map(str::trim).filter(|value| !value.is_empty()) else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
let is_archive_like = is_archive_object_name_for_content_encoding(object_name)
|
||||
|| content_type.is_some_and(is_archive_content_type_for_content_encoding);
|
||||
if !is_archive_like {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
Err(S3Error::with_message(
|
||||
S3ErrorCode::InvalidArgument,
|
||||
format!(
|
||||
"Content-Encoding '{content_encoding}' is not allowed for archive objects by default; set {ENV_ALLOW_ARCHIVE_CONTENT_ENCODING}=true to allow legacy behavior"
|
||||
),
|
||||
))
|
||||
}
|
||||
|
||||
/// Extracts metadata from headers and returns it as a HashMap with object name for MIME type detection.
|
||||
pub fn extract_metadata_from_mime_with_object_name(
|
||||
headers: &HeaderMap<HeaderValue>,
|
||||
@@ -692,6 +759,8 @@ fn get_content_sha256_cksum(headers: &HeaderMap<HeaderValue>, service_type: Serv
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use temp_env;
|
||||
|
||||
use super::*;
|
||||
use http::{HeaderMap, HeaderValue};
|
||||
use std::collections::HashMap;
|
||||
@@ -1358,6 +1427,30 @@ mod tests {
|
||||
assert_eq!(detect_content_type_from_object_name("noextension"), "application/octet-stream");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_validate_archive_content_encoding_rejects_archive_suffix_by_default() {
|
||||
let err = validate_archive_content_encoding("bundle.tar.gz", Some("application/gzip"), Some("gzip")).unwrap_err();
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_validate_archive_content_encoding_rejects_archive_mime_by_default() {
|
||||
let err = validate_archive_content_encoding("bundle", Some("application/zip"), Some("gzip")).unwrap_err();
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_validate_archive_content_encoding_allows_non_archive_precompressed_object() {
|
||||
validate_archive_content_encoding("logs/app.log.zst", Some("text/plain"), Some("zstd")).expect("non-archive");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_validate_archive_content_encoding_allows_legacy_opt_in() {
|
||||
temp_env::with_var(ENV_ALLOW_ARCHIVE_CONTENT_ENCODING, Some("true"), || {
|
||||
validate_archive_content_encoding("bundle.zip", Some("application/zip"), Some("gzip")).expect("legacy opt-in");
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_parse_copy_source_range() {
|
||||
// Test complete range: bytes=0-1023
|
||||
|
||||
Reference in New Issue
Block a user