diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index e42518962..d69935a08 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -26,12 +26,14 @@ use crate::storage::ecfs::{ }; use crate::storage::helper::OperationHelper; use crate::storage::*; +use http::StatusCode; use metrics::counter; use rustfs_ecstore::bucket::{ lifecycle::bucket_lifecycle_ops::validate_transition_tier, metadata::{ - BUCKET_ACL_CONFIG, BUCKET_LIFECYCLE_CONFIG, BUCKET_NOTIFICATION_CONFIG, BUCKET_POLICY_CONFIG, BUCKET_SSECONFIG, - BUCKET_TAGGING_CONFIG, BUCKET_VERSIONING_CONFIG, + BUCKET_ACL_CONFIG, 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_VERSIONING_CONFIG, }, metadata_sys, policy_sys::PolicySys, @@ -57,7 +59,7 @@ use s3s::dto::*; use s3s::xml; use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; use std::{fmt::Display, sync::Arc}; -use tracing::{debug, info, instrument, warn}; +use tracing::{debug, error, info, instrument, warn}; use urlencoding::encode; pub type BucketUsecaseResult = Result; @@ -325,6 +327,33 @@ impl DefaultBucketUsecase { Ok(S3Response::new(DeleteBucketEncryptionOutput::default())) } + #[instrument(level = "debug", skip(self))] + pub async fn execute_delete_bucket_cors( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let DeleteBucketCorsInput { bucket, .. } = req.input; + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + store + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .map_err(ApiError::from)?; + + metadata_sys::delete(&bucket, BUCKET_CORS_CONFIG) + .await + .map_err(ApiError::from)?; + + Ok(S3Response::new(DeleteBucketCorsOutput {})) + } + #[instrument(level = "debug", skip(self))] pub async fn execute_delete_bucket_lifecycle( &self, @@ -378,6 +407,34 @@ impl DefaultBucketUsecase { Ok(S3Response::new(DeleteBucketPolicyOutput {})) } + pub async fn execute_delete_bucket_replication( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let DeleteBucketReplicationInput { bucket, .. } = req.input; + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + store + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .map_err(ApiError::from)?; + metadata_sys::delete(&bucket, BUCKET_REPLICATION_CONFIG) + .await + .map_err(ApiError::from)?; + + // TODO: remove targets + 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, @@ -396,6 +453,33 @@ impl DefaultBucketUsecase { Ok(S3Response::new(DeleteBucketTaggingOutput {})) } + #[instrument(level = "debug", skip(self))] + pub async fn execute_delete_public_access_block( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let DeletePublicAccessBlockInput { bucket, .. } = req.input; + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + store + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .map_err(ApiError::from)?; + + metadata_sys::delete(&bucket, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG) + .await + .map_err(ApiError::from)?; + + Ok(S3Response::with_status(DeletePublicAccessBlockOutput::default(), StatusCode::NO_CONTENT)) + } + pub async fn execute_get_bucket_encryption( &self, req: S3Request, @@ -428,6 +512,42 @@ impl DefaultBucketUsecase { })) } + #[instrument(level = "debug", skip(self))] + pub async fn execute_get_bucket_cors(&self, req: S3Request) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let GetBucketCorsInput { bucket, .. } = req.input; + + let Some(store) = new_object_layer_fn() 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!("get_cors_config err {:?}", &err); + 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, @@ -632,6 +752,44 @@ impl DefaultBucketUsecase { Ok(S3Response::new(output)) } + pub async fn execute_get_bucket_replication( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let GetBucketReplicationInput { bucket, .. } = req.input; + + let Some(store) = new_object_layer_fn() 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(err) => { + error!("get_replication_config err {:?}", err); + if err == StorageError::ConfigNotFound { + return Err(S3Error::with_message( + S3ErrorCode::ReplicationConfigurationNotFoundError, + "replication not found".to_string(), + )); + } + 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, @@ -666,6 +824,44 @@ impl DefaultBucketUsecase { Ok(S3Response::new(GetBucketTaggingOutput { tag_set })) } + #[instrument(level = "debug", skip(self))] + pub async fn execute_get_public_access_block( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let GetPublicAccessBlockInput { bucket, .. } = req.input; + + let Some(store) = new_object_layer_fn() 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, @@ -889,6 +1085,100 @@ impl DefaultBucketUsecase { Ok(S3Response::new(PutBucketPolicyOutput {})) } + #[instrument(level = "debug", skip(self))] + pub async fn execute_put_bucket_cors(&self, req: S3Request) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let PutBucketCorsInput { + bucket, + cors_configuration, + .. + } = req.input; + + let Some(store) = new_object_layer_fn() 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)?; + metadata_sys::update(&bucket, BUCKET_CORS_CONFIG, data) + .await + .map_err(ApiError::from)?; + + Ok(S3Response::new(PutBucketCorsOutput::default())) + } + + pub async fn execute_put_bucket_replication( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let PutBucketReplicationInput { + bucket, + replication_configuration, + .. + } = req.input; + info!(bucket = %bucket, "updating bucket replication config"); + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + store + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .map_err(ApiError::from)?; + + // TODO: check enable, versioning enable + let data = serialize_config(&replication_configuration)?; + metadata_sys::update(&bucket, BUCKET_REPLICATION_CONFIG, data) + .await + .map_err(ApiError::from)?; + + Ok(S3Response::new(PutBucketReplicationOutput::default())) + } + + #[instrument(level = "debug", skip(self))] + pub async fn execute_put_public_access_block( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let PutPublicAccessBlockInput { + bucket, + public_access_block_configuration, + .. + } = req.input; + + let Some(store) = new_object_layer_fn() 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)?; + metadata_sys::update(&bucket, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, data) + .await + .map_err(ApiError::from)?; + + Ok(S3Response::new(PutPublicAccessBlockOutput::default())) + } + #[instrument(level = "debug", skip(self))] pub async fn execute_put_bucket_tagging( &self, @@ -1197,6 +1487,34 @@ mod tests { 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(); @@ -1222,6 +1540,62 @@ mod tests { 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_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_versioning_returns_internal_error_when_store_uninitialized() { let input = GetBucketVersioningInput::builder() @@ -1265,6 +1639,54 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InternalError); } + #[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() { + let input = PutBucketReplicationInput::builder() + .bucket("test-bucket".to_string()) + .replication_configuration(ReplicationConfiguration { + role: "arn:aws:iam::123456789012:role/test".to_string(), + rules: vec![], + }) + .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_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() diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 93b0b3418..e286082a0 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -28,16 +28,12 @@ use crate::storage::{ parse_object_lock_legal_hold, parse_object_lock_retention, validate_bucket_object_lock_enabled, }; use futures::StreamExt; -use http::{HeaderMap, StatusCode}; +use http::HeaderMap; use metrics::{counter, histogram}; use rustfs_ecstore::{ bucket::{ - metadata::{ - BUCKET_ACL_CONFIG, BUCKET_CORS_CONFIG, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, BUCKET_REPLICATION_CONFIG, - BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG, - }, + metadata::{BUCKET_ACL_CONFIG, BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG}, metadata_sys, - metadata_sys::get_replication_config, 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}, @@ -1013,22 +1009,8 @@ impl S3 for FS { #[instrument(level = "debug", skip(self))] async fn delete_bucket_cors(&self, req: S3Request) -> S3Result> { - let DeleteBucketCorsInput { bucket, .. } = req.input; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - store - .get_bucket_info(&bucket, &BucketOptions::default()) - .await - .map_err(ApiError::from)?; - - metadata_sys::delete(&bucket, BUCKET_CORS_CONFIG) - .await - .map_err(ApiError::from)?; - - Ok(S3Response::new(DeleteBucketCorsOutput {})) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_delete_bucket_cors(req).await } async fn delete_bucket_encryption( @@ -1060,24 +1042,8 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let DeleteBucketReplicationInput { bucket, .. } = req.input; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - store - .get_bucket_info(&bucket, &BucketOptions::default()) - .await - .map_err(ApiError::from)?; - metadata_sys::delete(&bucket, BUCKET_REPLICATION_CONFIG) - .await - .map_err(ApiError::from)?; - - // TODO: remove targets - error!("delete bucket"); - - Ok(S3Response::new(DeleteBucketReplicationOutput::default())) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_delete_bucket_replication(req).await } #[instrument(level = "debug", skip(self))] @@ -1094,22 +1060,8 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let DeletePublicAccessBlockInput { bucket, .. } = req.input; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - store - .get_bucket_info(&bucket, &BucketOptions::default()) - .await - .map_err(ApiError::from)?; - - metadata_sys::delete(&bucket, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG) - .await - .map_err(ApiError::from)?; - - Ok(S3Response::with_status(DeletePublicAccessBlockOutput::default(), StatusCode::NO_CONTENT)) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_delete_public_access_block(req).await } /// Delete an object @@ -1534,32 +1486,8 @@ impl S3 for FS { #[instrument(level = "debug", skip(self))] async fn get_bucket_cors(&self, req: S3Request) -> S3Result> { - let bucket = req.input.bucket.clone(); - // check bucket exists. - let _bucket = self - .head_bucket(req.map_input(|input| HeadBucketInput { - bucket: input.bucket, - expected_bucket_owner: None, - })) - .await?; - - 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!("get_cors_config err {:?}", &err); - return Err(ApiError::from(err).into()); - } - }; - - Ok(S3Response::new(GetBucketCorsOutput { - cors_rules: Some(cors_configuration.cors_rules), - })) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_get_bucket_cors(req).await } async fn get_bucket_encryption( @@ -1629,55 +1557,8 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let GetBucketReplicationInput { bucket, .. } = req.input; - - let Some(store) = new_object_layer_fn() 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 rcfg = match get_replication_config(&bucket).await { - Ok((cfg, _created)) => Some(cfg), - Err(err) => { - error!("get_replication_config err {:?}", err); - if err == StorageError::ConfigNotFound { - return Err(S3Error::with_message( - S3ErrorCode::ReplicationConfigurationNotFoundError, - "replication not found".to_string(), - )); - } - return Err(ApiError::from(err).into()); - } - }; - - if rcfg.is_none() { - return Err(S3Error::with_message( - S3ErrorCode::ReplicationConfigurationNotFoundError, - "replication not found".to_string(), - )); - } - - // Ok(S3Response::new(GetBucketReplicationOutput { - // replication_configuration: rcfg, - // })) - - if rcfg.is_some() { - Ok(S3Response::new(GetBucketReplicationOutput { - replication_configuration: rcfg, - })) - } else { - let rep = ReplicationConfiguration { - role: "".to_string(), - rules: vec![], - }; - Ok(S3Response::new(GetBucketReplicationOutput { - replication_configuration: Some(rep), - })) - } + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_get_bucket_replication(req).await } #[instrument(level = "debug", skip(self))] @@ -1691,31 +1572,8 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let bucket = req.input.bucket.clone(); - - let _bucket = self - .head_bucket(req.map_input(|input| HeadBucketInput { - bucket: input.bucket, - expected_bucket_owner: None, - })) - .await?; - - 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), - })) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_get_public_access_block(req).await } #[instrument(level = "debug", skip(self))] @@ -2432,28 +2290,8 @@ impl S3 for FS { #[instrument(level = "debug", skip(self))] async fn put_bucket_cors(&self, req: S3Request) -> S3Result> { - let PutBucketCorsInput { - bucket, - cors_configuration, - .. - } = req.input; - - let Some(store) = new_object_layer_fn() 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 = try_!(serialize(&cors_configuration)); - - metadata_sys::update(&bucket, BUCKET_CORS_CONFIG, data) - .await - .map_err(ApiError::from)?; - - Ok(S3Response::new(PutBucketCorsOutput::default())) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_put_bucket_cors(req).await } async fn put_bucket_encryption( @@ -2490,30 +2328,8 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let PutBucketReplicationInput { - bucket, - replication_configuration, - .. - } = req.input; - warn!("put bucket replication"); - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - store - .get_bucket_info(&bucket, &BucketOptions::default()) - .await - .map_err(ApiError::from)?; - - // TODO: check enable, versioning enable - let data = try_!(serialize(&replication_configuration)); - - metadata_sys::update(&bucket, BUCKET_REPLICATION_CONFIG, data) - .await - .map_err(ApiError::from)?; - - Ok(S3Response::new(PutBucketReplicationOutput::default())) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_put_bucket_replication(req).await } #[instrument(level = "debug", skip(self))] @@ -2521,28 +2337,8 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let PutPublicAccessBlockInput { - bucket, - public_access_block_configuration, - .. - } = req.input; - - let Some(store) = new_object_layer_fn() 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 = try_!(serialize(&public_access_block_configuration)); - - metadata_sys::update(&bucket, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, data) - .await - .map_err(ApiError::from)?; - - Ok(S3Response::new(PutPublicAccessBlockOutput::default())) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_put_public_access_block(req).await } #[instrument(level = "debug", skip(self))]