From 588631b02a630763af9378d4d1f15873e45cd3e9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=AE=89=E6=AD=A3=E8=B6=85?= Date: Mon, 23 Feb 2026 21:19:57 +0800 Subject: [PATCH] refactor(app): migrate object ACL/tagging flows to usecase (#1921) --- rustfs/src/app/object_usecase.rs | 412 ++++++++++++++++++++++++++++++- rustfs/src/storage/ecfs.rs | 346 ++------------------------ 2 files changed, 436 insertions(+), 322 deletions(-) diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 4eb51fffe..b7ce10f1b 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -37,6 +37,7 @@ use datafusion::arrow::{ }; use futures::StreamExt; use http::{HeaderMap, HeaderValue, StatusCode}; +use metrics::{counter, histogram}; use rustfs_ecstore::StorageAPI; use rustfs_ecstore::bucket::quota::checker::QuotaChecker; use rustfs_ecstore::bucket::{ @@ -51,7 +52,7 @@ use rustfs_ecstore::bucket::{ DeletedObjectReplicationInfo, get_must_replicate_options, must_replicate, schedule_replication, schedule_replication_delete, }, - tagging::decode_tags, + tagging::{decode_tags, encode_tags}, versioning_sys::BucketVersioningSys, }; use rustfs_ecstore::client::object_api_utils::to_s3s_etag; @@ -633,6 +634,187 @@ impl DefaultObjectUsecase { result } + pub async fn execute_put_object_acl(&self, req: S3Request) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let PutObjectAclInput { + bucket, + key, + acl, + access_control_policy, + grant_full_control, + grant_read, + grant_read_acp, + grant_write, + grant_write_acp, + version_id, + .. + } = req.input; + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers) + .await + .map_err(ApiError::from)?; + let info = store.get_object_info(&bucket, &key, &opts).await.map_err(ApiError::from)?; + + let bucket_owner = default_owner(); + let existing_owner = info + .user_defined + .get(INTERNAL_ACL_METADATA_KEY) + .and_then(|acl| serde_json::from_str::(acl).ok()) + .map(|acl| acl.owner) + .unwrap_or_else(|| bucket_owner.clone()); + + let mut stored_acl = access_control_policy + .as_ref() + .map(|policy| stored_acl_from_policy(policy, &existing_owner)) + .transpose()?; + + if stored_acl.is_none() { + stored_acl = stored_acl_from_grant_headers( + &existing_owner, + grant_read.map(|v| v.to_string()), + grant_write.map(|v| v.to_string()), + grant_read_acp.map(|v| v.to_string()), + grant_write_acp.map(|v| v.to_string()), + grant_full_control.map(|v| v.to_string()), + )?; + } + + if stored_acl.is_none() + && let Some(canned) = acl + { + stored_acl = Some(stored_acl_from_canned_object(canned.as_str(), &bucket_owner, &existing_owner)); + } + + let stored_acl = + stored_acl.unwrap_or_else(|| stored_acl_from_canned_object(ObjectCannedACL::PRIVATE, &bucket_owner, &existing_owner)); + + if let Ok((config, _)) = metadata_sys::get_public_access_block_config(&bucket).await + && config.block_public_acls.unwrap_or(false) + && stored_acl.grants.iter().any(is_public_grant) + { + return Err(s3_error!(AccessDenied, "Access Denied")); + } + + let acl_data = serialize_acl(&stored_acl)?; + let mut eval_metadata = HashMap::new(); + eval_metadata.insert(INTERNAL_ACL_METADATA_KEY.to_string(), String::from_utf8_lossy(&acl_data).to_string()); + + let popts = ObjectOptions { + mod_time: info.mod_time, + version_id: opts.version_id, + eval_metadata: Some(eval_metadata), + ..Default::default() + }; + + store.put_object_metadata(&bucket, &key, &popts).await.map_err(|e| { + error!("put_object_metadata failed, {}", e.to_string()); + s3_error!(InternalError, "{}", e.to_string()) + })?; + + Ok(S3Response::new(PutObjectAclOutput::default())) + } + + #[instrument(level = "debug", skip(self, req))] + pub async fn execute_put_object_tagging( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let start_time = std::time::Instant::now(); + let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutTagging, "s3:PutObjectTagging"); + let PutObjectTaggingInput { + bucket, + key: object, + tagging, + .. + } = req.input.clone(); + + if tagging.tag_set.len() > 10 { + error!("Tag set exceeds maximum of 10 tags: {}", tagging.tag_set.len()); + return Err(s3_error!(InvalidTag, "Cannot have more than 10 tags per object")); + } + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let mut tag_keys = std::collections::HashSet::with_capacity(tagging.tag_set.len()); + for tag in &tagging.tag_set { + let key = tag.key.as_ref().filter(|k| !k.is_empty()).ok_or_else(|| { + error!("Empty tag key"); + s3_error!(InvalidTag, "Tag key cannot be empty") + })?; + + if key.len() > 128 { + error!("Tag key too long: {} bytes", key.len()); + return Err(s3_error!(InvalidTag, "Tag key is too long, maximum allowed length is 128 characters")); + } + + let value = tag.value.as_ref().ok_or_else(|| { + error!("Null tag value"); + s3_error!(InvalidTag, "Tag value cannot be null") + })?; + + if value.len() > 256 { + error!("Tag value too long: {} bytes", value.len()); + return Err(s3_error!(InvalidTag, "Tag value is too long, maximum allowed length is 256 characters")); + } + + if !tag_keys.insert(key) { + error!("Duplicate tag key: {}", key); + return Err(s3_error!(InvalidTag, "Cannot provide multiple Tags with the same key")); + } + } + + let tags = encode_tags(tagging.tag_set); + debug!("Encoded tags: {}", tags); + + let version_id = req.input.version_id.clone(); + let opts = ObjectOptions { + version_id: parse_object_version_id(version_id)?.map(Into::into), + ..Default::default() + }; + + store.put_object_tags(&bucket, &object, &tags, &opts).await.map_err(|e| { + error!("Failed to put object tags: {}", e); + counter!("rustfs.put_object_tagging.failure").increment(1); + ApiError::from(e) + })?; + + let manager = get_concurrency_manager(); + let version_id = req.input.version_id.clone(); + let cache_key = ConcurrencyManager::make_cache_key(&bucket, &object, version_id.clone().as_deref()); + tokio::spawn(async move { + manager + .invalidate_cache_versioned(&bucket, &object, version_id.as_deref()) + .await; + debug!("Cache invalidated for tagged object: {}", cache_key); + }); + + counter!("rustfs.put_object_tagging.success").increment(1); + + let version_id_resp = req.input.version_id.clone().unwrap_or_default(); + helper = helper.version_id(version_id_resp); + + let result = Ok(S3Response::new(PutObjectTaggingOutput { + version_id: req.input.version_id.clone(), + })); + let _ = helper.complete(&result); + let duration = start_time.elapsed(); + histogram!("rustfs.object_tagging.operation.duration.seconds", "operation" => "put").record(duration.as_secs_f64()); + result + } + #[instrument( level = "debug", skip(self, req), @@ -1326,6 +1508,95 @@ impl DefaultObjectUsecase { result } + pub async fn execute_get_object_acl(&self, req: S3Request) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let GetObjectAclInput { + bucket, key, version_id, .. + } = req.input; + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers) + .await + .map_err(ApiError::from)?; + let info = store.get_object_info(&bucket, &key, &opts).await.map_err(ApiError::from)?; + + let bucket_owner = default_owner(); + let object_owner = info + .user_defined + .get(INTERNAL_ACL_METADATA_KEY) + .and_then(|acl| serde_json::from_str::(acl).ok()) + .map(|acl| acl.owner) + .unwrap_or_else(default_owner); + + let stored_acl = info + .user_defined + .get(INTERNAL_ACL_METADATA_KEY) + .map(|acl| parse_acl_json_or_canned_object(acl, &bucket_owner, &object_owner)) + .unwrap_or_else(|| stored_acl_from_canned_object(ObjectCannedACL::PRIVATE, &bucket_owner, &object_owner)); + + let mut sorted_grants = stored_acl.grants.clone(); + sorted_grants.sort_by_key(|grant| grant.grantee.grantee_type != "Group"); + let grants = sorted_grants.iter().map(stored_grant_to_dto).collect(); + + Ok(S3Response::new(GetObjectAclOutput { + grants: Some(grants), + owner: Some(stored_owner_to_dto(&stored_acl.owner)), + ..Default::default() + })) + } + + #[instrument(level = "debug", skip(self, req))] + pub async fn execute_get_object_tagging( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let start_time = std::time::Instant::now(); + let GetObjectTaggingInput { bucket, key: object, .. } = req.input; + + info!("Starting get_object_tagging for bucket: {}, object: {}", bucket, object); + + let Some(store) = new_object_layer_fn() else { + error!("Store not initialized"); + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let version_id = req.input.version_id.clone(); + let opts = ObjectOptions { + version_id: parse_object_version_id(version_id)?.map(Into::into), + ..Default::default() + }; + + let tags = store.get_object_tags(&bucket, &object, &opts).await.map_err(|e| { + if is_err_object_not_found(&e) { + error!("Object not found: {}", e); + return s3_error!(NoSuchKey); + } + error!("Failed to get object tags: {}", e); + ApiError::from(e).into() + })?; + + let tag_set = decode_tags(tags.as_str()); + debug!("Decoded tag set: {:?}", tag_set); + + counter!("rustfs.get_object_tagging.success").increment(1); + let duration = start_time.elapsed(); + histogram!("rustfs.object_tagging.operation.duration.seconds", "operation" => "get").record(duration.as_secs_f64()); + Ok(S3Response::new(GetObjectTaggingOutput { + tag_set, + version_id: req.input.version_id.clone(), + })) + } + #[instrument(level = "debug", skip(self, req))] pub async fn execute_copy_object(&self, req: S3Request) -> S3Result> { if let Some(context) = &self.context { @@ -1849,6 +2120,64 @@ impl DefaultObjectUsecase { result } + #[instrument(level = "debug", skip(self, req))] + pub async fn execute_delete_object_tagging( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let start_time = std::time::Instant::now(); + let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedDeleteTagging, "s3:DeleteObjectTagging"); + let DeleteObjectTaggingInput { + bucket, + key: object, + version_id, + .. + } = req.input.clone(); + + let Some(store) = new_object_layer_fn() else { + error!("Store not initialized"); + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let version_id_for_parse = version_id.clone(); + let opts = ObjectOptions { + version_id: parse_object_version_id(version_id_for_parse)?.map(Into::into), + ..Default::default() + }; + + store.delete_object_tags(&bucket, &object, &opts).await.map_err(|e| { + error!("Failed to delete object tags: {}", e); + ApiError::from(e) + })?; + + let manager = get_concurrency_manager(); + let version_id_clone = version_id.clone(); + tokio::spawn(async move { + manager + .invalidate_cache_versioned(&bucket, &object, version_id_clone.as_deref()) + .await; + debug!( + "Cache invalidated for deleted tagged object: bucket={}, object={}, version_id={:?}", + bucket, object, version_id_clone + ); + }); + + counter!("rustfs.delete_object_tagging.success").increment(1); + + let version_id_resp = version_id.clone().unwrap_or_default(); + helper = helper.version_id(version_id_resp); + + let result = Ok(S3Response::new(DeleteObjectTaggingOutput { version_id })); + let _ = helper.complete(&result); + let duration = start_time.elapsed(); + histogram!("rustfs.object_tagging.operation.duration.seconds", "operation" => "delete").record(duration.as_secs_f64()); + result + } + #[instrument(level = "debug", skip(self, req))] pub async fn execute_head_object(&self, req: S3Request) -> S3Result> { if let Some(context) = &self.context { @@ -2536,6 +2865,87 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } + #[tokio::test] + async fn execute_delete_object_tagging_returns_internal_error_when_store_uninitialized() { + let input = DeleteObjectTaggingInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::DELETE); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_delete_object_tagging(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[tokio::test] + async fn execute_get_object_acl_returns_internal_error_when_store_uninitialized() { + let input = GetObjectAclInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::GET); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_get_object_acl(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[tokio::test] + async fn execute_get_object_tagging_returns_internal_error_when_store_uninitialized() { + let input = GetObjectTaggingInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::GET); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_get_object_tagging(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[tokio::test] + async fn execute_put_object_acl_returns_internal_error_when_store_uninitialized() { + let input = PutObjectAclInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::PUT); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_put_object_acl(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[tokio::test] + async fn execute_put_object_tagging_returns_internal_error_when_store_uninitialized() { + let input = PutObjectTaggingInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .tagging(Tagging { + tag_set: vec![Tag { + key: Some("k".to_string()), + value: Some("v".to_string()), + }], + }) + .build() + .unwrap(); + + let req = build_request(input, Method::PUT); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_put_object_tagging(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + #[tokio::test] async fn execute_head_object_rejects_range_with_part_number() { let input = HeadObjectInput::builder() diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index e286082a0..14626e05d 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -16,7 +16,7 @@ use crate::app::bucket_usecase::DefaultBucketUsecase; use crate::app::multipart_usecase::DefaultMultipartUsecase; use crate::app::object_usecase::DefaultObjectUsecase; use crate::error::ApiError; -use crate::storage::concurrency::{ConcurrencyManager, get_concurrency_manager}; +use crate::storage::concurrency::get_concurrency_manager; use crate::storage::helper::OperationHelper; use crate::storage::options::get_content_sha256; use crate::storage::{ @@ -29,14 +29,12 @@ use crate::storage::{ }; use futures::StreamExt; use http::HeaderMap; -use metrics::{counter, histogram}; use rustfs_ecstore::{ bucket::{ metadata::{BUCKET_ACL_CONFIG, BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG}, metadata_sys, object_lock::objectlock_sys::{check_object_lock_for_deletion, check_retention_for_modification}, replication::{DeletedObjectReplicationInfo, check_replicate_delete, schedule_replication_delete}, - tagging::{decode_tags, encode_tags}, utils::serialize, versioning::VersioningApi, versioning_sys::BucketVersioningSys, @@ -82,7 +80,7 @@ use time::{OffsetDateTime, format_description::well_known::Rfc3339}; use tokio::io::{AsyncRead, AsyncSeek}; use tokio_tar::Archive; use tokio_util::io::StreamReader; -use tracing::{debug, error, info, instrument, warn}; +use tracing::{debug, error, instrument, warn}; use uuid::Uuid; macro_rules! try_ { @@ -468,7 +466,7 @@ fn grantee_from_grant(grantee: &Grantee) -> Result { } } -fn stored_acl_from_policy(policy: &AccessControlPolicy, owner_fallback: &StoredOwner) -> Result { +pub(crate) fn stored_acl_from_policy(policy: &AccessControlPolicy, owner_fallback: &StoredOwner) -> Result { let owner = policy .owner .as_ref() @@ -909,25 +907,6 @@ impl FS { result } - /// Auxiliary functions: parse version ID - /// - /// # Arguments - /// * `version_id` - An optional string representing the version ID to be parsed. - /// - /// # Returns - /// * `S3Result>` - A result containing an optional UUID if parsing is successful, or an S3 error if parsing fails. - fn parse_version_id(&self, version_id: Option) -> S3Result> { - if let Some(vid) = version_id { - let uuid = Uuid::parse_str(&vid).map_err(|e| { - error!("Invalid version ID: {}", e); - s3_error!(InvalidArgument, "Invalid version ID") - })?; - Ok(Some(uuid)) - } else { - Ok(None) - } - } - pub(crate) fn normalize_delete_objects_version_id( &self, version_id: Option, @@ -954,6 +933,18 @@ pub(crate) fn is_public_canned_acl(acl: &str) -> bool { matches!(acl, "public-read" | "public-read-write" | "authenticated-read") } +pub(crate) fn parse_object_version_id(version_id: Option) -> S3Result> { + if let Some(vid) = version_id { + let uuid = Uuid::parse_str(&vid).map_err(|e| { + error!("Invalid version ID: {}", e); + s3_error!(InvalidArgument, "Invalid version ID") + })?; + Ok(Some(uuid)) + } else { + Ok(None) + } +} + #[async_trait::async_trait] impl S3 for FS { #[instrument(level = "debug", skip(self))] @@ -1076,57 +1067,8 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let start_time = std::time::Instant::now(); - let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedDeleteTagging, "s3:DeleteObjectTagging"); - let DeleteObjectTaggingInput { - bucket, - key: object, - version_id, - .. - } = req.input.clone(); - - let Some(store) = new_object_layer_fn() else { - error!("Store not initialized"); - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - // Support versioned objects - let version_id_for_parse = version_id.clone(); - let opts = ObjectOptions { - version_id: self.parse_version_id(version_id_for_parse)?.map(Into::into), - ..Default::default() - }; - - // TODO: Replicate (keep the original TODO, if further replication logic is needed) - store.delete_object_tags(&bucket, &object, &opts).await.map_err(|e| { - error!("Failed to delete object tags: {}", e); - ApiError::from(e) - })?; - - // Invalidate cache for the deleted tagged object - let manager = get_concurrency_manager(); - let version_id_clone = version_id.clone(); - tokio::spawn(async move { - manager - .invalidate_cache_versioned(&bucket, &object, version_id_clone.as_deref()) - .await; - debug!( - "Cache invalidated for deleted tagged object: bucket={}, object={}, version_id={:?}", - bucket, object, version_id_clone - ); - }); - - // Add metrics - counter!("rustfs.delete_object_tagging.success").increment(1); - - let version_id_resp = version_id.clone().unwrap_or_default(); - helper = helper.version_id(version_id_resp); - - let result = Ok(S3Response::new(DeleteObjectTaggingOutput { version_id })); - let _ = helper.complete(&result); - let duration = start_time.elapsed(); - histogram!("rustfs.object_tagging.operation.duration.seconds", "operation" => "delete").record(duration.as_secs_f64()); - result + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_delete_object_tagging(req).await } /// Delete multiple objects @@ -1597,42 +1539,8 @@ impl S3 for FS { } async fn get_object_acl(&self, req: S3Request) -> S3Result> { - let GetObjectAclInput { - bucket, key, version_id, .. - } = req.input; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers) - .await - .map_err(ApiError::from)?; - let info = store.get_object_info(&bucket, &key, &opts).await.map_err(ApiError::from)?; - - let bucket_owner = default_owner(); - let object_owner = info - .user_defined - .get(INTERNAL_ACL_METADATA_KEY) - .and_then(|acl| serde_json::from_str::(acl).ok()) - .map(|acl| acl.owner) - .unwrap_or_else(default_owner); - - let stored_acl = info - .user_defined - .get(INTERNAL_ACL_METADATA_KEY) - .map(|acl| parse_acl_json_or_canned_object(acl, &bucket_owner, &object_owner)) - .unwrap_or_else(|| stored_acl_from_canned_object(ObjectCannedACL::PRIVATE, &bucket_owner, &object_owner)); - - let mut sorted_grants = stored_acl.grants.clone(); - sorted_grants.sort_by_key(|grant| grant.grantee.grantee_type != "Group"); - let grants = sorted_grants.iter().map(stored_grant_to_dto).collect(); - - Ok(S3Response::new(GetObjectAclOutput { - grants: Some(grants), - owner: Some(stored_owner_to_dto(&stored_acl.owner)), - ..Default::default() - })) + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_get_object_acl(req).await } async fn get_object_attributes( @@ -1808,43 +1716,8 @@ impl S3 for FS { #[instrument(level = "debug", skip(self))] async fn get_object_tagging(&self, req: S3Request) -> S3Result> { - let start_time = std::time::Instant::now(); - let GetObjectTaggingInput { bucket, key: object, .. } = req.input; - - info!("Starting get_object_tagging for bucket: {}, object: {}", bucket, object); - - let Some(store) = new_object_layer_fn() else { - error!("Store not initialized"); - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - // Support versioned objects - let version_id = req.input.version_id.clone(); - let opts = ObjectOptions { - version_id: self.parse_version_id(version_id)?.map(Into::into), - ..Default::default() - }; - - let tags = store.get_object_tags(&bucket, &object, &opts).await.map_err(|e| { - if is_err_object_not_found(&e) { - error!("Object not found: {}", e); - return s3_error!(NoSuchKey); - } - error!("Failed to get object tags: {}", e); - ApiError::from(e).into() - })?; - - let tag_set = decode_tags(tags.as_str()); - debug!("Decoded tag set: {:?}", tag_set); - - // Add metrics - counter!("rustfs.get_object_tagging.success").increment(1); - let duration = start_time.elapsed(); - histogram!("rustfs.object_tagging.operation.duration.seconds", "operation" => "put").record(duration.as_secs_f64()); - Ok(S3Response::new(GetObjectTaggingOutput { - tag_set, - version_id: req.input.version_id.clone(), - })) + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_get_object_tagging(req).await } #[instrument(level = "debug", skip(self, _req))] @@ -2363,86 +2236,8 @@ impl S3 for FS { } async fn put_object_acl(&self, req: S3Request) -> S3Result> { - let PutObjectAclInput { - bucket, - key, - acl, - access_control_policy, - grant_full_control, - grant_read, - grant_read_acp, - grant_write, - grant_write_acp, - version_id, - .. - } = req.input; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers) - .await - .map_err(ApiError::from)?; - let info = store.get_object_info(&bucket, &key, &opts).await.map_err(ApiError::from)?; - - let bucket_owner = default_owner(); - let existing_owner = info - .user_defined - .get(INTERNAL_ACL_METADATA_KEY) - .and_then(|acl| serde_json::from_str::(acl).ok()) - .map(|acl| acl.owner) - .unwrap_or_else(|| bucket_owner.clone()); - - let mut stored_acl = access_control_policy - .as_ref() - .map(|policy| stored_acl_from_policy(policy, &existing_owner)) - .transpose()?; - - if stored_acl.is_none() { - stored_acl = stored_acl_from_grant_headers( - &existing_owner, - grant_read.map(|v| v.to_string()), - grant_write.map(|v| v.to_string()), - grant_read_acp.map(|v| v.to_string()), - grant_write_acp.map(|v| v.to_string()), - grant_full_control.map(|v| v.to_string()), - )?; - } - - if stored_acl.is_none() - && let Some(canned) = acl - { - stored_acl = Some(stored_acl_from_canned_object(canned.as_str(), &bucket_owner, &existing_owner)); - } - - let stored_acl = - stored_acl.unwrap_or_else(|| stored_acl_from_canned_object(ObjectCannedACL::PRIVATE, &bucket_owner, &existing_owner)); - - if let Ok((config, _)) = metadata_sys::get_public_access_block_config(&bucket).await - && config.block_public_acls.unwrap_or(false) - && stored_acl.grants.iter().any(is_public_grant) - { - return Err(s3_error!(AccessDenied, "Access Denied")); - } - - let acl_data = serialize_acl(&stored_acl)?; - let mut eval_metadata = HashMap::new(); - eval_metadata.insert(INTERNAL_ACL_METADATA_KEY.to_string(), String::from_utf8_lossy(&acl_data).to_string()); - - let popts = ObjectOptions { - mod_time: info.mod_time, - version_id: opts.version_id, - eval_metadata: Some(eval_metadata), - ..Default::default() - }; - - store.put_object_metadata(&bucket, &key, &popts).await.map_err(|e| { - error!("put_object_metadata failed, {}", e.to_string()); - s3_error!(InternalError, "{}", e.to_string()) - })?; - - Ok(S3Response::new(PutObjectAclOutput::default())) + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_put_object_acl(req).await } async fn put_object_legal_hold( @@ -2642,99 +2437,8 @@ impl S3 for FS { #[instrument(level = "debug", skip(self, req))] async fn put_object_tagging(&self, req: S3Request) -> S3Result> { - let start_time = std::time::Instant::now(); - let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutTagging, "s3:PutObjectTagging"); - let PutObjectTaggingInput { - bucket, - key: object, - tagging, - .. - } = req.input.clone(); - - if tagging.tag_set.len() > 10 { - // TOTO: Note that Amazon S3 limits the maximum number of tags to 10 tags per object. - // Reference: https://docs.aws.amazon.com/AmazonS3/latest/userguide/object-tagging.html - // Reference: https://docs.aws.amazon.com/zh_cn/AmazonS3/latest/API/API_PutObjectTagging.html - // https://github.com/minio/mint/blob/master/run/core/aws-sdk-go-v2/main.go#L1647 - error!("Tag set exceeds maximum of 10 tags: {}", tagging.tag_set.len()); - return Err(s3_error!(InvalidTag, "Cannot have more than 10 tags per object")); - } - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let mut tag_keys = std::collections::HashSet::with_capacity(tagging.tag_set.len()); - for tag in &tagging.tag_set { - let key = tag.key.as_ref().filter(|k| !k.is_empty()).ok_or_else(|| { - error!("Empty tag key"); - s3_error!(InvalidTag, "Tag key cannot be empty") - })?; - - if key.len() > 128 { - error!("Tag key too long: {} bytes", key.len()); - return Err(s3_error!(InvalidTag, "Tag key is too long, maximum allowed length is 128 characters")); - } - - // allow to set the value of a tag to an empty string, but cannot set it to a null value. - // Reference:https://docs.aws.amazon.com/zh_cn/AWSEC2/latest/UserGuide/Using_Tags.html#tag-restrictions - let value = tag.value.as_ref().ok_or_else(|| { - error!("Null tag value"); - s3_error!(InvalidTag, "Tag value cannot be null") - })?; - - if value.len() > 256 { - error!("Tag value too long: {} bytes", value.len()); - return Err(s3_error!(InvalidTag, "Tag value is too long, maximum allowed length is 256 characters")); - } - - if !tag_keys.insert(key) { - error!("Duplicate tag key: {}", key); - return Err(s3_error!(InvalidTag, "Cannot provide multiple Tags with the same key")); - } - } - - let tags = encode_tags(tagging.tag_set); - debug!("Encoded tags: {}", tags); - - // TODO: getOpts, Replicate - // Support versioned objects - let version_id = req.input.version_id.clone(); - let opts = ObjectOptions { - version_id: self.parse_version_id(version_id)?.map(Into::into), - ..Default::default() - }; - - store.put_object_tags(&bucket, &object, &tags, &opts).await.map_err(|e| { - error!("Failed to put object tags: {}", e); - counter!("rustfs.put_object_tagging.failure").increment(1); - ApiError::from(e) - })?; - - // Invalidate cache for the tagged object - let manager = get_concurrency_manager(); - let version_id = req.input.version_id.clone(); - let cache_key = ConcurrencyManager::make_cache_key(&bucket, &object, version_id.clone().as_deref()); - tokio::spawn(async move { - manager - .invalidate_cache_versioned(&bucket, &object, version_id.as_deref()) - .await; - debug!("Cache invalidated for tagged object: {}", cache_key); - }); - - // Add metrics - counter!("rustfs.put_object_tagging.success").increment(1); - - let version_id_resp = req.input.version_id.clone().unwrap_or_default(); - helper = helper.version_id(version_id_resp); - - let result = Ok(S3Response::new(PutObjectTaggingOutput { - version_id: req.input.version_id.clone(), - })); - let _ = helper.complete(&result); - let duration = start_time.elapsed(); - histogram!("rustfs.object_tagging.operation.duration.seconds", "operation" => "put").record(duration.as_secs_f64()); - result + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_put_object_tagging(req).await } async fn restore_object(&self, req: S3Request) -> S3Result> {