// 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. //! Bucket application use-case contracts. use super::storage_api::bucket_usecase::ECStore; use super::storage_api::bucket_usecase::StorageObjectInfo as ObjectInfo; #[cfg(test)] use super::storage_api::bucket_usecase::access::ReqInfo; use super::storage_api::bucket_usecase::access::{ authorize_request, bucket_config_mutation_incarnation, log_list_buckets_iam_implicit_deny, prepare_list_buckets_iam_authorization, req_info_ref, }; #[cfg(test)] use super::storage_api::bucket_usecase::bucket::target::BucketTarget; use super::storage_api::bucket_usecase::bucket::{ ObjectLockConfigExt as _, VersioningConfigExt as _, lifecycle::bucket_lifecycle_ops::{ enqueue_expiry_for_existing_objects, enqueue_transition_for_existing_objects, run_stale_multipart_upload_cleanup_once, validate_lifecycle_config, validate_transition_tier, }, metadata::{ BUCKET_CORS_CONFIG, BUCKET_LIFECYCLE_CONFIG, BUCKET_NOTIFICATION_CONFIG, BUCKET_POLICY_CONFIG, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, BUCKET_REPLICATION_CONFIG, BUCKET_SSECONFIG, BUCKET_TAGGING_CONFIG, BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, }, metadata_sys, policy_sys::PolicySys, replication::{ ReplicationTargetValidationError, invalid_replication_config_status_field, replication_target_arns, should_remove_replication_target, unsupported_replication_config_field, validate_replication_config_structure, validate_replication_config_target_arns, }, target::{BucketTargetType, BucketTargets}, utils::serialize, versioning_sys::BucketVersioningSys, }; use super::storage_api::bucket_usecase::contract::bucket::{ BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, }; use super::storage_api::bucket_usecase::contract::list::{ ListObjectVersionsInfo as StorageListObjectVersionsInfo, ListObjectsV2Info as StorageListObjectsV2Info, ListOperations as _, }; use super::storage_api::bucket_usecase::error::StorageError; use super::storage_api::bucket_usecase::helper::{OperationHelper, spawn_background_with_context}; use super::storage_api::bucket_usecase::object_utils::to_s3s_etag; use super::storage_api::bucket_usecase::s3_api::bucket::{ ListObjectVersionsParams, ListObjectsV2Params, build_list_buckets_output, build_list_object_versions_output, build_list_objects_output, build_list_objects_v2_output, parse_list_object_versions_params, parse_list_objects_v2_params, rustfs_owner, }; use super::storage_api::bucket_usecase::{ get_validated_store, process_lambda_configurations, process_queue_configurations, process_topic_configurations, request_context, validate_list_object_unordered_with_delimiter, }; use crate::admin::handlers::site_replication::{ site_replication_bucket_meta_hook, site_replication_delete_bucket_hook, site_replication_make_bucket_hook, }; use crate::app::object_data_cache::invalidate_object_data_cache_bucket_after_delete; use crate::app::runtime_sources::{ AppContext, current_app_context, current_encryption_service, current_notification_system, current_notify_interface_for_context, current_object_data_cache_for_context, current_object_store_handle_for_context, }; use crate::auth::get_condition_values_with_client_info; use crate::error::ApiError; use crate::shared_types::RemoteAddr; use crate::storage::storage_api::lock_bucket_targets_metadata; use http::StatusCode; use metrics::counter; use rustfs_config::RUSTFS_REGION; use rustfs_io_metrics::record_s3_op; use rustfs_madmin::{SITE_REPL_API_VERSION, SRBucketMeta}; use rustfs_policy::policy::{ action::{Action, S3Action}, {BucketPolicy, BucketPolicyArgs, Effect, Validator}, }; use rustfs_s3_ops::S3Operation; use rustfs_targets::{ EventName, arn::{ARN, TargetIDError}, }; use rustfs_trusted_proxies::ClientInfo; use rustfs_utils::http::{SUFFIX_FORCE_DELETE, get_header}; use rustfs_utils::obj::extract_user_defined_metadata; use rustfs_utils::string::parse_bool; use s3s::dto::{ BucketLifecycleConfiguration, BucketLocationConstraint, BucketVersioningStatus, CommonPrefix, CreateBucketInput, CreateBucketOutput, DeleteBucketCorsInput, DeleteBucketCorsOutput, DeleteBucketEncryptionInput, DeleteBucketEncryptionOutput, DeleteBucketInput, DeleteBucketLifecycleInput, DeleteBucketLifecycleOutput, DeleteBucketOutput, DeleteBucketPolicyInput, DeleteBucketPolicyOutput, DeleteBucketReplicationInput, DeleteBucketReplicationOutput, DeleteBucketTaggingInput, DeleteBucketTaggingOutput, DeleteMarkerEntry, DeletePublicAccessBlockInput, DeletePublicAccessBlockOutput, EncodingType, ExpirationStatus, GetBucketCorsInput, GetBucketCorsOutput, GetBucketEncryptionInput, GetBucketEncryptionOutput, GetBucketLifecycleConfigurationInput, GetBucketLifecycleConfigurationOutput, GetBucketLocationInput, GetBucketLocationOutput, GetBucketNotificationConfigurationInput, GetBucketNotificationConfigurationOutput, GetBucketPolicyInput, GetBucketPolicyOutput, GetBucketPolicyStatusInput, GetBucketPolicyStatusOutput, GetBucketReplicationInput, GetBucketReplicationOutput, GetBucketTaggingInput, GetBucketTaggingOutput, GetBucketVersioningInput, GetBucketVersioningOutput, GetPublicAccessBlockInput, GetPublicAccessBlockOutput, HeadBucketInput, HeadBucketOutput, LifecycleRule, ListBucketsInput, ListBucketsOutput, ListObjectVersionsInput, ListObjectVersionsOutput, ListObjectsInput, ListObjectsOutput, ListObjectsV2Input, ListObjectsV2Output, MetadataEntry, NotificationConfiguration, NotificationConfigurationFilter, Object, ObjectLockConfiguration, ObjectStorageClass, ObjectVersion, ObjectVersionStorageClass, PolicyStatus, PutBucketCorsInput, PutBucketCorsOutput, PutBucketEncryptionInput, PutBucketEncryptionOutput, PutBucketLifecycleConfigurationInput, PutBucketLifecycleConfigurationOutput, PutBucketNotificationConfigurationInput, PutBucketNotificationConfigurationOutput, PutBucketPolicyInput, PutBucketPolicyOutput, PutBucketReplicationInput, PutBucketReplicationOutput, PutBucketTaggingInput, PutBucketTaggingOutput, PutBucketVersioningInput, PutBucketVersioningOutput, PutPublicAccessBlockInput, PutPublicAccessBlockOutput, ReplicationConfiguration, ServerSideEncryption, Tagging, Timestamp, UserMetadata, VersioningConfiguration, }; use s3s::region::Region; use s3s::xml; use s3s::xml::{SerResult, SerializeContent}; use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; use std::{ borrow::Cow, collections::{HashMap, HashSet}, fmt::Display, future::Future, io::Write, sync::{Arc, LazyLock}, }; use tokio::sync::{Semaphore, TryAcquireError}; use tracing::{Instrument as _, debug, error, info, instrument, warn}; const LOG_COMPONENT_APP: &str = "app"; const LOG_SUBSYSTEM_BUCKET: &str = "bucket"; const BUCKET_OPERATION_CONCURRENCY: usize = 8; static BUCKET_OPERATION_ADMISSION: LazyLock> = LazyLock::new(|| Arc::new(Semaphore::new(BUCKET_OPERATION_CONCURRENCY))); use urlencoding::encode; type ListObjectVersionsInfo = StorageListObjectVersionsInfo; type ListObjectsV2Info = StorageListObjectsV2Info; const XMLNS_S3: &str = "http://s3.amazonaws.com/doc/2006-03-01/"; #[derive(Clone, Debug, Default, PartialEq)] pub(crate) struct ObjectInternalInfo { pub k: i32, pub m: i32, } #[derive(Clone, Debug, Default, PartialEq)] pub(crate) struct ObjectMetadataExtension { pub user_metadata: Option, pub user_tags: Option, pub internal: Option, } #[derive(Clone, Debug, PartialEq)] pub(crate) enum ListObjectVersionMetadataEntry { DeleteMarker(DeleteMarkerEntry, ObjectMetadataExtension), Version(ObjectVersion, ObjectMetadataExtension), } #[derive(Clone, Debug, Default, PartialEq)] pub(crate) struct ListObjectVersionsMetadataOutput { pub output: ListObjectVersionsOutput, pub entries: Vec, } #[derive(Clone, Debug, Default, PartialEq)] pub(crate) struct ObjectMetadataEntry { pub object: Object, pub extension: ObjectMetadataExtension, } #[derive(Clone, Debug, Default, PartialEq)] pub(crate) struct ListObjectsV2MetadataOutput { pub output: ListObjectsV2Output, pub contents: Vec, } impl xml::Serialize for ListObjectVersionsMetadataOutput { fn serialize(&self, s: &mut xml::Serializer) -> SerResult { s.content_with_ns("ListVersionsResult", XMLNS_S3, self) } } impl SerializeContent for ListObjectVersionsMetadataOutput { fn serialize_content(&self, s: &mut xml::Serializer) -> SerResult { if let Some(iter) = &self.output.common_prefixes { s.flattened_list("CommonPrefixes", iter)?; } if let Some(ref val) = self.output.delimiter { s.content("Delimiter", val)?; } if let Some(ref val) = self.output.encoding_type { s.content("EncodingType", val)?; } if let Some(ref val) = self.output.is_truncated { s.content("IsTruncated", val)?; } if let Some(ref val) = self.output.key_marker { s.content("KeyMarker", val)?; } if let Some(ref val) = self.output.max_keys { s.content("MaxKeys", val)?; } if let Some(ref val) = self.output.name { s.content("Name", val)?; } if let Some(ref val) = self.output.next_key_marker { s.content("NextKeyMarker", val)?; } if let Some(ref val) = self.output.next_version_id_marker { s.content("NextVersionIdMarker", val)?; } if let Some(ref val) = self.output.prefix { s.content("Prefix", val)?; } if let Some(ref val) = self.output.version_id_marker { s.content("VersionIdMarker", val)?; } for entry in &self.entries { match entry { ListObjectVersionMetadataEntry::DeleteMarker(marker, extension) => { s.content( "DeleteMarker", &VersionMetadataContent { value: marker, extension, }, )?; } ListObjectVersionMetadataEntry::Version(version, extension) => { s.content( "Version", &VersionMetadataContent { value: version, extension, }, )?; } } } Ok(()) } } impl xml::Serialize for ListObjectsV2MetadataOutput { fn serialize(&self, s: &mut xml::Serializer) -> SerResult { s.content_with_ns("ListBucketResult", XMLNS_S3, self) } } impl SerializeContent for ListObjectsV2MetadataOutput { fn serialize_content(&self, s: &mut xml::Serializer) -> SerResult { if let Some(ref val) = self.output.name { s.content("Name", val)?; } if let Some(ref val) = self.output.prefix { s.content("Prefix", val)?; } if let Some(ref val) = self.output.max_keys { s.content("MaxKeys", val)?; } if let Some(ref val) = self.output.key_count { s.content("KeyCount", val)?; } if let Some(ref val) = self.output.continuation_token { s.content("ContinuationToken", val)?; } if let Some(ref val) = self.output.is_truncated { s.content("IsTruncated", val)?; } if let Some(ref val) = self.output.next_continuation_token { s.content("NextContinuationToken", val)?; } for entry in &self.contents { s.content("Contents", entry)?; } if let Some(iter) = &self.output.common_prefixes { s.flattened_list("CommonPrefixes", iter)?; } if let Some(ref val) = self.output.delimiter { s.content("Delimiter", val)?; } if let Some(ref val) = self.output.encoding_type { s.content("EncodingType", val)?; } if let Some(ref val) = self.output.start_after { s.content("StartAfter", val)?; } Ok(()) } } struct VersionMetadataContent<'a, T> { value: &'a T, extension: &'a ObjectMetadataExtension, } impl SerializeContent for VersionMetadataContent<'_, T> { fn serialize_content(&self, s: &mut xml::Serializer) -> SerResult { self.value.serialize_content(s)?; self.extension.serialize_content(s) } } impl SerializeContent for ObjectMetadataEntry { fn serialize_content(&self, s: &mut xml::Serializer) -> SerResult { self.object.serialize_content(s)?; self.extension.serialize_content(s) } } /// Whether `name` is a legal XML 1.0 element name. /// /// User metadata keys are emitted as XML *element names* in the `metadata=true` /// listing (`value`). Keys derived from HTTP headers can legally /// contain characters that are illegal in an XML `Name` (spaces, `$`, `%`, `#`, /// leading digits, control bytes, ...). Emitting such a key verbatim as a tag /// produces a malformed document that the console's XML parser rejects wholesale, /// blanking the entire prefix in the Web UI while plain `ListObjectsV2` (which /// never serializes user metadata) keeps working. See issue #2743. fn is_valid_xml_name(name: &str) -> bool { let mut chars = name.chars(); let Some(first) = chars.next() else { return false; }; let is_name_start = |c: char| c == '_' || c == ':' || c.is_alphabetic(); let is_name_char = |c: char| is_name_start(c) || c == '-' || c == '.' || c.is_numeric(); is_name_start(first) && chars.all(is_name_char) } /// Drop characters that are illegal in XML 1.0 text content. /// /// XML 1.0 permits tab/newline/carriage-return but forbids the other C0 control /// characters. The serializer escapes `< > & ' "` but does not strip these, so a /// metadata value carrying e.g. a `\u{1}` byte would still emit an unparsable /// document. Returns a borrowed slice when nothing needs stripping. fn sanitize_xml_text(value: &str) -> Cow<'_, str> { let is_illegal = |c: char| (c as u32) < 0x20 && c != '\t' && c != '\n' && c != '\r'; if value.contains(is_illegal) { Cow::Owned(value.chars().filter(|c| !is_illegal(*c)).collect()) } else { Cow::Borrowed(value) } } impl SerializeContent for ObjectMetadataExtension { fn serialize_content(&self, s: &mut xml::Serializer) -> SerResult { if let Some(metadata) = &self.user_metadata { s.element("UserMetadata", |s| { for entry in metadata { if let (Some(name), Some(value)) = (&entry.name, &entry.value) { // Skip entries whose key is not a valid XML element name and strip // XML-illegal control characters from the value, so a single poison // object can never corrupt the whole listing document. See issue #2743. if is_valid_xml_name(name) { let sanitized = sanitize_xml_text(value); s.content(name, sanitized.as_ref())?; } } } Ok(()) })?; } if let Some(tags) = &self.user_tags { s.content("UserTags", tags)?; } if let Some(internal) = &self.internal { s.content("Internal", internal)?; } Ok(()) } } impl SerializeContent for ObjectInternalInfo { fn serialize_content(&self, s: &mut xml::Serializer) -> SerResult { s.content("K", &self.k)?; s.content("M", &self.m) } } fn serialize_config(value: &T) -> S3Result> { serialize(value).map_err(to_internal_error) } async fn update_bucket_config_for_incarnation( bucket: &str, config_file: &str, data: Vec, expected_incarnation_id: Option, ) -> Result { match expected_incarnation_id { Some(incarnation_id) => metadata_sys::update_if_incarnation(bucket, config_file, data, incarnation_id).await, None => metadata_sys::update(bucket, config_file, data).await, } } async fn delete_bucket_config_for_incarnation( bucket: &str, config_file: &str, expected_incarnation_id: Option, ) -> Result { match expected_incarnation_id { Some(incarnation_id) => metadata_sys::delete_if_incarnation(bucket, config_file, incarnation_id).await, None => metadata_sys::delete(bucket, config_file).await, } } fn to_internal_error(err: impl Display) -> S3Error { S3Error::with_message(S3ErrorCode::InternalError, format!("{err}")) } fn is_valid_notification_filter_value(value: &str) -> bool { if value.len() > 1024 || value.contains('\\') { return false; } !value.split('/').any(|segment| segment == "." || segment == "..") } fn invalid_filter_value_message(cfg_scope: &str, value: &str) -> String { format!("invalid notification filter value (len={}) ({cfg_scope})", value.len()) } fn invalid_filter_name_message(cfg_scope: &str, name: &str) -> String { format!( "invalid notification filter name (len={}) (only 'prefix'/'suffix' are supported) ({cfg_scope})", name.len() ) } fn validate_notification_filter_rules( filter: Option<&NotificationConfigurationFilter>, cfg_kind: &str, cfg_id: Option<&str>, ) -> S3Result<()> { let Some(filter) = filter else { return Ok(()); }; let Some(s3key_filter) = filter.key.as_ref() else { return Ok(()); }; let Some(rules) = s3key_filter.filter_rules.as_ref() else { return Ok(()); }; let mut has_prefix = false; let mut has_suffix = false; let cfg_scope = cfg_id.map_or_else(|| cfg_kind.to_string(), |id| format!("{cfg_kind} id={id}")); for rule in rules { let Some(name) = rule.name.as_ref() else { return Err(s3_error!(InvalidArgument, "invalid notification filter rule: missing Name ({cfg_scope})")); }; let Some(value) = rule.value.as_ref() else { return Err(s3_error!( InvalidArgument, "invalid notification filter rule: missing Value ({cfg_scope})" )); }; if !is_valid_notification_filter_value(value) { return Err(s3_error!(InvalidArgument, "{}", invalid_filter_value_message(&cfg_scope, value))); } if name.as_str().eq_ignore_ascii_case("prefix") { if has_prefix { return Err(s3_error!(InvalidArgument, "duplicate notification filter name 'prefix' ({cfg_scope})")); } has_prefix = true; } else if name.as_str().eq_ignore_ascii_case("suffix") { if has_suffix { return Err(s3_error!(InvalidArgument, "duplicate notification filter name 'suffix' ({cfg_scope})")); } has_suffix = true; } else { return Err(s3_error!(InvalidArgument, "{}", invalid_filter_name_message(&cfg_scope, name.as_str()))); } } Ok(()) } fn validate_notification_configuration_filters(notification_configuration: &NotificationConfiguration) -> S3Result<()> { if let Some(queue_configs) = notification_configuration.queue_configurations.as_ref() { for cfg in queue_configs { validate_notification_filter_rules(cfg.filter.as_ref(), "QueueConfiguration", cfg.id.as_deref())?; } } if let Some(topic_configs) = notification_configuration.topic_configurations.as_ref() { for cfg in topic_configs { validate_notification_filter_rules(cfg.filter.as_ref(), "TopicConfiguration", cfg.id.as_deref())?; } } if let Some(lambda_configs) = notification_configuration.lambda_function_configurations.as_ref() { for cfg in lambda_configs { validate_notification_filter_rules(cfg.filter.as_ref(), "LambdaFunctionConfiguration", cfg.id.as_deref())?; } } Ok(()) } fn sr_bucket_meta_item(bucket: String, item_type: &str) -> SRBucketMeta { SRBucketMeta { bucket, r#type: item_type.to_string(), updated_at: Some(time::OffsetDateTime::now_utc()), api_version: Some(SITE_REPL_API_VERSION.to_string()), ..Default::default() } } fn notify_bucket_metadata_reload( bucket: String, operation: &'static str, request_context: Option, scanner_maintenance_change: bool, ) { record_local_scanner_maintenance_reload(&bucket, scanner_maintenance_change); spawn_background_with_context(request_context, async move { if let Some(notification_sys) = current_notification_system() { let result = if scanner_maintenance_change { notification_sys.load_bucket_metadata_for_scanner_maintenance(&bucket).await } else { notification_sys.load_bucket_metadata(&bucket).await }; if let Err(err) = result { warn!(bucket = %bucket, error = %err, "failed to notify peers after {operation}"); } } }); } fn record_local_scanner_maintenance_reload(bucket: &str, scanner_maintenance_change: bool) { if scanner_maintenance_change { rustfs_scanner::record_scanner_maintenance_change(bucket); } } /// Notify peers to drop their cached metadata for a bucket that was just deleted, /// so they stop serving stale bucket configuration. Runs in the background to /// avoid blocking the delete response. fn notify_bucket_metadata_delete(bucket: String, request_context: Option) { spawn_background_with_context(request_context, async move { if let Some(notification_sys) = current_notification_system() { for peer_err in notification_sys .delete_bucket_metadata(&bucket) .await .into_iter() .filter(|e| e.err.is_some()) { warn!( bucket = %bucket, host = %peer_err.host, error = ?peer_err.err, "failed to notify peer to delete bucket metadata" ); } } }); } fn validate_replication_config_targets(targets: &BucketTargets, config: &ReplicationConfiguration) -> S3Result<()> { let configured_arns = targets .targets .iter() .filter(|target| target.target_type == BucketTargetType::ReplicationService) .map(|target| target.arn.as_str()); match validate_replication_config_target_arns(configured_arns, config) { Ok(()) => Ok(()), Err(err) => { let message = match err { ReplicationTargetValidationError::RoleWithMultipleDestinations => { "replication config with Role cannot define multiple destination targets" } ReplicationTargetValidationError::StaleTarget => "replication config has a stale target", }; Err(S3Error::with_message(S3ErrorCode::InvalidRequest, message)) } } } fn validate_replication_config_capabilities(config: &ReplicationConfiguration) -> S3Result<()> { if let Err(err) = validate_replication_config_structure(config) { return Err(S3Error::with_message(S3ErrorCode::InvalidRequest, err.message())); } if let Some(field) = invalid_replication_config_status_field(config) { return Err(S3Error::with_message( S3ErrorCode::InvalidRequest, format!("replication field {field} has an invalid status"), )); } if let Some(field) = unsupported_replication_config_field(config) { return Err(S3Error::with_message( S3ErrorCode::InvalidRequest, format!("replication field {field} is not supported by this RustFS version"), )); } Ok(()) } async fn validate_bucket_replication_update(bucket: &str, config: &ReplicationConfiguration) -> S3Result<()> { if !BucketVersioningSys::enabled(bucket).await { return Err(s3_error!( InvalidRequest, "bucket versioning must be enabled before replication can be configured" )); } let targets = metadata_sys::get_bucket_targets_config(bucket) .await .map_err(|err| match err { StorageError::ConfigNotFound => { S3Error::with_message(S3ErrorCode::InvalidRequest, "replication target configuration not found".to_string()) } other => ApiError::from(other).into(), })?; validate_replication_config_targets(&targets, config) } async fn replication_targets_without_config_targets( bucket: &str, config: &ReplicationConfiguration, ) -> S3Result> { let target_arns = replication_target_arns(config); if target_arns.is_empty() { return Ok(None); } let mut targets = match metadata_sys::get_bucket_targets_config(bucket).await { Ok(targets) => targets, Err(StorageError::ConfigNotFound) => return Ok(None), Err(err) => return Err(ApiError::from(err).into()), }; let removed = remove_replication_targets_from_config_targets(&mut targets, &target_arns); if removed == 0 { return Ok(None); } Ok(Some((targets, removed))) } fn remove_replication_targets_from_config_targets(targets: &mut BucketTargets, target_arns: &HashSet) -> usize { let original_len = targets.targets.len(); targets.targets.retain(|target| { !should_remove_replication_target( target.arn.as_str(), target.target_type == BucketTargetType::ReplicationService, target_arns, ) }); original_len - targets.targets.len() } async fn write_replication_targets_after_config_delete( bucket: &str, targets: &BucketTargets, removed: usize, expected_incarnation_id: Option, ) -> S3Result<()> { let json_targets = serde_json::to_vec(&targets).map_err(to_internal_error)?; update_bucket_config_for_incarnation(bucket, BUCKET_TARGETS_FILE, json_targets, expected_incarnation_id) .await .map_err(ApiError::from)?; info!(bucket = %bucket, removed, "removed replication remote targets referenced by deleted bucket replication config"); Ok(()) } async fn restore_replication_config_after_target_cleanup_failure( bucket: &str, config: &ReplicationConfiguration, cleanup_err: S3Error, expected_incarnation_id: Option, ) -> S3Error { match serialize(config) { Ok(data) => { if let Err(restore_err) = update_bucket_config_for_incarnation(bucket, BUCKET_REPLICATION_CONFIG, data, expected_incarnation_id).await { error!( bucket = %bucket, error = ?restore_err, cleanup_error = ?cleanup_err, "failed to restore bucket replication config after target cleanup failure" ); } } Err(restore_err) => { error!( bucket = %bucket, error = ?restore_err, cleanup_error = ?cleanup_err, "failed to serialize bucket replication config for restore after target cleanup failure" ); } } cleanup_err } fn versioning_configuration_has_object_lock_incompatible_settings(config: &VersioningConfiguration) -> bool { config.suspended() || config.exclude_folders.unwrap_or(false) || config .excluded_prefixes .as_ref() .is_some_and(|excluded_prefixes| !excluded_prefixes.is_empty()) } async fn validate_bucket_versioning_update(bucket: &str, config: &VersioningConfiguration) -> S3Result<()> { if config .status .as_ref() .is_some_and(|status| !matches!(status.as_str(), BucketVersioningStatus::ENABLED | BucketVersioningStatus::SUSPENDED)) { return Err(S3Error::with_message( S3ErrorCode::InvalidArgument, "bucket versioning configuration has an invalid status", )); } match metadata_sys::get_object_lock_config(bucket).await { Ok((object_lock_config, _)) => { if object_lock_config.enabled() && versioning_configuration_has_object_lock_incompatible_settings(config) { return Err(S3Error::with_message( S3ErrorCode::InvalidBucketState, "An Object Lock configuration is present on this bucket, versioning cannot be suspended.".to_string(), )); } } Err(StorageError::ConfigNotFound) => {} Err(err) => return Err(ApiError::from(err).into()), } // AWS S3 and MinIO both refuse to suspend versioning while a replication // configuration exists: suspension would start minting null versions that // the replication engine (versioned by contract) can never converge. if config.suspended() { match metadata_sys::get_replication_config(bucket).await { Ok(_) => { return Err(S3Error::with_message( S3ErrorCode::InvalidBucketState, "A replication configuration is present on this bucket, bucket wide versioning cannot be suspended." .to_string(), )); } Err(StorageError::ConfigNotFound) => {} Err(err) => return Err(ApiError::from(err).into()), } } Ok(()) } #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] struct ObjectMetadataPermissions { metadata_allowed: bool, tags_allowed: bool, } fn encode_list_versions_value(value: &str, encoding_type: Option<&EncodingType>) -> String { if encoding_type.is_some_and(|encoding| encoding.as_str() == EncodingType::URL) { encode(value).into_owned() } else { value.to_string() } } fn encode_list_objects_v2_value(value: &str, encoding_type: Option<&EncodingType>) -> String { if encoding_type.is_some_and(|encoding| encoding.as_str() == EncodingType::URL) { value .split('/') .map(|part| encode(part).into_owned()) .collect::>() .join("/") } else { value.to_string() } } fn build_metadata_extension_user_metadata(user_defined: &HashMap) -> Option { let mut items = extract_user_defined_metadata(user_defined) .into_iter() .filter(|(key, _)| !key.is_empty()) .map(|(key, value)| MetadataEntry { name: Some(key), value: Some(value), }) .collect::>(); items.sort_by(|left, right| left.name.cmp(&right.name)); if items.is_empty() { None } else { Some(items) } } fn metadata_count_to_i32(value: usize) -> i32 { i32::try_from(value).unwrap_or(i32::MAX) } async fn is_list_objects_metadata_action_allowed( req: &S3Request, bucket: &str, object: &str, action: S3Action, ) -> S3Result { let mut auth_req = S3Request { input: (), method: req.method.clone(), uri: req.uri.clone(), headers: req.headers.clone(), extensions: req.extensions.clone(), credentials: req.credentials.clone(), region: req.region.clone(), service: req.service.clone(), trailing_headers: req.trailing_headers.clone(), }; let mut req_info = req_info_ref(req)?.clone(); req_info.bucket = Some(bucket.to_string()); req_info.object = Some(object.to_string()); req_info.version_id = None; // Denial here is an expected filter outcome, not an error (issue #5740). req_info.suppress_denial_log = true; auth_req.extensions.insert(req_info); match authorize_request(&mut auth_req, Action::S3Action(action)).await { Ok(()) => Ok(true), Err(err) if err.code() == &S3ErrorCode::AccessDenied => Ok(false), Err(err) => Err(err), } } async fn collect_list_objects_metadata_permissions( req: &S3Request, bucket: &str, objects: &[ObjectInfo], ) -> S3Result> { let mut permissions = HashMap::new(); for object in objects { if object.name.is_empty() || permissions.contains_key(&object.name) { continue; } let metadata_allowed = is_list_objects_metadata_action_allowed(req, bucket, &object.name, S3Action::GetObjectAction).await?; let tags_allowed = is_list_objects_metadata_action_allowed(req, bucket, &object.name, S3Action::GetObjectTaggingAction).await?; permissions.insert( object.name.clone(), ObjectMetadataPermissions { metadata_allowed, tags_allowed, }, ); } Ok(permissions) } fn build_list_object_versions_metadata_output( object_infos: ListObjectVersionsInfo, bucket: &str, params: &ListObjectVersionsParams, encoding_type: Option<&EncodingType>, permissions: &HashMap, ) -> ListObjectVersionsMetadataOutput { let owner = rustfs_owner(); let common_prefixes = object_infos .prefixes .into_iter() .map(|prefix_value| CommonPrefix { prefix: Some(encode_list_versions_value(&prefix_value, encoding_type)), }) .collect::>(); let entries = object_infos .objects .into_iter() .filter(|object| !object.name.is_empty()) .map(|object| { let object_name = encode_list_versions_value(&object.name, encoding_type); let version_id = object .version_id .map(|version| version.to_string()) .unwrap_or_else(|| "null".to_string()); let permission = permissions.get(&object.name).copied().unwrap_or_default(); let user_metadata = if permission.metadata_allowed { build_metadata_extension_user_metadata(&object.user_defined) } else { None }; let user_tags = if permission.tags_allowed && !object.user_tags.is_empty() { Some((*object.user_tags).clone()) } else { None }; let internal = if permission.metadata_allowed && (object.data_blocks > 0 || object.parity_blocks > 0) { Some(ObjectInternalInfo { k: metadata_count_to_i32(object.data_blocks), m: metadata_count_to_i32(object.parity_blocks), }) } else { None }; let extension = ObjectMetadataExtension { user_metadata, user_tags, internal, }; if object.delete_marker { ListObjectVersionMetadataEntry::DeleteMarker( DeleteMarkerEntry { key: Some(object_name), last_modified: object.mod_time.map(Timestamp::from), owner: Some(owner.clone()), version_id: Some(version_id), is_latest: Some(object.is_latest), }, extension, ) } else { ListObjectVersionMetadataEntry::Version( ObjectVersion { key: Some(object_name), last_modified: object.mod_time.map(Timestamp::from), size: Some(object.size), version_id: Some(version_id), is_latest: Some(object.is_latest), e_tag: object.etag.clone().map(|etag| to_s3s_etag(&etag)), storage_class: Some(ObjectVersionStorageClass::from( object .storage_class .unwrap_or_else(|| ObjectVersionStorageClass::STANDARD.to_string()), )), owner: Some(owner.clone()), ..Default::default() }, extension, ) } }) .collect::>(); let next_key_marker = object_infos .next_marker .filter(|marker| !marker.is_empty()) .map(|marker| encode_list_versions_value(&marker, encoding_type)); let next_version_id_marker = object_infos.next_version_idmarker.filter(|marker| !marker.is_empty()); ListObjectVersionsMetadataOutput { output: ListObjectVersionsOutput { common_prefixes: Some(common_prefixes), delimiter: params .delimiter .clone() .map(|value| encode_list_versions_value(&value, encoding_type)), encoding_type: encoding_type.cloned(), is_truncated: Some(object_infos.is_truncated), key_marker: Some(encode_list_versions_value( params.key_marker.as_deref().unwrap_or_default(), encoding_type, )), max_keys: Some(params.max_keys), name: Some(bucket.to_owned()), next_key_marker, next_version_id_marker, prefix: Some(encode_list_versions_value(¶ms.prefix, encoding_type)), request_charged: None, version_id_marker: Some(params.version_id_marker.clone().unwrap_or_default()), ..Default::default() }, entries, } } fn build_list_objects_v2_metadata_output( object_infos: ListObjectsV2Info, bucket: &str, params: &ListObjectsV2Params, encoding_type: Option<&EncodingType>, fetch_owner: bool, permissions: &HashMap, ) -> ListObjectsV2MetadataOutput { let owner = rustfs_owner(); let contents = object_infos .objects .iter() .filter(|object| !object.name.is_empty()) .map(|object| { let permission = permissions.get(&object.name).copied().unwrap_or_default(); let user_metadata = if permission.metadata_allowed { build_metadata_extension_user_metadata(&object.user_defined) } else { None }; let user_tags = if permission.tags_allowed && !object.user_tags.is_empty() { Some((*object.user_tags).clone()) } else { None }; let internal = if permission.metadata_allowed && (object.data_blocks > 0 || object.parity_blocks > 0) { Some(ObjectInternalInfo { k: metadata_count_to_i32(object.data_blocks), m: metadata_count_to_i32(object.parity_blocks), }) } else { None }; ObjectMetadataEntry { object: Object { key: Some(encode_list_objects_v2_value(&object.name, encoding_type)), last_modified: object.mod_time.map(Timestamp::from), size: Some(object.get_actual_size().unwrap_or_default()), e_tag: object.etag.clone().map(|etag| to_s3s_etag(&etag)), storage_class: Some(ObjectStorageClass::from( object .storage_class .clone() .unwrap_or_else(|| ObjectStorageClass::STANDARD.to_string()), )), owner: fetch_owner.then_some(owner.clone()), ..Default::default() }, extension: ObjectMetadataExtension { user_metadata, user_tags, internal, }, } }) .collect::>(); let common_prefixes = object_infos .prefixes .into_iter() .map(|prefix| CommonPrefix { prefix: Some(encode_list_objects_v2_value(&prefix, encoding_type)), }) .collect::>(); let key_count = metadata_count_to_i32(contents.len() + common_prefixes.len()); let next_continuation_token = object_infos .next_continuation_token .map(|token| base64_simd::STANDARD.encode_to_string(token.as_bytes())); ListObjectsV2MetadataOutput { output: ListObjectsV2Output { name: Some(bucket.to_owned()), prefix: Some(params.prefix.clone()), max_keys: Some(params.max_keys), key_count: Some(key_count), continuation_token: params.response_continuation_token.clone(), is_truncated: Some(object_infos.is_truncated), next_continuation_token, common_prefixes: Some(common_prefixes), delimiter: params.delimiter.clone(), encoding_type: encoding_type.cloned(), start_after: params.response_start_after.clone(), ..Default::default() }, contents, } } fn create_bucket_exists_response(is_owner: bool) -> S3Result> { if is_owner { return Ok(S3Response::new(CreateBucketOutput::default())); } Err(s3_error!( BucketAlreadyExists, "The requested bucket name is not available. The bucket namespace is shared by all users of the system. Please select a different name and try again." )) } fn resolve_notification_region(global_region: Option, request_region: Option) -> String { global_region .or(request_region) .map(|region| region.to_string()) .unwrap_or_else(|| RUSTFS_REGION.to_string()) } const ERR_LIFECYCLE_RULE_STATUS: &str = "Rule status must be either Enabled or Disabled"; fn assign_lifecycle_rule_ids(rules: &mut [LifecycleRule]) { let mut rule_ids: HashSet = HashSet::new(); for rule in rules.iter() { if let Some(id) = rule.id.as_ref() { rule_ids.insert(id.to_string()); } } for (idx, rule) in rules.iter_mut().enumerate() { if rule.id.is_none() { let mut suffix = 0usize; let mut generated_id = format!("rule-{}", idx); while rule_ids.contains(&generated_id) { suffix += 1; generated_id = format!("rule-{idx}-{suffix}"); } rule_ids.insert(generated_id.clone()); rule.id = Some(generated_id); } } } fn validate_lifecycle_rule_status(rules: &[LifecycleRule]) -> std::result::Result<(), &'static str> { for rule in rules { if rule.status != ExpirationStatus::from_static(ExpirationStatus::ENABLED) && rule.status != ExpirationStatus::from_static(ExpirationStatus::DISABLED) { return Err(ERR_LIFECYCLE_RULE_STATUS); } } Ok(()) } fn lifecycle_has_transition_rules(config: &BucketLifecycleConfiguration) -> bool { config.rules.iter().any(|rule| { rule.status == ExpirationStatus::from_static(ExpirationStatus::ENABLED) && (rule.transitions.as_ref().is_some_and(|transitions| { transitions.iter().any(|transition| { transition .storage_class .as_ref() .is_some_and(|storage_class| !storage_class.as_str().is_empty()) }) }) || rule.noncurrent_version_transitions.as_ref().is_some_and(|transitions| { transitions.iter().any(|transition| { transition .storage_class .as_ref() .is_some_and(|storage_class| !storage_class.as_str().is_empty()) }) })) }) } fn lifecycle_has_expiry_rules(config: &BucketLifecycleConfiguration) -> bool { config.rules.iter().any(|rule| { rule.status == ExpirationStatus::from_static(ExpirationStatus::ENABLED) && (rule.expiration.is_some() || rule.del_marker_expiration.is_some() || rule.noncurrent_version_expiration.is_some()) }) } fn lifecycle_has_abort_multipart_rules(config: &BucketLifecycleConfiguration) -> bool { config.rules.iter().any(|rule| { rule.status == ExpirationStatus::from_static(ExpirationStatus::ENABLED) && rule.abort_incomplete_multipart_upload.is_some() }) } #[derive(Clone, Default)] pub struct DefaultBucketUsecase { context: Option>, } async fn await_bucket_usecase_on_fresh_task(operation: &'static str, future: F) -> S3Result where T: Send + 'static, F: Future> + Send + 'static, { await_bucket_usecase_on_fresh_task_with_admission(operation, Arc::clone(&BUCKET_OPERATION_ADMISSION), future).await } async fn await_bucket_usecase_on_fresh_task_with_admission( operation: &'static str, admission: Arc, future: F, ) -> S3Result where T: Send + 'static, F: Future> + Send + 'static, { let permit = admission.try_acquire_owned().map_err(|err| match err { TryAcquireError::NoPermits => { S3Error::with_message(S3ErrorCode::SlowDown, format!("{operation} concurrency limit reached; retry later")) } TryAcquireError::Closed => S3Error::with_message(S3ErrorCode::InternalError, format!("{operation} admission closed")), })?; tokio::spawn( async move { let _permit = permit; future.await } .in_current_span(), ) .await .map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("{operation} task failed: {err}")))? } impl DefaultBucketUsecase { #[cfg(test)] pub fn without_context() -> Self { Self { context: None } } pub fn from_global() -> Self { Self { context: current_app_context(), } } /// Build the use-case bound to an explicit application context /// (backlog#1052 S6): the per-server request path passes its own context /// so the use-case resolves that server's store; `None` falls back to the /// ambient default. pub fn with_context(context: Option>) -> Self { Self { context } } fn global_region(&self) -> Option { self.context.as_ref().and_then(|context| context.region().get()) } fn object_store(&self) -> Option> { current_object_store_handle_for_context(self.context.as_deref()) } #[instrument( level = "debug", skip(self, req), fields(start_time=?time::OffsetDateTime::now_utc()) )] pub async fn execute_create_bucket(&self, req: S3Request) -> S3Result> { let usecase = self.clone(); await_bucket_usecase_on_fresh_task("bucket creation", async move { usecase.execute_create_bucket_inner(req).await }).await } async fn execute_create_bucket_inner(&self, req: S3Request) -> S3Result> { let helper = OperationHelper::new(&req, EventName::BucketCreated, S3Operation::CreateBucket); let requester_is_owner = match req_info_ref(&req) { Ok(r) => r.is_owner, Err(_) => { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Missing request info".to_string())); } }; let CreateBucketInput { bucket, object_lock_enabled_for_bucket, .. } = req.input; let lock_enabled = object_lock_enabled_for_bucket.is_some_and(|v| v); let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let make_result = store .make_bucket( &bucket, &MakeBucketOptions { force_create: false, lock_enabled, ..Default::default() }, ) .await; match make_result { Ok(()) => { // Invalidate the bucket validation cache so subsequent GETs // see the newly created bucket immediately. crate::storage::invalidate_bucket_validation_cache(&bucket); } Err(StorageError::BucketExists(_)) => { // Per S3 spec: bucket namespace is global. Owner recreating returns 200 OK; // non-owner gets 409 BucketAlreadyExists. let result = create_bucket_exists_response(requester_is_owner); let _ = helper.complete(&result); return result; } Err(e) => return Err(ApiError::from(e).into()), } if let Err(err) = site_replication_make_bucket_hook(&bucket, lock_enabled).await { warn!(bucket = %bucket, error = ?err, "site replication make bucket hook failed"); } let output = CreateBucketOutput::default(); counter!("rustfs_create_bucket_total").increment(1); let result = Ok(S3Response::new(output)); let _ = helper.complete(&result); rustfs_scanner::record_dirty_usage_bucket(&bucket); result } #[instrument(level = "debug", skip(self, req))] pub async fn execute_delete_bucket(&self, req: S3Request) -> S3Result> { let usecase = self.clone(); await_bucket_usecase_on_fresh_task("bucket deletion", async move { usecase.execute_delete_bucket_inner(req).await }).await } async fn execute_delete_bucket_inner( &self, mut req: S3Request, ) -> S3Result> { let helper = OperationHelper::new(&req, EventName::BucketRemoved, S3Operation::DeleteBucket); let input = req.input.clone(); let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let force_str = get_header(&req.headers, SUFFIX_FORCE_DELETE) .map(|v| v.into_owned()) .unwrap_or_default(); let force = parse_bool(&force_str).unwrap_or_default(); if force { authorize_request(&mut req, Action::S3Action(S3Action::ForceDeleteBucketAction)).await?; } store .delete_bucket( &input.bucket, &DeleteBucketOptions { force, ..Default::default() }, ) .await .map_err(ApiError::from)?; // Drop every cached object body for the now-deleted bucket so dead // bytes do not sit resident until TTL. Covers both the normal and the // force-delete path, which share this single delete_bucket call // (ODC-28, backlog#1133). let cache_adapter = current_object_data_cache_for_context(self.context.as_deref()); let _ = invalidate_object_data_cache_bucket_after_delete(&cache_adapter, &input.bucket).await; // Invalidate bucket validation cache crate::storage::invalidate_bucket_validation_cache(&input.bucket); // Re-evaluate lifecycle and replication after bucket removal. rustfs_scanner::record_scanner_maintenance_change(&input.bucket); if let Err(err) = site_replication_delete_bucket_hook(&input.bucket, force).await { warn!(bucket = %input.bucket, error = ?err, "site replication delete bucket hook failed"); } // Notify peers to drop their cached metadata for the now-deleted bucket. let request_context = req.extensions.get::().cloned(); notify_bucket_metadata_delete(input.bucket.clone(), request_context); let result = Ok(S3Response::new(DeleteBucketOutput {})); let _ = helper.complete(&result); result } #[instrument(level = "debug", skip(self, req))] pub async fn execute_head_bucket(&self, req: S3Request) -> S3Result> { record_s3_op(S3Operation::HeadBucket); let input = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&input.bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; Ok(S3Response::new(HeadBucketOutput::default())) } #[instrument(level = "debug", skip(self, req))] pub async fn execute_get_bucket_location( &self, req: S3Request, ) -> S3Result> { let input = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&input.bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; if let Some(region) = self.global_region() { return Ok(S3Response::new(GetBucketLocationOutput { location_constraint: Some(BucketLocationConstraint::from(region.to_string())), })); } Ok(S3Response::new(GetBucketLocationOutput::default())) } #[instrument(level = "debug", skip(self))] pub async fn execute_list_buckets(&self, req: S3Request) -> S3Result> { let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; if req.credentials.as_ref().is_none_or(|cred| cred.access_key.is_empty()) { return Err(S3Error::with_message(S3ErrorCode::AccessDenied, "Access Denied")); } let iam_authorization = prepare_list_buckets_iam_authorization(&req).await?; let bucket_infos = if iam_authorization.is_allowed("", S3Action::ListAllMyBucketsAction).await { store.list_bucket(&BucketOptions::default()).await.map_err(ApiError::from)? } else { log_list_buckets_iam_implicit_deny(&req)?; let bucket_infos = store.list_bucket(&BucketOptions::default()).await.map_err(ApiError::from)?; let mut visible_bucket_infos = Vec::new(); for info in bucket_infos { if iam_authorization.is_allowed(&info.name, S3Action::ListBucketAction).await || iam_authorization .is_allowed(&info.name, S3Action::GetBucketLocationAction) .await { visible_bucket_infos.push(info); } } if visible_bucket_infos.is_empty() { return Err(ApiError::access_denied().into()); } visible_bucket_infos }; Ok(S3Response::new(build_list_buckets_output(&bucket_infos))) } pub async fn execute_delete_bucket_encryption( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let DeleteBucketEncryptionInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; delete_bucket_config_for_incarnation(&bucket, BUCKET_SSECONFIG, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "delete bucket encryption", request_context, false); let item = sr_bucket_meta_item(bucket.clone(), "sse-config"); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket encryption delete hook failed"); } Ok(S3Response::with_status(DeleteBucketEncryptionOutput::default(), StatusCode::NO_CONTENT)) } #[instrument(level = "debug", skip(self))] pub async fn execute_delete_bucket_cors( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let DeleteBucketCorsInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; delete_bucket_config_for_incarnation(&bucket, BUCKET_CORS_CONFIG, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "delete bucket cors", request_context, false); let item = sr_bucket_meta_item(bucket.clone(), "cors-config"); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket cors delete hook failed"); } Ok(S3Response::new(DeleteBucketCorsOutput {})) } #[instrument(level = "debug", skip(self))] pub async fn execute_delete_bucket_lifecycle( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let DeleteBucketLifecycleInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; delete_bucket_config_for_incarnation(&bucket, BUCKET_LIFECYCLE_CONFIG, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "delete bucket lifecycle", request_context, true); let item = sr_bucket_meta_item(bucket.clone(), "lc-config"); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket lifecycle delete hook failed"); } Ok(S3Response::new(DeleteBucketLifecycleOutput::default())) } pub async fn execute_delete_bucket_policy( &self, req: S3Request, ) -> S3Result> { record_s3_op(S3Operation::DeleteBucketPolicy); let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let DeleteBucketPolicyInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; delete_bucket_config_for_incarnation(&bucket, BUCKET_POLICY_CONFIG, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "delete bucket policy", request_context, false); let item = sr_bucket_meta_item(bucket.clone(), "policy"); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket policy delete hook failed"); } Ok(S3Response::new(DeleteBucketPolicyOutput {})) } pub async fn execute_delete_bucket_replication( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let DeleteBucketReplicationInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let targets_guard = lock_bucket_targets_metadata(&bucket).await; let replication_config = match metadata_sys::get_replication_config(&bucket).await { Ok((config, _)) => Some(config), Err(StorageError::ConfigNotFound) => None, Err(err) => return Err(ApiError::from(err).into()), }; let updated_targets = if let Some(config) = replication_config.as_ref() { replication_targets_without_config_targets(&bucket, config).await? } else { None }; delete_bucket_config_for_incarnation(&bucket, BUCKET_REPLICATION_CONFIG, expected_incarnation_id) .await .map_err(ApiError::from)?; if let Some((targets, removed)) = updated_targets && let Err(err) = write_replication_targets_after_config_delete(&bucket, &targets, removed, expected_incarnation_id).await { if let Some(config) = replication_config.as_ref() { return Err(restore_replication_config_after_target_cleanup_failure( &bucket, config, err, expected_incarnation_id, ) .await); } return Err(err); } drop(targets_guard); notify_bucket_metadata_reload(bucket.clone(), "delete bucket replication", request_context, true); let item = sr_bucket_meta_item(bucket.clone(), "replication-config"); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket replication-config delete hook failed"); } info!(bucket = %bucket, "deleted bucket replication config"); Ok(S3Response::new(DeleteBucketReplicationOutput::default())) } #[instrument(level = "debug", skip(self))] pub async fn execute_delete_bucket_tagging( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let DeleteBucketTaggingInput { bucket, .. } = req.input; delete_bucket_config_for_incarnation(&bucket, BUCKET_TAGGING_CONFIG, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "delete bucket tagging", request_context, false); let item = sr_bucket_meta_item(bucket.clone(), "tags"); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket tagging delete hook failed"); } rustfs_scanner::record_dirty_usage_bucket(&bucket); Ok(S3Response::new(DeleteBucketTaggingOutput {})) } #[instrument(level = "debug", skip(self))] pub async fn execute_delete_public_access_block( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let DeletePublicAccessBlockInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; delete_bucket_config_for_incarnation(&bucket, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "delete public access block", request_context, false); Ok(S3Response::with_status(DeletePublicAccessBlockOutput::default(), StatusCode::NO_CONTENT)) } pub async fn execute_get_bucket_encryption( &self, req: S3Request, ) -> S3Result> { let GetBucketEncryptionInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let server_side_encryption_configuration = match metadata_sys::get_sse_config(&bucket).await { Ok((cfg, _)) => Some(cfg), Err(err) => { if err == StorageError::ConfigNotFound { return Err(s3_error!(ServerSideEncryptionConfigurationNotFoundError)); } warn!( component = LOG_COMPONENT_APP, subsystem = LOG_SUBSYSTEM_BUCKET, event = "bucket_sse_config_load_failed", bucket = %bucket, error = ?err, "Failed to load bucket SSE configuration" ); return Err(ApiError::from(err).into()); } }; Ok(S3Response::new(GetBucketEncryptionOutput { server_side_encryption_configuration, })) } #[instrument(level = "debug", skip(self))] pub async fn execute_get_bucket_cors(&self, req: S3Request) -> S3Result> { let GetBucketCorsInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let cors_configuration = match metadata_sys::get_cors_config(&bucket).await { Ok((config, _)) => config, Err(err) => { if err == StorageError::ConfigNotFound { return Err(S3Error::with_message( S3ErrorCode::NoSuchCORSConfiguration, "The CORS configuration does not exist".to_string(), )); } warn!( component = LOG_COMPONENT_APP, subsystem = LOG_SUBSYSTEM_BUCKET, event = "bucket_cors_config_load_failed", bucket = %bucket, error = ?err, "Failed to load bucket CORS configuration" ); return Err(ApiError::from(err).into()); } }; Ok(S3Response::new(GetBucketCorsOutput { cors_rules: Some(cors_configuration.cors_rules), })) } #[instrument(level = "debug", skip(self))] pub async fn execute_get_bucket_lifecycle_configuration( &self, req: S3Request, ) -> S3Result> { let GetBucketLifecycleConfigurationInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let rules = match metadata_sys::get_lifecycle_config(&bucket).await { Ok((cfg, _)) => cfg.rules, Err(_) => { return Err(s3_error!(NoSuchLifecycleConfiguration)); } }; Ok(S3Response::new(GetBucketLifecycleConfigurationOutput { rules: Some(rules), ..Default::default() })) } pub async fn execute_get_bucket_notification_configuration( &self, req: S3Request, ) -> S3Result> { let GetBucketNotificationConfigurationInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let has_notification_config = metadata_sys::get_notification_config(&bucket).await.unwrap_or_else(|err| { warn!( component = LOG_COMPONENT_APP, subsystem = LOG_SUBSYSTEM_BUCKET, event = "bucket_notification_config_load_failed", bucket = %bucket, error = ?err, "Failed to load bucket notification configuration" ); None }); if let Some(NotificationConfiguration { event_bridge_configuration, lambda_function_configurations, queue_configurations, topic_configurations, }) = has_notification_config { Ok(S3Response::new(GetBucketNotificationConfigurationOutput { event_bridge_configuration, lambda_function_configurations, queue_configurations, topic_configurations, })) } else { Ok(S3Response::new(GetBucketNotificationConfigurationOutput::default())) } } pub async fn execute_get_bucket_policy( &self, req: S3Request, ) -> S3Result> { let GetBucketPolicyInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let (policy_str, _) = match metadata_sys::get_bucket_policy_raw(&bucket).await { Ok(res) => res, Err(err) => { if StorageError::ConfigNotFound == err { return Err(s3_error!(NoSuchBucketPolicy)); } return Err(S3Error::with_message(S3ErrorCode::InternalError, err.to_string())); } }; Ok(S3Response::new(GetBucketPolicyOutput { policy: Some(policy_str), })) } pub async fn execute_get_bucket_policy_status( &self, req: S3Request, ) -> S3Result> { let GetBucketPolicyStatusInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let remote_addr = req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)); let client_info = req.extensions.get::(); let conditions = get_condition_values_with_client_info( &req.headers, &rustfs_credentials::Credentials::default(), None, None, remote_addr, client_info, ); let read_allowed = PolicySys::try_is_allowed(&BucketPolicyArgs { bucket: &bucket, action: Action::S3Action(S3Action::ListBucketAction), is_owner: false, account: "", groups: &None, conditions: &conditions, object: "", }) .await .map_err(ApiError::from)?; let write_allowed = PolicySys::try_is_allowed(&BucketPolicyArgs { bucket: &bucket, action: Action::S3Action(S3Action::PutObjectAction), is_owner: false, account: "", groups: &None, conditions: &conditions, object: "", }) .await .map_err(ApiError::from)?; let mut is_public = read_allowed || write_allowed; let ignore_public_acls = match metadata_sys::get_public_access_block_config(&bucket).await { Ok((config, _)) => config.ignore_public_acls.unwrap_or(false), Err(_) => false, }; let _ = ignore_public_acls; let policy_public = match metadata_sys::get_bucket_policy(&bucket).await { Ok((cfg, _)) => cfg.statements.iter().any(|statement| { matches!(statement.effect, Effect::Allow) && statement.principal.is_match("*") && statement.conditions.is_empty() && statement.actions.is_match(&Action::S3Action(S3Action::ListBucketAction)) }), Err(err) => { if err == StorageError::ConfigNotFound { false } else { return Err(ApiError::from(err).into()); } } }; if policy_public { is_public = true; } let output = GetBucketPolicyStatusOutput { policy_status: Some(PolicyStatus { is_public: Some(is_public), }), }; Ok(S3Response::new(output)) } pub async fn execute_get_bucket_replication( &self, req: S3Request, ) -> S3Result> { let GetBucketReplicationInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let replication_configuration = match metadata_sys::get_replication_config(&bucket).await { Ok((cfg, _created)) => cfg, Err(StorageError::ConfigNotFound) => { return Err(S3Error::with_message( S3ErrorCode::ReplicationConfigurationNotFoundError, "replication not found".to_string(), )); } Err(err) => { warn!( component = LOG_COMPONENT_APP, subsystem = LOG_SUBSYSTEM_BUCKET, event = "bucket_replication_config_load_failed", bucket = %bucket, error = ?err, "Failed to load bucket replication configuration" ); return Err(ApiError::from(err).into()); } }; Ok(S3Response::new(GetBucketReplicationOutput { replication_configuration: Some(replication_configuration), })) } #[instrument(level = "debug", skip(self))] pub async fn execute_get_bucket_tagging( &self, req: S3Request, ) -> S3Result> { let GetBucketTaggingInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let Tagging { tag_set } = match metadata_sys::get_tagging_config(&bucket).await { Ok((tags, _)) => tags, Err(err) => { if err == StorageError::ConfigNotFound { return Err(S3Error::with_message(S3ErrorCode::NoSuchTagSet, "The TagSet does not exist".to_string())); } warn!( component = LOG_COMPONENT_APP, subsystem = LOG_SUBSYSTEM_BUCKET, event = "bucket_tagging_config_load_failed", bucket = %bucket, error = ?err, "Failed to load bucket tagging configuration" ); return Err(ApiError::from(err).into()); } }; Ok(S3Response::new(GetBucketTaggingOutput { tag_set })) } #[instrument(level = "debug", skip(self))] pub async fn execute_get_public_access_block( &self, req: S3Request, ) -> S3Result> { let GetPublicAccessBlockInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let config = match metadata_sys::get_public_access_block_config(&bucket).await { Ok((config, _)) => config, Err(err) => { if err == StorageError::ConfigNotFound { return Err(S3Error::with_message( S3ErrorCode::Custom("NoSuchPublicAccessBlockConfiguration".into()), "Public access block configuration does not exist".to_string(), )); } return Err(ApiError::from(err).into()); } }; Ok(S3Response::new(GetPublicAccessBlockOutput { public_access_block_configuration: Some(config), })) } #[instrument(level = "debug", skip(self))] pub async fn execute_get_bucket_versioning( &self, req: S3Request, ) -> S3Result> { let GetBucketVersioningInput { bucket, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let VersioningConfiguration { status, .. } = BucketVersioningSys::get(&bucket).await.map_err(ApiError::from)?; Ok(S3Response::new(GetBucketVersioningOutput { status, ..Default::default() })) } pub async fn execute_put_bucket_encryption( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let PutBucketEncryptionInput { bucket, mut server_side_encryption_configuration, .. } = req.input; // When SSE-KMS is set without a specific key ID, populate the default // KMS key so that GetBucketEncryption responses include it. Clients like // mc rely on the presence of KMSMasterKeyID to distinguish SSE-KMS from // SSE-S3 in their display logic. if let Some(rule) = server_side_encryption_configuration.rules.first_mut() && let Some(ref mut by_default) = rule.apply_server_side_encryption_by_default && by_default.sse_algorithm.as_str() == ServerSideEncryption::AWS_KMS && by_default.kms_master_key_id.as_deref().is_none_or(str::is_empty) { let service = current_encryption_service() .await .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "KMS service not initialized".to_string()))?; let default_key = service .get_default_key_id() .filter(|key| !key.is_empty()) .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "KMS default key not configured".to_string()))?; by_default.kms_master_key_id = Some(default_key.clone()); } info!("sse_config {:?}", &server_side_encryption_configuration); let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let data = serialize_config(&server_side_encryption_configuration)?; update_bucket_config_for_incarnation(&bucket, BUCKET_SSECONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "put bucket encryption", request_context, false); let mut item = sr_bucket_meta_item(bucket.clone(), "sse-config"); item.sse_config = Some( serialize_config(&server_side_encryption_configuration) .and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?, ); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket encryption hook failed"); } Ok(S3Response::new(PutBucketEncryptionOutput::default())) } #[instrument(level = "debug", skip(self))] pub async fn execute_put_bucket_lifecycle_configuration( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let PutBucketLifecycleConfigurationInput { bucket, lifecycle_configuration, .. } = req.input; let Some(mut input_cfg) = lifecycle_configuration else { return Err(s3_error!(InvalidArgument)) }; assign_lifecycle_rule_ids(&mut input_cfg.rules); if let Err(err) = validate_lifecycle_rule_status(&input_cfg.rules) { return Err(S3Error::with_message(S3ErrorCode::MalformedXML, format!("Malformed XML: {err}"))); } let rcfg = match metadata_sys::get_object_lock_config(&bucket).await { Ok((cfg, _)) => cfg, Err(StorageError::ConfigNotFound) => ObjectLockConfiguration::default(), Err(err) => { warn!( component = LOG_COMPONENT_APP, subsystem = LOG_SUBSYSTEM_BUCKET, event = "bucket_object_lock_config_load_failed", bucket = %bucket, error = ?err, "Failed to load bucket object lock configuration" ); return Err(ApiError::from(err).into()); } }; if let Err(err) = validate_lifecycle_config(&input_cfg, &rcfg).await { return Err(s3_error!(InvalidArgument, "{err}")); } if let Err(err) = validate_transition_tier(&input_cfg).await { return Err(s3_error!(InvalidArgument, "{err}")); } input_cfg.expiry_updated_at = Some(Timestamp::from(time::OffsetDateTime::now_utc())); let data = serialize_config(&input_cfg)?; update_bucket_config_for_incarnation(&bucket, BUCKET_LIFECYCLE_CONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "put bucket lifecycle", request_context, true); let mut item = sr_bucket_meta_item(bucket.clone(), "lc-config"); item.expiry_lc_config = Some(serialize_config(&input_cfg).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?); item.expiry_updated_at = item.updated_at; if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket lifecycle hook failed"); } if lifecycle_has_transition_rules(&input_cfg) && let Some(store) = self.object_store() { let bucket_name = bucket.clone(); let request_context = req.extensions.get::().cloned(); spawn_background_with_context(request_context, async move { if let Err(err) = enqueue_transition_for_existing_objects(store, &bucket_name).await { warn!(bucket = %bucket_name, error = ?err, "failed to enqueue transition for existing objects"); } }); } if lifecycle_has_expiry_rules(&input_cfg) && let Some(store) = self.object_store() { let bucket_name = bucket.clone(); let request_context = req.extensions.get::().cloned(); spawn_background_with_context(request_context, async move { if let Err(err) = enqueue_expiry_for_existing_objects(store, &bucket_name).await { warn!(bucket = %bucket_name, error = ?err, "failed to enqueue expiry for existing objects"); } }); } if lifecycle_has_abort_multipart_rules(&input_cfg) && let Some(store) = self.object_store() { let bucket_name = bucket.clone(); let request_context = req.extensions.get::().cloned(); spawn_background_with_context(request_context, async move { let deleted = run_stale_multipart_upload_cleanup_once(store).await; debug!(bucket = %bucket_name, deleted, "completed lifecycle abort multipart cleanup trigger"); }); } Ok(S3Response::new(PutBucketLifecycleConfigurationOutput::default())) } pub async fn execute_put_bucket_notification_configuration( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_region = req.region.clone(); let request_context = req.extensions.get::().cloned(); let PutBucketNotificationConfigurationInput { bucket, notification_configuration, .. } = req.input; validate_notification_configuration_filters(¬ification_configuration)?; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let data = serialize_config(¬ification_configuration)?; update_bucket_config_for_incarnation(&bucket, BUCKET_NOTIFICATION_CONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "put bucket notification", request_context, false); let region = resolve_notification_region(self.global_region(), request_region); let notify = current_notify_interface_for_context(self.context.as_deref()); let clear_rules = notify.clear_bucket_notification_rules(&bucket); let parse_rules = async { let mut event_rules = Vec::new(); process_queue_configurations(&mut event_rules, notification_configuration.queue_configurations.clone(), |arn_str| { ARN::parse(arn_str) .map(|arn| arn.target_id) .map_err(|e| TargetIDError::InvalidFormat(e.to_string())) })?; process_topic_configurations(&mut event_rules, notification_configuration.topic_configurations.clone(), |arn_str| { ARN::parse(arn_str) .map(|arn| arn.target_id) .map_err(|e| TargetIDError::InvalidFormat(e.to_string())) })?; process_lambda_configurations( &mut event_rules, notification_configuration.lambda_function_configurations.clone(), |arn_str| { ARN::parse(arn_str) .map(|arn| arn.target_id) .map_err(|e| TargetIDError::InvalidFormat(e.to_string())) }, )?; Ok::<_, TargetIDError>(event_rules) }; let (clear_result, event_rules_result) = tokio::join!(clear_rules, parse_rules); clear_result.map_err(|e| s3_error!(InternalError, "Failed to clear rules: {e}"))?; let event_rules = event_rules_result.map_err(|e| s3_error!(InvalidArgument, "Invalid ARN in notification configuration: {e}"))?; warn!("notify event rules: {:?}", &event_rules); notify .add_event_specific_rules(&bucket, region.as_str(), &event_rules) .await .map_err(|e| s3_error!(InternalError, "Failed to add rules: {e}"))?; Ok(S3Response::new(PutBucketNotificationConfigurationOutput {})) } pub async fn execute_put_bucket_policy( &self, req: S3Request, ) -> S3Result> { record_s3_op(S3Operation::PutBucketPolicy); let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let PutBucketPolicyInput { bucket, policy, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let cfg: BucketPolicy = serde_json::from_str(&policy).map_err(|e| s3_error!(InvalidArgument, "parse policy failed {:?}", e))?; if let Err(err) = cfg.is_valid() { warn!("put_bucket_policy err input {:?}, {:?}", &policy, err); return Err(s3_error!(MalformedPolicy)); } let is_public_policy = cfg .statements .iter() .any(|statement| matches!(statement.effect, Effect::Allow) && statement.principal.is_match("*")); if is_public_policy { match metadata_sys::get_public_access_block_config(&bucket).await { Ok((config, _)) => { if config.block_public_policy.unwrap_or(false) { return Err(s3_error!(AccessDenied, "Access Denied")); } } Err(err) => { if err != StorageError::ConfigNotFound { warn!( component = LOG_COMPONENT_APP, subsystem = LOG_SUBSYSTEM_BUCKET, event = "bucket_public_access_block_load_failed", bucket = %bucket, error = ?err, "Failed to load bucket public access block configuration" ); return Err(ApiError::from(err).into()); } } } } let data = policy.as_bytes().to_vec(); update_bucket_config_for_incarnation(&bucket, BUCKET_POLICY_CONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "put bucket policy", request_context, false); let mut item = sr_bucket_meta_item(bucket.clone(), "policy"); item.policy = Some(serde_json::from_str(&policy).map_err(|e| s3_error!(InvalidArgument, "parse policy failed {:?}", e))?); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket policy hook failed"); } Ok(S3Response::new(PutBucketPolicyOutput {})) } #[instrument(level = "debug", skip(self))] pub async fn execute_put_bucket_cors(&self, req: S3Request) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let PutBucketCorsInput { bucket, cors_configuration, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let data = serialize_config(&cors_configuration)?; update_bucket_config_for_incarnation(&bucket, BUCKET_CORS_CONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "put bucket cors", request_context, false); let mut item = sr_bucket_meta_item(bucket.clone(), "cors-config"); item.cors = Some(serialize_config(&cors_configuration).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket cors hook failed"); } Ok(S3Response::new(PutBucketCorsOutput::default())) } pub async fn execute_put_bucket_replication( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let PutBucketReplicationInput { bucket, replication_configuration, .. } = req.input; info!(bucket = %bucket, "updating bucket replication config"); validate_replication_config_capabilities(&replication_configuration)?; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let targets_guard = lock_bucket_targets_metadata(&bucket).await; validate_bucket_replication_update(&bucket, &replication_configuration).await?; let data = serialize_config(&replication_configuration)?; update_bucket_config_for_incarnation(&bucket, BUCKET_REPLICATION_CONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; drop(targets_guard); notify_bucket_metadata_reload(bucket.clone(), "put bucket replication", request_context, true); let mut item = sr_bucket_meta_item(bucket.clone(), "replication-config"); item.replication_config = Some( serialize_config(&replication_configuration).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?, ); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket replication-config hook failed"); } Ok(S3Response::new(PutBucketReplicationOutput::default())) } #[instrument(level = "debug", skip(self))] pub async fn execute_put_public_access_block( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let PutPublicAccessBlockInput { bucket, public_access_block_configuration, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let data = serialize_config(&public_access_block_configuration)?; update_bucket_config_for_incarnation(&bucket, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "put public access block", request_context, false); Ok(S3Response::new(PutPublicAccessBlockOutput::default())) } #[instrument(level = "debug", skip(self))] pub async fn execute_put_bucket_tagging( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let PutBucketTaggingInput { bucket, tagging, .. } = req.input; let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; let data = serialize_config(&tagging)?; update_bucket_config_for_incarnation(&bucket, BUCKET_TAGGING_CONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "put bucket tagging", request_context, false); let mut item = sr_bucket_meta_item(bucket.clone(), "tags"); item.tags = Some(serialize_config(&tagging).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket tagging hook failed"); } rustfs_scanner::record_dirty_usage_bucket(&bucket); Ok(S3Response::new(PutBucketTaggingOutput::default())) } #[instrument(level = "debug", skip(self))] pub async fn execute_put_bucket_versioning( &self, req: S3Request, ) -> S3Result> { let expected_incarnation_id = bucket_config_mutation_incarnation(&req, &req.input.bucket)?; let request_context = req.extensions.get::().cloned(); let PutBucketVersioningInput { bucket, versioning_configuration, .. } = req.input; validate_bucket_versioning_update(&bucket, &versioning_configuration).await?; let data = serialize_config(&versioning_configuration)?; update_bucket_config_for_incarnation(&bucket, BUCKET_VERSIONING_CONFIG, data, expected_incarnation_id) .await .map_err(ApiError::from)?; notify_bucket_metadata_reload(bucket.clone(), "put bucket versioning", request_context, false); let mut item = sr_bucket_meta_item(bucket.clone(), "version-config"); item.versioning = Some( serialize_config(&versioning_configuration).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?, ); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket versioning hook failed"); } rustfs_scanner::record_dirty_usage_bucket(&bucket); Ok(S3Response::new(PutBucketVersioningOutput {})) } #[instrument(level = "trace", skip(self, req))] pub async fn execute_list_objects_v2(&self, req: S3Request) -> S3Result> { // warn!("list_objects_v2 req {:?}", &req.input); let ListObjectsV2Input { bucket, continuation_token, delimiter, encoding_type, fetch_owner, max_keys, prefix, start_after, .. } = req.input; let params = parse_list_objects_v2_params(prefix, delimiter, max_keys, continuation_token, start_after)?; validate_list_object_unordered_with_delimiter(params.delimiter.as_ref(), req.uri.query())?; let store = get_validated_store(&bucket).await?; let incl_deleted = get_header(&req.headers, rustfs_utils::http::SUFFIX_INCLUDE_DELETED) .map(|v| v.as_ref() == "true") .unwrap_or_default(); let object_infos = store .list_objects_v2( &bucket, ¶ms.prefix, params.decoded_continuation_token.clone(), params.delimiter.clone(), params.max_keys, fetch_owner.unwrap_or_default(), params.start_after_for_query.clone(), incl_deleted, ) .await .map_err(ApiError::from)?; let output = build_list_objects_v2_output( object_infos, fetch_owner.unwrap_or_default(), params.max_keys, bucket, params.prefix, params.delimiter, encoding_type, params.response_continuation_token, params.response_start_after, ); Ok(S3Response::new(output)) } pub(crate) async fn execute_list_objects_v2m( &self, req: S3Request, ) -> S3Result> { let input = req.input.clone(); let ListObjectsV2Input { bucket, continuation_token, delimiter, encoding_type, fetch_owner, max_keys, prefix, start_after, .. } = input; let params = parse_list_objects_v2_params(prefix, delimiter, max_keys, continuation_token, start_after)?; validate_list_object_unordered_with_delimiter(params.delimiter.as_ref(), req.uri.query())?; let store = get_validated_store(&bucket).await?; let incl_deleted = get_header(&req.headers, rustfs_utils::http::SUFFIX_INCLUDE_DELETED) .map(|value| value.as_ref() == "true") .unwrap_or_default(); let object_infos = store .list_objects_v2( &bucket, ¶ms.prefix, params.decoded_continuation_token.clone(), params.delimiter.clone(), params.max_keys, fetch_owner.unwrap_or_default(), params.start_after_for_query.clone(), incl_deleted, ) .await .map_err(ApiError::from)?; let permissions = collect_list_objects_metadata_permissions(&req, &bucket, &object_infos.objects).await?; let output = build_list_objects_v2_metadata_output( object_infos, &bucket, ¶ms, encoding_type.as_ref(), fetch_owner.unwrap_or_default(), &permissions, ); Ok(S3Response::new(output)) } pub async fn execute_list_object_versions( &self, req: S3Request, ) -> S3Result> { let ListObjectVersionsInput { bucket, delimiter, encoding_type, key_marker, version_id_marker, max_keys, prefix, .. } = req.input; let params = parse_list_object_versions_params(prefix, delimiter, key_marker, version_id_marker, max_keys)?; let store = get_validated_store(&bucket).await?; let object_infos = store .list_object_versions( &bucket, ¶ms.prefix, params.key_marker.clone(), params.version_id_marker.clone(), params.delimiter.clone(), params.max_keys, ) .await .map_err(ApiError::from)?; let output = build_list_object_versions_output(object_infos, bucket, ¶ms, encoding_type.as_ref()); Ok(S3Response::new(output)) } pub(crate) async fn execute_list_object_versions_m( &self, req: S3Request, ) -> S3Result> { let input = req.input.clone(); let ListObjectVersionsInput { bucket, delimiter, encoding_type, key_marker, version_id_marker, max_keys, prefix, .. } = input; let params = parse_list_object_versions_params(prefix, delimiter, key_marker, version_id_marker, max_keys)?; let store = get_validated_store(&bucket).await?; let object_infos = store .list_object_versions( &bucket, ¶ms.prefix, params.key_marker.clone(), params.version_id_marker.clone(), params.delimiter.clone(), params.max_keys, ) .await .map_err(ApiError::from)?; let permissions = collect_list_objects_metadata_permissions(&req, &bucket, &object_infos.objects).await?; let output = build_list_object_versions_metadata_output(object_infos, &bucket, ¶ms, encoding_type.as_ref(), &permissions); Ok(S3Response::new(output)) } #[instrument(level = "debug", skip(self, req))] pub async fn execute_list_objects(&self, req: S3Request) -> S3Result> { let request_marker = req.input.marker.clone(); let v2_resp = self.execute_list_objects_v2(req.map_input(Into::into)).await?; Ok(v2_resp.map_output(|v2| build_list_objects_output(v2, request_marker))) } } #[cfg(test)] mod tests { use super::*; use http::{Extensions, HeaderMap, Method, Uri}; use s3s::dto::ReplicationRuleStatus; use s3s::dto::{ BucketVersioningStatus, CORSConfiguration, Destination, ExcludedPrefix, FilterRule, FilterRuleName, LifecycleExpiration, NoncurrentVersionTransition, PublicAccessBlockConfiguration, QueueConfiguration, ReplicationRule, S3KeyFilter, ServerSideEncryptionConfiguration, Tag, Transition, TransitionStorageClass, }; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::time::Duration; use tokio::sync::Notify; fn s3_op_total(op: S3Operation) -> u64 { rustfs_io_metrics::s3_op_metrics_snapshot() .into_iter() .find(|snapshot| snapshot.op == op.as_str()) .map(|snapshot| snapshot.total) .unwrap_or_default() } #[tokio::test] async fn bucket_usecase_task_finishes_post_commit_hooks_after_parent_cancellation() { let admission = Arc::new(Semaphore::new(1)); let committed = Arc::new(Notify::new()); let committed_wait = committed.notified(); let release_hook = Arc::new(Notify::new()); let finished = Arc::new(Notify::new()); let finished_wait = finished.notified(); let hook_ran = Arc::new(AtomicBool::new(false)); let committed_for_task = committed.clone(); let release_hook_for_task = release_hook.clone(); let finished_for_task = finished.clone(); let hook_ran_for_task = hook_ran.clone(); let parent = tokio::spawn(await_bucket_usecase_on_fresh_task_with_admission( "test bucket operation", admission.clone(), async move { committed_for_task.notify_one(); release_hook_for_task.notified().await; hook_ran_for_task.store(true, Ordering::SeqCst); finished_for_task.notify_one(); Ok(()) }, )); committed_wait.await; parent.abort(); let _ = parent.await; release_hook.notify_waiters(); tokio::time::timeout(Duration::from_secs(1), finished_wait) .await .expect("the post-commit hook should finish after request cancellation"); assert!( hook_ran.load(Ordering::SeqCst), "the fresh task must own both the storage mutation and its post-commit hooks" ); assert_eq!(admission.available_permits(), 1); } #[tokio::test] async fn saturated_bucket_usecase_admission_does_not_start_work() { let admission = Arc::new(Semaphore::new(1)); let held = admission .clone() .acquire_owned() .await .expect("test admission should remain open"); let started = Arc::new(AtomicBool::new(false)); let started_for_task = started.clone(); let result = tokio::time::timeout( Duration::from_secs(1), await_bucket_usecase_on_fresh_task_with_admission("test bucket operation", admission, async move { started_for_task.store(true, Ordering::SeqCst); Ok(()) }), ) .await .expect("saturated admission must fail within the bounded timeout"); drop(held); assert!( !started.load(Ordering::SeqCst), "saturated admission must not start a detached bucket operation" ); assert_eq!(result.expect_err("saturated admission must fail fast").code(), &S3ErrorCode::SlowDown); } #[tokio::test] async fn bucket_usecase_admission_rejects_excess_detached_tasks() { let admission = Arc::new(Semaphore::new(BUCKET_OPERATION_CONCURRENCY)); let release = Arc::new(Semaphore::new(0)); let started = Arc::new(AtomicUsize::new(0)); let mut tasks = Vec::with_capacity(BUCKET_OPERATION_CONCURRENCY); for _ in 0..BUCKET_OPERATION_CONCURRENCY { let admission_for_task = admission.clone(); let release_for_task = release.clone(); let started_for_task = started.clone(); tasks.push(tokio::spawn(await_bucket_usecase_on_fresh_task_with_admission( "test bucket operation", admission_for_task, async move { started_for_task.fetch_add(1, Ordering::SeqCst); let _release = release_for_task .acquire() .await .expect("test release gate should remain open"); Ok(()) }, ))); } tokio::time::timeout(Duration::from_secs(1), async { while started.load(Ordering::SeqCst) != BUCKET_OPERATION_CONCURRENCY { tokio::task::yield_now().await; } }) .await .expect("all admitted operations should start"); let ninth_started = Arc::new(AtomicBool::new(false)); let ninth_started_for_task = ninth_started.clone(); let ninth_result = tokio::time::timeout( Duration::from_secs(1), await_bucket_usecase_on_fresh_task_with_admission("test bucket operation", admission.clone(), async move { ninth_started_for_task.store(true, Ordering::SeqCst); Ok(()) }), ) .await .expect("an excess bucket transaction must fail within the bounded timeout"); assert!( !ninth_started.load(Ordering::SeqCst), "the ninth bucket transaction must not start when admission is saturated" ); assert_eq!( ninth_result.expect_err("the ninth bucket transaction must fail fast").code(), &S3ErrorCode::SlowDown ); release.add_permits(BUCKET_OPERATION_CONCURRENCY); for task in tasks { task.await .expect("bucket transaction parent should join") .expect("bucket transaction should succeed"); } await_bucket_usecase_on_fresh_task_with_admission("test bucket operation", admission, async { Ok(()) }) .await .expect("admission must recover after active transactions finish"); } #[tokio::test] async fn bucket_usecase_task_maps_panics_to_internal_errors() { async fn panicking_bucket_operation() -> S3Result<()> { panic!("injected bucket task panic"); } let result = await_bucket_usecase_on_fresh_task("test bucket operation", panicking_bucket_operation()).await; let err = result.expect_err("a bucket task panic must be returned as an error"); assert_eq!(err.code(), &S3ErrorCode::InternalError); assert!(err.message().is_some_and(|message| message.contains("task failed"))); } fn build_request(input: T, method: Method) -> S3Request { S3Request { input, method, uri: Uri::from_static("/"), headers: HeaderMap::new(), extensions: Extensions::new(), credentials: None, region: None, service: None, trailing_headers: None, } } fn build_request_with_req_info(input: T, method: Method, req_info: ReqInfo) -> S3Request { let mut req = build_request(input, method); req.extensions.insert(req_info); req } fn usecase_method_source<'a>(source: &'a str, method: &str) -> &'a str { let start_marker = format!("pub async fn {method}"); let start = source.find(&start_marker).expect("method should exist"); let rest = &source[start + start_marker.len()..]; let end = rest.find("\n pub async fn ").unwrap_or(rest.len()); &rest[..end] } #[test] fn bucket_create_and_delete_use_admitted_fresh_tasks() { let source = include_str!("bucket_usecase.rs"); for (method, operation, inner_method) in [ ("execute_create_bucket", "bucket creation", "execute_create_bucket_inner"), ("execute_delete_bucket", "bucket deletion", "execute_delete_bucket_inner"), ] { let body = usecase_method_source(source, method); assert!( body.contains(&format!("await_bucket_usecase_on_fresh_task(\"{operation}\"")), "{method} must use bounded detached-task admission" ); assert!( body.contains(&format!("{inner_method}(req).await")), "{method} must run its mutation inside the admitted task" ); } } #[test] fn bucket_metadata_config_changes_notify_peer_metadata_reload() { let source = include_str!("bucket_usecase.rs"); for (method, operation, scanner_maintenance_change) in [ ("execute_delete_bucket_policy", "delete bucket policy", false), ("execute_put_bucket_policy", "put bucket policy", false), ("execute_delete_public_access_block", "delete public access block", false), ("execute_put_public_access_block", "put public access block", false), ("execute_delete_bucket_lifecycle", "delete bucket lifecycle", true), ("execute_put_bucket_lifecycle_configuration", "put bucket lifecycle", true), ("execute_put_bucket_versioning", "put bucket versioning", false), ("execute_delete_bucket_tagging", "delete bucket tagging", false), ("execute_put_bucket_tagging", "put bucket tagging", false), ("execute_delete_bucket_replication", "delete bucket replication", true), ("execute_put_bucket_replication", "put bucket replication", true), ("execute_delete_bucket_cors", "delete bucket cors", false), ("execute_put_bucket_cors", "put bucket cors", false), ("execute_delete_bucket_encryption", "delete bucket encryption", false), ("execute_put_bucket_encryption", "put bucket encryption", false), ("execute_put_bucket_notification_configuration", "put bucket notification", false), ] { let body = usecase_method_source(source, method); assert!( body.contains("notify_bucket_metadata_reload("), "{method} should notify peers to reload cached bucket metadata" ); assert!( body.contains(operation), "{method} should identify the bucket metadata operation in reload logs" ); let expected_reload = format!( "notify_bucket_metadata_reload(bucket.clone(), \"{operation}\", request_context, {scanner_maintenance_change});" ); assert!( body.contains(&expected_reload), "{method} should propagate scanner_maintenance_change={scanner_maintenance_change}" ); } } #[test] #[serial_test::serial] fn scanner_maintenance_reload_marks_only_scanner_owned_config_changes() { const BUCKET: &str = "scanner-maintenance-reload-test"; rustfs_scanner::clear_dirty_usage_bucket(BUCKET); let before = rustfs_scanner::scanner_maintenance_generation(); record_local_scanner_maintenance_reload(BUCKET, false); assert_eq!(rustfs_scanner::scanner_maintenance_generation(), before); record_local_scanner_maintenance_reload(BUCKET, true); assert_eq!(rustfs_scanner::scanner_maintenance_generation(), before.saturating_add(1)); rustfs_scanner::clear_dirty_usage_bucket(BUCKET); } fn replication_rule_for_target(arn: &str) -> ReplicationRule { ReplicationRule { delete_marker_replication: None, delete_replication: None, destination: Destination { bucket: arn.to_string(), ..Default::default() }, existing_object_replication: None, filter: None, id: Some("rule-1".to_string()), prefix: None, priority: Some(1), source_selection_criteria: None, status: ReplicationRuleStatus::from_static(ReplicationRuleStatus::ENABLED), } } #[test] fn replication_target_arns_use_role_when_present() { let role = "arn:rustfs:replication:us-east-1:source:bucket"; let destination = "arn:rustfs:replication:us-east-1:target:bucket"; let config = ReplicationConfiguration { role: format!(" {role} "), rules: vec![replication_rule_for_target(destination)], }; let arns = replication_target_arns(&config); assert!(arns.contains(role)); assert!(!arns.contains(destination)); } #[test] fn replication_target_arns_use_rule_destinations_without_role() { let destination = "arn:rustfs:replication:us-east-1:target:bucket"; let config = ReplicationConfiguration { role: String::new(), rules: vec![replication_rule_for_target(destination)], }; let arns = replication_target_arns(&config); assert!(arns.contains(destination)); } fn replication_targets_with_arn(arns: &[&str]) -> BucketTargets { BucketTargets { targets: arns .iter() .map(|arn| BucketTarget { arn: (*arn).to_string(), target_type: BucketTargetType::ReplicationService, ..Default::default() }) .collect(), } } #[test] fn validate_replication_config_targets_accepts_matching_destination_arns() { let arn = "arn:rustfs:replication:us-east-1:target:bucket"; let targets = replication_targets_with_arn(&[arn]); let config = ReplicationConfiguration { role: String::new(), rules: vec![replication_rule_for_target(arn)], }; validate_replication_config_targets(&targets, &config).expect("matching target should pass validation"); } #[test] fn validate_replication_config_targets_rejects_stale_destination_arns() { let targets = replication_targets_with_arn(&["arn:rustfs:replication:us-east-1:target:bucket-a"]); let config = ReplicationConfiguration { role: String::new(), rules: vec![replication_rule_for_target( "arn:rustfs:replication:us-east-1:target:bucket-b", )], }; let err = validate_replication_config_targets(&targets, &config).expect_err("stale target should fail validation"); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); } #[test] fn validate_replication_config_targets_accepts_matching_role_arn() { let arn = "arn:rustfs:replication:us-east-1:role-target:bucket"; let targets = replication_targets_with_arn(&[arn]); let config = ReplicationConfiguration { role: format!(" {arn} "), rules: vec![replication_rule_for_target("arn:rustfs:replication:us-east-1:ignored:bucket")], }; validate_replication_config_targets(&targets, &config).expect("matching role ARN should pass validation"); } #[test] fn validate_replication_config_targets_rejects_role_with_multiple_destinations() { let role = "arn:rustfs:replication:us-east-1:role-target:bucket"; let targets = replication_targets_with_arn(&[role]); let config = ReplicationConfiguration { role: role.to_string(), rules: vec![ replication_rule_for_target("arn:rustfs:replication:us-east-1:target-a:bucket"), replication_rule_for_target("arn:rustfs:replication:us-east-1:target-b:bucket"), ], }; let err = validate_replication_config_targets(&targets, &config) .expect_err("role plus multiple destinations should be rejected"); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); } #[test] fn validate_replication_config_targets_trims_destination_arns() { let arn = "arn:rustfs:replication:us-east-1:target:bucket"; let targets = replication_targets_with_arn(&[arn]); let config = ReplicationConfiguration { role: String::new(), rules: vec![replication_rule_for_target( " arn:rustfs:replication:us-east-1:target:bucket ", )], }; validate_replication_config_targets(&targets, &config).expect("trimmed destination ARN should match configured target"); } #[test] fn validate_replication_config_targets_ignores_disabled_rules() { let targets = replication_targets_with_arn(&[]); let mut rule = replication_rule_for_target("arn:rustfs:replication:us-east-1:stale:bucket"); rule.status = ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED); let config = ReplicationConfiguration { role: String::new(), rules: vec![rule], }; validate_replication_config_targets(&targets, &config).expect("disabled rules should not require live targets"); } #[test] fn validate_replication_config_capabilities_names_unsupported_field() { let mut rule = replication_rule_for_target("arn:rustfs:replication:us-east-1:target:bucket"); let destination_key_id = "arn:aws:kms:us-east-1:123456789012:key/opaque-key-id"; rule.destination.encryption_configuration = Some(s3s::dto::EncryptionConfiguration { replica_kms_key_id: Some(destination_key_id.to_string()), }); let config = ReplicationConfiguration { role: String::new(), rules: vec![rule], }; let err = validate_replication_config_capabilities(&config) .expect_err("destination encryption must be rejected until the execution path supports it"); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); assert!( err.to_string() .contains("Destination.EncryptionConfiguration is not supported") ); assert!(!err.to_string().contains(destination_key_id)); } #[test] fn validate_replication_config_capabilities_rejects_structural_defects_before_write() { let mut first = replication_rule_for_target("arn:rustfs:replication:us-east-1:target:bucket"); first.priority = Some(1); let mut second = replication_rule_for_target("arn:rustfs:replication:us-east-1:target:bucket"); second.priority = Some(1); let config = ReplicationConfiguration { role: String::new(), rules: vec![first, second], }; let err = validate_replication_config_capabilities(&config) .expect_err("duplicate rule priorities must be rejected before persistence"); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); assert!(err.to_string().contains("Priority must be unique")); } #[test] fn validate_replication_config_capabilities_rejects_invalid_status_before_write() { let mut rule = replication_rule_for_target("arn:rustfs:replication:us-east-1:target:bucket"); rule.status = ReplicationRuleStatus::from_static("Invalid"); let config = ReplicationConfiguration { role: String::new(), rules: vec![rule], }; let err = validate_replication_config_capabilities(&config) .expect_err("an invalid string-backed replication status must be rejected before persistence"); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); assert!(err.to_string().contains("Rule.Status has an invalid status")); } #[tokio::test] async fn validate_bucket_versioning_update_rejects_invalid_status_before_metadata_lookup() { let config = VersioningConfiguration { status: Some(BucketVersioningStatus::from_static("Invalid")), ..Default::default() }; let err = validate_bucket_versioning_update("unregistered-test-bucket", &config) .await .expect_err("an invalid string-backed versioning status must be rejected before persistence"); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[test] fn remove_replication_targets_from_config_targets_only_removes_referenced_replication_targets() { let removed_arn = "arn:rustfs:replication:us-east-1:removed:bucket"; let kept_replication_arn = "arn:rustfs:replication:us-east-1:kept:bucket"; let kept_ilm_arn = "arn:rustfs:ilm:us-east-1:kept:bucket"; let mut targets = BucketTargets { targets: vec![ BucketTarget { arn: removed_arn.to_string(), target_type: BucketTargetType::ReplicationService, ..Default::default() }, BucketTarget { arn: kept_replication_arn.to_string(), target_type: BucketTargetType::ReplicationService, ..Default::default() }, BucketTarget { arn: kept_ilm_arn.to_string(), target_type: BucketTargetType::IlmService, ..Default::default() }, ], }; let target_arns = HashSet::from([removed_arn.to_string(), kept_ilm_arn.to_string()]); let removed = remove_replication_targets_from_config_targets(&mut targets, &target_arns); assert_eq!(removed, 1); let remaining_arns = targets .targets .iter() .map(|target| target.arn.as_str()) .collect::>(); assert!(!remaining_arns.contains(removed_arn)); assert!(remaining_arns.contains(kept_replication_arn)); assert!(remaining_arns.contains(kept_ilm_arn)); } #[test] fn versioning_configuration_has_object_lock_incompatible_settings_rejects_suspended() { let config = VersioningConfiguration { status: Some(BucketVersioningStatus::from_static(BucketVersioningStatus::SUSPENDED)), ..Default::default() }; assert!(versioning_configuration_has_object_lock_incompatible_settings(&config)); } #[test] fn versioning_configuration_has_object_lock_incompatible_settings_rejects_exclude_folders() { let config = VersioningConfiguration { exclude_folders: Some(true), ..Default::default() }; assert!(versioning_configuration_has_object_lock_incompatible_settings(&config)); } #[test] fn versioning_configuration_has_object_lock_incompatible_settings_rejects_excluded_prefixes() { let config = VersioningConfiguration { excluded_prefixes: Some(vec![ExcludedPrefix { prefix: Some("archive/".to_string()), }]), ..Default::default() }; assert!(versioning_configuration_has_object_lock_incompatible_settings(&config)); } #[test] fn resolve_notification_region_prefers_global_region() { let binding = resolve_notification_region(Some("us-east-1".parse().unwrap()), Some("ap-southeast-1".parse().unwrap())); assert_eq!(binding, "us-east-1"); } #[test] fn resolve_notification_region_falls_back_to_request_region() { let binding = resolve_notification_region(None, Some("ap-southeast-1".parse().unwrap())); assert_eq!(binding, "ap-southeast-1"); } #[test] fn resolve_notification_region_defaults_value() { let binding = resolve_notification_region(None, None); assert_eq!(binding, RUSTFS_REGION); } #[test] fn create_bucket_exists_response_returns_ok_for_owner() { let response = create_bucket_exists_response(true).expect("owner recreate should succeed"); assert_eq!(response.output.location, None); } #[test] fn create_bucket_exists_response_returns_bucket_already_exists_for_non_owner() { let err = create_bucket_exists_response(false).expect_err("non-owner recreate should fail"); assert_eq!(err.code(), &S3ErrorCode::BucketAlreadyExists); } #[test] fn build_request_with_req_info_preserves_owner_state() { let input = CreateBucketInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request_with_req_info( input, Method::PUT, ReqInfo { is_owner: true, ..Default::default() }, ); assert!(req_info_ref(&req).expect("req info should be present").is_owner); } #[tokio::test] async fn execute_create_bucket_returns_internal_error_when_store_uninitialized() { let input = CreateBucketInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let err = Box::pin(usecase.execute_create_bucket(req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_delete_bucket_returns_internal_error_when_store_uninitialized() { let input = DeleteBucketInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::DELETE); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_delete_bucket(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_delete_bucket_cors_returns_internal_error_when_store_uninitialized() { let input = DeleteBucketCorsInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::DELETE); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_delete_bucket_cors(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_delete_bucket_replication_returns_internal_error_when_store_uninitialized() { let input = DeleteBucketReplicationInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::DELETE); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_delete_bucket_replication(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_head_bucket_returns_internal_error_when_store_uninitialized() { let input = HeadBucketInput::builder().bucket("test-bucket".to_string()).build().unwrap(); let req = build_request(input, Method::HEAD); let usecase = DefaultBucketUsecase::without_context(); let before = s3_op_total(S3Operation::HeadBucket); let err = usecase.execute_head_bucket(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); assert_eq!(s3_op_total(S3Operation::HeadBucket), before + 1); } #[tokio::test] async fn execute_delete_bucket_policy_records_s3_operation_before_store_lookup() { let input = DeleteBucketPolicyInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::DELETE); let usecase = DefaultBucketUsecase::without_context(); let before = s3_op_total(S3Operation::DeleteBucketPolicy); let err = usecase.execute_delete_bucket_policy(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); assert_eq!(s3_op_total(S3Operation::DeleteBucketPolicy), before + 1); } #[tokio::test] async fn execute_get_bucket_policy_returns_internal_error_when_store_uninitialized() { let input = GetBucketPolicyInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_get_bucket_policy(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_get_bucket_location_returns_internal_error_when_store_uninitialized() { let input = GetBucketLocationInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_get_bucket_location(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_get_bucket_cors_returns_internal_error_when_store_uninitialized() { let input = GetBucketCorsInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_get_bucket_cors(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_delete_public_access_block_returns_internal_error_when_store_uninitialized() { let input = DeletePublicAccessBlockInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::DELETE); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_delete_public_access_block(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_get_bucket_encryption_returns_internal_error_when_store_uninitialized() { let input = GetBucketEncryptionInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_get_bucket_encryption(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_get_bucket_replication_returns_internal_error_when_store_uninitialized() { let input = GetBucketReplicationInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_get_bucket_replication(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_get_public_access_block_returns_internal_error_when_store_uninitialized() { let input = GetPublicAccessBlockInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_get_public_access_block(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_get_bucket_tagging_returns_internal_error_when_store_uninitialized() { let input = GetBucketTaggingInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_get_bucket_tagging(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_get_bucket_versioning_returns_internal_error_when_store_uninitialized() { let input = GetBucketVersioningInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_get_bucket_versioning(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[test] fn normalize_lifecycle_rules_generate_rule_ids_for_missing_values() { let mut rules = vec![ LifecycleRule { status: ExpirationStatus::from_static(ExpirationStatus::ENABLED), expiration: Some(LifecycleExpiration { days: Some(30), ..Default::default() }), abort_incomplete_multipart_upload: None, del_marker_expiration: None, filter: None, id: None, noncurrent_version_expiration: None, noncurrent_version_transitions: None, prefix: None, transitions: None, }, LifecycleRule { status: ExpirationStatus::from_static(ExpirationStatus::ENABLED), expiration: Some(LifecycleExpiration { days: Some(60), ..Default::default() }), abort_incomplete_multipart_upload: None, del_marker_expiration: None, filter: None, id: Some("rule-1".to_string()), noncurrent_version_expiration: None, noncurrent_version_transitions: None, prefix: None, transitions: None, }, LifecycleRule { status: ExpirationStatus::from_static(ExpirationStatus::ENABLED), expiration: Some(LifecycleExpiration { days: Some(90), ..Default::default() }), abort_incomplete_multipart_upload: None, del_marker_expiration: None, filter: None, id: None, noncurrent_version_expiration: None, noncurrent_version_transitions: None, prefix: None, transitions: None, }, ]; assign_lifecycle_rule_ids(&mut rules); assert_eq!(rules[0].id.as_deref(), Some("rule-0")); assert_eq!(rules[1].id.as_deref(), Some("rule-1")); assert_eq!(rules[2].id.as_deref(), Some("rule-2")); } #[test] fn validate_lifecycle_rule_status_rejects_invalid_status() { let rules = vec![LifecycleRule { status: ExpirationStatus::from_static("enabled"), expiration: Some(LifecycleExpiration { days: Some(30), ..Default::default() }), abort_incomplete_multipart_upload: None, del_marker_expiration: None, filter: None, id: None, noncurrent_version_expiration: None, noncurrent_version_transitions: None, prefix: None, transitions: None, }]; assert_eq!(validate_lifecycle_rule_status(&rules).unwrap_err(), ERR_LIFECYCLE_RULE_STATUS); } #[test] fn lifecycle_has_transition_rules_ignores_disabled_rules() { let config = BucketLifecycleConfiguration { expiry_updated_at: None, rules: vec![LifecycleRule { status: ExpirationStatus::from_static(ExpirationStatus::DISABLED), expiration: None, abort_incomplete_multipart_upload: None, del_marker_expiration: None, filter: None, id: Some("disabled-transition".to_string()), noncurrent_version_expiration: None, noncurrent_version_transitions: None, prefix: None, transitions: Some(vec![Transition { days: Some(1), date: None, storage_class: Some(TransitionStorageClass::from_static("WARM")), }]), }], }; assert!(!lifecycle_has_transition_rules(&config)); } #[test] fn lifecycle_has_transition_rules_accepts_enabled_noncurrent_transitions() { let config = BucketLifecycleConfiguration { expiry_updated_at: None, rules: vec![LifecycleRule { status: ExpirationStatus::from_static(ExpirationStatus::ENABLED), expiration: None, abort_incomplete_multipart_upload: None, del_marker_expiration: None, filter: None, id: Some("enabled-noncurrent-transition".to_string()), noncurrent_version_expiration: None, noncurrent_version_transitions: Some(vec![NoncurrentVersionTransition { noncurrent_days: Some(1), newer_noncurrent_versions: None, storage_class: Some(TransitionStorageClass::from_static("WARM")), }]), prefix: None, transitions: None, }], }; assert!(lifecycle_has_transition_rules(&config)); } #[test] fn lifecycle_has_abort_multipart_rules_accepts_enabled_abort_rule() { let config = BucketLifecycleConfiguration { expiry_updated_at: None, rules: vec![LifecycleRule { status: ExpirationStatus::from_static(ExpirationStatus::ENABLED), expiration: None, abort_incomplete_multipart_upload: Some(s3s::dto::AbortIncompleteMultipartUpload { days_after_initiation: Some(0), }), del_marker_expiration: None, filter: None, id: Some("enabled-abort".to_string()), noncurrent_version_expiration: None, noncurrent_version_transitions: None, prefix: None, transitions: None, }], }; assert!(lifecycle_has_abort_multipart_rules(&config)); assert!(!lifecycle_has_expiry_rules(&config)); } #[tokio::test] async fn execute_list_buckets_returns_internal_error_when_store_uninitialized() { let input = ListBucketsInput::builder().build().unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_buckets(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_list_object_versions_returns_internal_error_when_store_uninitialized() { let input = ListObjectVersionsInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_object_versions(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_list_object_versions_m_returns_internal_error_when_store_uninitialized() { let input = ListObjectVersionsInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_object_versions_m(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[test] fn build_list_object_versions_metadata_output_maps_metadata_and_preserves_entry_order() { use time::macros::datetime; use uuid::Uuid; let object_infos = ListObjectVersionsInfo { is_truncated: true, next_marker: Some("obj-z".to_string()), next_version_idmarker: Some("null".to_string()), prefixes: vec!["logs/".to_string()], objects: vec![ ObjectInfo { bucket: "demo-bucket".to_string(), name: "obj-a".to_string(), mod_time: Some(datetime!(2025-01-01 00:00 UTC)), size: 11, user_defined: Arc::new(HashMap::from([("project".to_string(), "alpha".to_string())])), parity_blocks: 2, data_blocks: 4, version_id: Some(Uuid::nil()), user_tags: Arc::new("env=prod".to_string()), is_latest: true, etag: Some("0123456789abcdef0123456789abcdef".to_string()), ..Default::default() }, ObjectInfo { bucket: "demo-bucket".to_string(), name: "obj-b".to_string(), mod_time: Some(datetime!(2025-01-02 00:00 UTC)), delete_marker: true, user_defined: Arc::new(HashMap::from([("marker".to_string(), "true".to_string())])), version_id: None, ..Default::default() }, ], }; let permissions = HashMap::from([ ( "obj-a".to_string(), ObjectMetadataPermissions { metadata_allowed: true, tags_allowed: true, }, ), ( "obj-b".to_string(), ObjectMetadataPermissions { metadata_allowed: true, tags_allowed: false, }, ), ]); let params = ListObjectVersionsParams { prefix: "pre".to_string(), delimiter: Some("/".to_string()), key_marker: Some("start marker".to_string()), version_id_marker: Some("vid-1".to_string()), max_keys: 1000, }; let output = build_list_object_versions_metadata_output( object_infos, "demo-bucket", ¶ms, Some(&EncodingType::from_static(EncodingType::URL)), &permissions, ); assert_eq!(output.output.name.as_deref(), Some("demo-bucket")); assert_eq!(output.output.prefix.as_deref(), Some("pre")); assert_eq!(output.output.key_marker.as_deref(), Some("start%20marker")); assert_eq!(output.output.next_key_marker.as_deref(), Some("obj-z")); assert_eq!(output.output.next_version_id_marker.as_deref(), Some("null")); assert_eq!(output.entries.len(), 2); match &output.entries[0] { ListObjectVersionMetadataEntry::Version(version, extension) => { assert_eq!(version.key.as_deref(), Some("obj-a")); assert_eq!(version.version_id.as_deref(), Some(Uuid::nil().to_string().as_str())); assert_eq!(extension.user_tags.as_deref(), Some("env=prod")); assert_eq!(extension.internal, Some(ObjectInternalInfo { k: 4, m: 2 })); assert_eq!( extension.user_metadata.clone(), Some(vec![MetadataEntry { name: Some("project".to_string()), value: Some("alpha".to_string()), }]) ); } other => panic!("expected version entry, got {other:?}"), } match &output.entries[1] { ListObjectVersionMetadataEntry::DeleteMarker(marker, extension) => { assert_eq!(marker.key.as_deref(), Some("obj-b")); assert_eq!(marker.version_id.as_deref(), Some("null")); assert!(extension.user_tags.is_none()); assert_eq!( extension.user_metadata.clone(), Some(vec![MetadataEntry { name: Some("marker".to_string()), value: Some("true".to_string()), }]) ); } other => panic!("expected delete marker entry, got {other:?}"), } } #[test] fn build_list_object_versions_metadata_output_uses_params_and_hides_metadata_without_permissions() { use time::macros::datetime; let object_infos = ListObjectVersionsInfo { is_truncated: false, next_marker: Some(String::new()), next_version_idmarker: Some(String::new()), prefixes: vec!["logs and more/".to_string()], objects: vec![ObjectInfo { bucket: "demo-bucket".to_string(), name: "logs and more/object one.txt".to_string(), mod_time: Some(datetime!(2025-01-04 00:00 UTC)), size: 7, user_defined: Arc::new(HashMap::from([("secret".to_string(), "value".to_string())])), user_tags: Arc::new("env=prod".to_string()), parity_blocks: 1, data_blocks: 2, ..Default::default() }], }; let params = ListObjectVersionsParams { prefix: "logs and more/".to_string(), delimiter: Some(" ".to_string()), key_marker: Some("marker value".to_string()), version_id_marker: None, max_keys: 25, }; let output = build_list_object_versions_metadata_output( object_infos, "demo-bucket", ¶ms, Some(&EncodingType::from_static(EncodingType::URL)), &HashMap::new(), ); assert_eq!(output.output.name.as_deref(), Some("demo-bucket")); assert_eq!(output.output.prefix.as_deref(), Some("logs%20and%20more%2F")); assert_eq!(output.output.delimiter.as_deref(), Some("%20")); assert_eq!(output.output.key_marker.as_deref(), Some("marker%20value")); assert_eq!(output.output.version_id_marker.as_deref(), Some("")); assert_eq!(output.output.next_key_marker, None); assert_eq!(output.output.next_version_id_marker, None); match &output.entries[0] { ListObjectVersionMetadataEntry::Version(version, extension) => { assert_eq!(version.key.as_deref(), Some("logs%20and%20more%2Fobject%20one.txt")); assert!(extension.user_metadata.is_none()); assert!(extension.user_tags.is_none()); assert!(extension.internal.is_none()); } other => panic!("expected version entry, got {other:?}"), } } #[tokio::test] async fn execute_list_objects_returns_internal_error_when_store_uninitialized() { let input = ListObjectsInput::builder().bucket("test-bucket".to_string()).build().unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_objects(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_list_objects_v2m_returns_internal_error_when_store_uninitialized() { let input = ListObjectsV2Input::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_objects_v2m(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[test] fn build_list_objects_v2_metadata_output_maps_metadata_and_key_count() { use time::macros::datetime; let object_infos = ListObjectsV2Info { is_truncated: true, next_continuation_token: Some("next-token".to_string()), objects: vec![ObjectInfo { bucket: "demo-bucket".to_string(), name: "logs/obj a.txt".to_string(), mod_time: Some(datetime!(2025-01-03 00:00 UTC)), size: 11, user_defined: Arc::new(HashMap::from([("project".to_string(), "alpha".to_string())])), parity_blocks: 2, data_blocks: 4, user_tags: Arc::new("env=prod".to_string()), etag: Some("0123456789abcdef0123456789abcdef".to_string()), ..Default::default() }], prefixes: vec!["logs/archive/".to_string()], ..Default::default() }; let permissions = HashMap::from([( "logs/obj a.txt".to_string(), ObjectMetadataPermissions { metadata_allowed: true, tags_allowed: true, }, )]); let params = ListObjectsV2Params { prefix: "logs/".to_string(), max_keys: 1000, delimiter: Some("/".to_string()), response_start_after: Some("logs/start after".to_string()), start_after_for_query: None, response_continuation_token: Some("start token".to_string()), decoded_continuation_token: None, }; let output = build_list_objects_v2_metadata_output( object_infos, "demo-bucket", ¶ms, Some(&EncodingType::from_static(EncodingType::URL)), true, &permissions, ); assert_eq!(output.output.name.as_deref(), Some("demo-bucket")); assert_eq!(output.output.prefix.as_deref(), Some("logs/")); assert_eq!(output.output.continuation_token.as_deref(), Some("start token")); assert_eq!(output.output.start_after.as_deref(), Some("logs/start after")); assert_eq!(output.output.next_continuation_token.as_deref(), Some("bmV4dC10b2tlbg==")); assert_eq!(output.output.key_count, Some(2)); assert_eq!(output.contents.len(), 1); assert_eq!(output.output.common_prefixes.as_ref().map(Vec::len), Some(1)); let entry = output.contents.first().unwrap(); assert_eq!(entry.object.key.as_deref(), Some("logs/obj%20a.txt")); assert_eq!(entry.extension.user_tags.as_deref(), Some("env=prod")); assert_eq!(entry.extension.internal, Some(ObjectInternalInfo { k: 4, m: 2 })); assert!(entry.object.owner.is_some()); assert_eq!( entry.extension.user_metadata.clone(), Some(vec![MetadataEntry { name: Some("project".to_string()), value: Some("alpha".to_string()), }]) ); let prefix = output.output.common_prefixes.as_ref().unwrap().first().unwrap(); assert_eq!(prefix.prefix.as_deref(), Some("logs/archive/")); } #[test] fn list_objects_v2_metadata_output_serializes_minio_extension_xml() { use time::macros::datetime; let object_infos = ListObjectsV2Info { is_truncated: false, objects: vec![ObjectInfo { bucket: "demo-bucket".to_string(), name: "logs/obj-a.txt".to_string(), mod_time: Some(datetime!(2025-01-03 00:00 UTC)), size: 11, user_defined: Arc::new(HashMap::from([("project".to_string(), "alpha".to_string())])), parity_blocks: 2, data_blocks: 4, user_tags: Arc::new("env=prod&project=alpha".to_string()), ..Default::default() }], ..Default::default() }; let permissions = HashMap::from([( "logs/obj-a.txt".to_string(), ObjectMetadataPermissions { metadata_allowed: true, tags_allowed: true, }, )]); let params = ListObjectsV2Params { prefix: "logs/".to_string(), max_keys: 1000, delimiter: None, response_start_after: None, start_after_for_query: None, response_continuation_token: None, decoded_continuation_token: None, }; let output = build_list_objects_v2_metadata_output(object_infos, "demo-bucket", ¶ms, None, false, &permissions); let xml = String::from_utf8(serialize_config(&output).expect("metadata output should serialize")) .expect("metadata output should be UTF-8"); assert!(xml.contains("")); assert!(xml.contains("alpha")); assert!(xml.contains("env=prod&project=alpha")); assert!(xml.contains("42")); assert!(!xml.contains("")); } #[test] fn is_valid_xml_name_rejects_non_name_keys() { // Regression for #2743: keys that are illegal as XML element names. for good in ["project", "_x", "a-b.c", "Content_Type", "ns:key", "键名"] { assert!(is_valid_xml_name(good), "{good:?} should be a valid XML name"); } for bad in ["", "bad key", "1abc", "a$b", "a%b", "a#b", "a*b", "a|b", "a\u{1}b", "-lead"] { assert!(!is_valid_xml_name(bad), "{bad:?} should be rejected as an XML name"); } } #[test] fn sanitize_xml_text_strips_illegal_control_chars() { // Regression for #2743: XML 1.0 forbids C0 controls other than tab/newline/CR. assert_eq!(sanitize_xml_text("clean value").as_ref(), "clean value"); assert_eq!(sanitize_xml_text("keep\ttab\nnl\rcr").as_ref(), "keep\ttab\nnl\rcr"); assert_eq!(sanitize_xml_text("bad\u{1}\u{7}\u{1b}chars").as_ref(), "badchars"); } #[test] fn list_objects_v2_metadata_output_survives_poison_metadata_key() { // Regression for #2743: a single object carrying an XML-unsafe user-metadata key // (space in the key) or an illegal control char in the value must not corrupt the // whole listing document. The offending key is dropped; well-formed siblings remain. use time::macros::datetime; let object_infos = ListObjectsV2Info { is_truncated: false, objects: vec![ObjectInfo { bucket: "demo-bucket".to_string(), name: "logs/obj-a.txt".to_string(), mod_time: Some(datetime!(2025-01-03 00:00 UTC)), size: 11, user_defined: Arc::new(HashMap::from([ ("project".to_string(), "alpha".to_string()), ("bad key".to_string(), "should-be-dropped".to_string()), ("note".to_string(), "line1\u{1}line2".to_string()), ])), parity_blocks: 2, data_blocks: 4, ..Default::default() }], ..Default::default() }; let permissions = HashMap::from([( "logs/obj-a.txt".to_string(), ObjectMetadataPermissions { metadata_allowed: true, tags_allowed: true, }, )]); let params = ListObjectsV2Params { prefix: "logs/".to_string(), max_keys: 1000, delimiter: None, response_start_after: None, start_after_for_query: None, response_continuation_token: None, decoded_continuation_token: None, }; let output = build_list_objects_v2_metadata_output(object_infos, "demo-bucket", ¶ms, None, false, &permissions); let xml = String::from_utf8(serialize_config(&output).expect("metadata output should serialize")) .expect("metadata output should be UTF-8"); // Good sibling metadata survives. assert!(xml.contains("alpha"), "well-formed key must remain: {xml}"); // The XML-unsafe key is dropped entirely (never emitted as a malformed tag). assert!(!xml.contains("bad key"), "poison key must not appear: {xml}"); // The control char is stripped from the value. assert!(xml.contains("line1line2"), "control char must be stripped: {xml}"); assert!(!xml.contains('\u{1}'), "no raw control chars in output"); // The document is well-formed and re-parses cleanly. let mut reader = quick_xml::Reader::from_str(&xml); loop { match reader.read_event() { Ok(quick_xml::events::Event::Eof) => break, Ok(_) => {} Err(err) => panic!("serialized listing must be well-formed XML, got {err} in {xml}"), } } } #[test] fn build_list_objects_v2_metadata_output_uses_params_and_hides_owner_without_fetch_owner() { use time::macros::datetime; let object_infos = ListObjectsV2Info { is_truncated: false, next_continuation_token: None, objects: vec![ObjectInfo { bucket: "demo-bucket".to_string(), name: "logs and more/object one.txt".to_string(), mod_time: Some(datetime!(2025-01-05 00:00 UTC)), size: 13, user_defined: Arc::new(HashMap::from([("secret".to_string(), "value".to_string())])), user_tags: Arc::new("env=prod".to_string()), parity_blocks: 1, data_blocks: 2, ..Default::default() }], prefixes: vec!["logs and more/archive/".to_string()], ..Default::default() }; let params = ListObjectsV2Params { prefix: "logs and more/".to_string(), max_keys: 25, delimiter: Some("/".to_string()), response_start_after: Some("logs and more/start after".to_string()), start_after_for_query: Some("decoded start after".to_string()), response_continuation_token: Some("opaque token".to_string()), decoded_continuation_token: Some("decoded token".to_string()), }; let output = build_list_objects_v2_metadata_output( object_infos, "demo-bucket", ¶ms, Some(&EncodingType::from_static(EncodingType::URL)), false, &HashMap::new(), ); assert_eq!(output.output.name.as_deref(), Some("demo-bucket")); assert_eq!(output.output.prefix.as_deref(), Some("logs and more/")); assert_eq!(output.output.delimiter.as_deref(), Some("/")); assert_eq!(output.output.continuation_token.as_deref(), Some("opaque token")); assert_eq!(output.output.start_after.as_deref(), Some("logs and more/start after")); assert_eq!(output.output.key_count, Some(2)); assert_eq!(output.output.encoding_type.as_ref().map(EncodingType::as_str), Some(EncodingType::URL)); let entry = output.contents.first().unwrap(); assert_eq!(entry.object.key.as_deref(), Some("logs%20and%20more/object%20one.txt")); assert!(entry.object.owner.is_none()); assert!(entry.extension.user_metadata.is_none()); assert!(entry.extension.user_tags.is_none()); assert!(entry.extension.internal.is_none()); let prefix = output.output.common_prefixes.as_ref().unwrap().first().unwrap(); assert_eq!(prefix.prefix.as_deref(), Some("logs%20and%20more/archive/")); } #[tokio::test] async fn execute_put_bucket_lifecycle_configuration_rejects_missing_configuration() { let input = PutBucketLifecycleConfigurationInput::builder() .bucket("test-bucket".to_string()) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_put_bucket_lifecycle_configuration(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[test] fn validate_notification_configuration_filters_rejects_invalid_filter_name() { let raw_name = "unsupported".repeat(100); let cfg = NotificationConfiguration { queue_configurations: Some(vec![QueueConfiguration { id: Some("q1".to_string()), queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(), events: vec!["s3:ObjectCreated:*".to_string().into()], filter: Some(NotificationConfigurationFilter { key: Some(S3KeyFilter { filter_rules: Some(vec![FilterRule { name: Some(FilterRuleName::from(raw_name.clone())), value: Some("uploads/".to_string()), }]), }), }), }]), ..Default::default() }; let err = validate_notification_configuration_filters(&cfg).unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); let msg = err.message().unwrap_or_default(); assert!(msg.contains("len="), "error message should include summarized length"); assert!(!msg.contains(&raw_name), "error message should not echo full raw filter name"); } #[test] fn validate_notification_configuration_filters_accepts_case_insensitive_filter_names() { let cfg = NotificationConfiguration { queue_configurations: Some(vec![QueueConfiguration { id: Some("q1".to_string()), queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(), events: vec!["s3:ObjectCreated:*".to_string().into()], filter: Some(NotificationConfigurationFilter { key: Some(S3KeyFilter { filter_rules: Some(vec![ FilterRule { name: Some(FilterRuleName::from("Prefix".to_string())), value: Some("uploads/".to_string()), }, FilterRule { name: Some(FilterRuleName::from("Suffix".to_string())), value: Some(".csv".to_string()), }, ]), }), }), }]), ..Default::default() }; validate_notification_configuration_filters(&cfg).expect("capitalized filter names should be accepted"); } #[test] fn validate_notification_configuration_filters_rejects_duplicate_prefix_rules() { let cfg = NotificationConfiguration { queue_configurations: Some(vec![QueueConfiguration { id: Some("q1".to_string()), queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(), events: vec!["s3:ObjectCreated:*".to_string().into()], filter: Some(NotificationConfigurationFilter { key: Some(S3KeyFilter { filter_rules: Some(vec![ FilterRule { name: Some(FilterRuleName::from_static(FilterRuleName::PREFIX)), value: Some("uploads/".to_string()), }, FilterRule { name: Some(FilterRuleName::from_static(FilterRuleName::PREFIX)), value: Some("images/".to_string()), }, ]), }), }), }]), ..Default::default() }; let err = validate_notification_configuration_filters(&cfg).unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[test] fn validate_notification_configuration_filters_rejects_invalid_filter_value() { let cfg = NotificationConfiguration { queue_configurations: Some(vec![QueueConfiguration { id: Some("q1".to_string()), queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(), events: vec!["s3:ObjectCreated:*".to_string().into()], filter: Some(NotificationConfigurationFilter { key: Some(S3KeyFilter { filter_rules: Some(vec![FilterRule { name: Some(FilterRuleName::from_static(FilterRuleName::SUFFIX)), value: Some("../secret".to_string()), }]), }), }), }]), ..Default::default() }; let err = validate_notification_configuration_filters(&cfg).unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); let msg = err.message().unwrap_or_default(); assert!(msg.contains("len="), "error message should include summarized length"); assert!(!msg.contains("../secret"), "error message should not echo full raw filter value"); } #[test] fn validate_notification_configuration_filters_rejects_missing_filter_name() { let cfg = NotificationConfiguration { queue_configurations: Some(vec![QueueConfiguration { id: Some("q1".to_string()), queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(), events: vec!["s3:ObjectCreated:*".to_string().into()], filter: Some(NotificationConfigurationFilter { key: Some(S3KeyFilter { filter_rules: Some(vec![FilterRule { name: None, value: Some("uploads/".to_string()), }]), }), }), }]), ..Default::default() }; let err = validate_notification_configuration_filters(&cfg).unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[test] fn validate_notification_configuration_filters_rejects_missing_filter_value() { let cfg = NotificationConfiguration { queue_configurations: Some(vec![QueueConfiguration { id: Some("q1".to_string()), queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(), events: vec!["s3:ObjectCreated:*".to_string().into()], filter: Some(NotificationConfigurationFilter { key: Some(S3KeyFilter { filter_rules: Some(vec![FilterRule { name: Some(FilterRuleName::from_static(FilterRuleName::PREFIX)), value: None, }]), }), }), }]), ..Default::default() }; let err = validate_notification_configuration_filters(&cfg).unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[test] fn validate_notification_configuration_filters_rejects_duplicate_suffix_rules() { let cfg = NotificationConfiguration { queue_configurations: Some(vec![QueueConfiguration { id: Some("q1".to_string()), queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(), events: vec!["s3:ObjectCreated:*".to_string().into()], filter: Some(NotificationConfigurationFilter { key: Some(S3KeyFilter { filter_rules: Some(vec![ FilterRule { name: Some(FilterRuleName::from_static(FilterRuleName::SUFFIX)), value: Some(".csv".to_string()), }, FilterRule { name: Some(FilterRuleName::from_static(FilterRuleName::SUFFIX)), value: Some(".log".to_string()), }, ]), }), }), }]), ..Default::default() }; let err = validate_notification_configuration_filters(&cfg).unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[tokio::test] async fn execute_put_bucket_policy_returns_internal_error_when_store_uninitialized() { let input = PutBucketPolicyInput::builder() .bucket("test-bucket".to_string()) .policy("{}".to_string()) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let before = s3_op_total(S3Operation::PutBucketPolicy); let err = usecase.execute_put_bucket_policy(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); assert_eq!(s3_op_total(S3Operation::PutBucketPolicy), before + 1); } #[tokio::test] async fn execute_put_bucket_notification_configuration_rejects_invalid_filter_before_store_lookup() { let input = PutBucketNotificationConfigurationInput::builder() .bucket("test-bucket".to_string()) .notification_configuration(NotificationConfiguration { queue_configurations: Some(vec![QueueConfiguration { id: Some("q1".to_string()), queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(), events: vec!["s3:ObjectCreated:*".to_string().into()], filter: Some(NotificationConfigurationFilter { key: Some(S3KeyFilter { filter_rules: Some(vec![FilterRule { name: Some(FilterRuleName::from("unsupported".to_string())), value: Some("uploads/".to_string()), }]), }), }), }]), ..Default::default() }) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_put_bucket_notification_configuration(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[tokio::test] async fn execute_put_bucket_cors_returns_internal_error_when_store_uninitialized() { let input = PutBucketCorsInput::builder() .bucket("test-bucket".to_string()) .cors_configuration(CORSConfiguration::default()) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_put_bucket_cors(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_put_bucket_replication_returns_internal_error_when_store_uninitialized() { // The config must clear the structural/capability validators so the // request actually reaches the store lookup this test pins. let input = PutBucketReplicationInput::builder() .bucket("test-bucket".to_string()) .replication_configuration(ReplicationConfiguration { role: "arn:aws:iam::123456789012:role/test".to_string(), rules: vec![replication_rule_for_target("arn:rustfs:replication:us-east-1:target:bucket")], }) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_put_bucket_replication(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_put_bucket_replication_rejects_unsupported_fields_before_store_or_metadata_write() { let mut rule = replication_rule_for_target("arn:rustfs:replication:us-east-1:target:bucket"); rule.destination.account = Some("123456789012".to_string()); let input = PutBucketReplicationInput::builder() .bucket("test-bucket".to_string()) .replication_configuration(ReplicationConfiguration { role: String::new(), rules: vec![rule], }) .build() .unwrap(); let err = DefaultBucketUsecase::without_context() .execute_put_bucket_replication(build_request(input, Method::PUT)) .await .expect_err("unsupported fields must be rejected before store access"); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); assert_eq!( err.message(), Some("replication field Destination.Account is not supported by this RustFS version") ); } #[tokio::test] async fn execute_put_bucket_encryption_returns_internal_error_when_store_uninitialized() { let input = PutBucketEncryptionInput::builder() .bucket("test-bucket".to_string()) .server_side_encryption_configuration(ServerSideEncryptionConfiguration::default()) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_put_bucket_encryption(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_put_bucket_tagging_returns_internal_error_when_store_uninitialized() { let input = PutBucketTaggingInput::builder() .bucket("test-bucket".to_string()) .tagging(Tagging { tag_set: vec![Tag { key: Some("env".to_string()), value: Some("prod".to_string()), }], }) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_put_bucket_tagging(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_put_public_access_block_returns_internal_error_when_store_uninitialized() { let input = PutPublicAccessBlockInput::builder() .bucket("test-bucket".to_string()) .public_access_block_configuration(PublicAccessBlockConfiguration::default()) .build() .unwrap(); let req = build_request(input, Method::PUT); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_put_public_access_block(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } #[tokio::test] async fn execute_list_objects_v2_rejects_negative_max_keys() { let input = ListObjectsV2Input::builder() .bucket("test-bucket".to_string()) .max_keys(Some(-1)) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_objects_v2(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[tokio::test] async fn execute_list_objects_v2_rejects_invalid_continuation_token_before_store_lookup() { let input = ListObjectsV2Input::builder() .bucket("test-bucket".to_string()) .continuation_token(Some("%%%".to_string())) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_objects_v2(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); assert_eq!(err.message(), Some("Invalid continuation token")); } #[tokio::test] async fn execute_list_objects_v2m_rejects_negative_max_keys() { let input = ListObjectsV2Input::builder() .bucket("test-bucket".to_string()) .max_keys(Some(-1)) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_objects_v2m(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[tokio::test] async fn execute_list_objects_v2m_rejects_invalid_continuation_token_before_store_lookup() { let input = ListObjectsV2Input::builder() .bucket("test-bucket".to_string()) .continuation_token(Some("%%%".to_string())) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_objects_v2m(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); assert_eq!(err.message(), Some("Invalid continuation token")); } #[tokio::test] async fn execute_list_object_versions_rejects_negative_max_keys_before_store_lookup() { let input = ListObjectVersionsInput::builder() .bucket("test-bucket".to_string()) .max_keys(Some(-1)) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_object_versions(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } #[tokio::test] async fn execute_list_object_versions_m_rejects_negative_max_keys_before_store_lookup() { let input = ListObjectVersionsInput::builder() .bucket("test-bucket".to_string()) .max_keys(Some(-1)) .build() .unwrap(); let req = build_request(input, Method::GET); let usecase = DefaultBucketUsecase::without_context(); let err = usecase.execute_list_object_versions_m(req).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } }