diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index d69935a08..90da26b6f 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -19,13 +19,14 @@ use crate::app::context::{AppContext, get_global_app_context}; use crate::auth::get_condition_values; use crate::error::ApiError; use crate::server::RemoteAddr; -use crate::storage::access::authorize_request; +use crate::storage::access::{ReqInfo, authorize_request}; use crate::storage::ecfs::{ - default_owner, is_public_grant, parse_acl_json_or_canned_bucket, serialize_acl, stored_acl_from_canned_bucket, - stored_acl_from_grant_headers, + RUSTFS_OWNER, default_owner, is_public_grant, parse_acl_json_or_canned_bucket, serialize_acl, stored_acl_from_canned_bucket, + stored_acl_from_grant_headers, stored_acl_from_policy, stored_grant_to_dto, stored_owner_to_dto, }; use crate::storage::helper::OperationHelper; use crate::storage::*; +use futures::StreamExt; use http::StatusCode; use metrics::counter; use rustfs_ecstore::bucket::{ @@ -234,6 +235,72 @@ impl DefaultBucketUsecase { result } + pub async fn execute_put_bucket_acl(&self, req: S3Request) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let PutBucketAclInput { + bucket, + acl, + access_control_policy, + grant_full_control, + grant_read, + grant_read_acp, + grant_write, + grant_write_acp, + .. + } = 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 owner = default_owner(); + let mut stored_acl = access_control_policy + .as_ref() + .map(|policy| stored_acl_from_policy(policy, &owner)) + .transpose()?; + + if stored_acl.is_none() { + stored_acl = stored_acl_from_grant_headers( + &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_bucket(canned.as_str(), &owner)); + } + + let stored_acl = stored_acl.unwrap_or_else(|| stored_acl_from_canned_bucket(BucketCannedACL::PRIVATE, &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 data = serialize_acl(&stored_acl)?; + metadata_sys::update(&bucket, BUCKET_ACL_CONFIG, data) + .await + .map_err(ApiError::from)?; + + Ok(S3Response::new(PutBucketAclOutput::default())) + } + #[instrument(level = "debug", skip(self, req))] pub async fn execute_delete_bucket(&self, mut req: S3Request) -> S3Result> { if let Some(context) = &self.context { @@ -301,6 +368,140 @@ impl DefaultBucketUsecase { Ok(S3Response::new(HeadBucketOutput::default())) } + pub async fn execute_get_bucket_acl(&self, req: S3Request) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let GetBucketAclInput { 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 owner = default_owner(); + let stored_acl = match metadata_sys::get_bucket_acl_config(&bucket).await { + Ok((acl, _)) => parse_acl_json_or_canned_bucket(&acl, &owner), + Err(err) => { + if err != StorageError::ConfigNotFound { + return Err(ApiError::from(err).into()); + } + stored_acl_from_canned_bucket(BucketCannedACL::PRIVATE, &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(GetBucketAclOutput { + grants: Some(grants), + owner: Some(stored_owner_to_dto(&stored_acl.owner)), + })) + } + + #[instrument(level = "debug", skip(self, req))] + pub async fn execute_get_bucket_location( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let input = 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(&input.bucket, &BucketOptions::default()) + .await + .map_err(ApiError::from)?; + + if let Some(region) = rustfs_ecstore::global::get_global_region() { + return Ok(S3Response::new(GetBucketLocationOutput { + location_constraint: Some(BucketLocationConstraint::from(region)), + })); + } + + Ok(S3Response::new(GetBucketLocationOutput::default())) + } + + #[instrument(level = "debug", skip(self))] + pub async fn execute_list_buckets(&self, req: S3Request) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let mut req = req; + + if req.credentials.as_ref().is_none_or(|cred| cred.access_key.is_empty()) { + return Err(S3Error::with_message(S3ErrorCode::AccessDenied, "Access Denied")); + } + + let bucket_infos = if let Err(e) = authorize_request(&mut req, Action::S3Action(S3Action::ListAllMyBucketsAction)).await { + if e.code() != &S3ErrorCode::AccessDenied { + return Err(e); + } + + let mut list_bucket_infos = store.list_bucket(&BucketOptions::default()).await.map_err(ApiError::from)?; + + list_bucket_infos = futures::stream::iter(list_bucket_infos) + .filter_map(|info| async { + let mut req_clone = req.clone(); + let req_info = req_clone.extensions.get_mut::().expect("ReqInfo not found"); + req_info.bucket = Some(info.name.clone()); + + if authorize_request(&mut req_clone, Action::S3Action(S3Action::ListBucketAction)) + .await + .is_ok() + || authorize_request(&mut req_clone, Action::S3Action(S3Action::GetBucketLocationAction)) + .await + .is_ok() + { + Some(info) + } else { + None + } + }) + .collect() + .await; + + if list_bucket_infos.is_empty() { + return Err(S3Error::with_message(S3ErrorCode::AccessDenied, "Access Denied")); + } + list_bucket_infos + } else { + store.list_bucket(&BucketOptions::default()).await.map_err(ApiError::from)? + }; + + let buckets: Vec = bucket_infos + .iter() + .map(|v| Bucket { + creation_date: v.created.map(Timestamp::from), + name: Some(v.name.clone()), + ..Default::default() + }) + .collect(); + + Ok(S3Response::new(ListBucketsOutput { + buckets: Some(buckets), + owner: Some(RUSTFS_OWNER.to_owned()), + ..Default::default() + })) + } + pub async fn execute_delete_bucket_encryption( &self, req: S3Request, @@ -1540,6 +1741,34 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InternalError); } + #[tokio::test] + async fn execute_get_bucket_acl_returns_internal_error_when_store_uninitialized() { + let input = GetBucketAclInput::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_acl(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() @@ -1610,6 +1839,17 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InternalError); } + #[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_put_bucket_lifecycle_configuration_rejects_missing_configuration() { let input = PutBucketLifecycleConfigurationInput::builder() @@ -1672,6 +1912,20 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InternalError); } + #[tokio::test] + async fn execute_put_bucket_acl_returns_internal_error_when_store_uninitialized() { + let input = PutBucketAclInput::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_acl(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() diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 9d0aca8de..ecf747428 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -30,8 +30,6 @@ use futures::StreamExt; use http::HeaderMap; use rustfs_ecstore::{ bucket::{ - metadata::BUCKET_ACL_CONFIG, - metadata_sys, object_lock::objectlock_sys::check_object_lock_for_deletion, replication::{DeletedObjectReplicationInfo, check_replicate_delete, schedule_replication_delete}, versioning::VersioningApi, @@ -44,7 +42,6 @@ use rustfs_ecstore::{ new_object_layer_fn, set_disk::MAX_PARTS_COUNT, store_api::{ - BucketOptions, ObjectIO, ObjectInfo, ObjectOptions, @@ -1377,36 +1374,8 @@ impl S3 for FS { } async fn get_bucket_acl(&self, req: S3Request) -> S3Result> { - let GetBucketAclInput { 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 owner = default_owner(); - let stored_acl = match metadata_sys::get_bucket_acl_config(&bucket).await { - Ok((acl, _)) => parse_acl_json_or_canned_bucket(&acl, &owner), - Err(err) => { - if err != StorageError::ConfigNotFound { - return Err(ApiError::from(err).into()); - } - stored_acl_from_canned_bucket(BucketCannedACL::PRIVATE, &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(GetBucketAclOutput { - grants: Some(grants), - owner: Some(stored_owner_to_dto(&stored_acl.owner)), - })) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_get_bucket_acl(req).await } #[instrument(level = "debug", skip(self))] @@ -1435,26 +1404,8 @@ impl S3 for FS { /// Get bucket location #[instrument(level = "debug", skip(self, req))] async fn get_bucket_location(&self, req: S3Request) -> S3Result> { - // mc get 1 - let input = 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(&input.bucket, &BucketOptions::default()) - .await - .map_err(ApiError::from)?; - - if let Some(region) = rustfs_ecstore::global::get_global_region() { - return Ok(S3Response::new(GetBucketLocationOutput { - location_constraint: Some(BucketLocationConstraint::from(region)), - })); - } - - let output = GetBucketLocationOutput::default(); - Ok(S3Response::new(output)) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_get_bucket_location(req).await } async fn get_bucket_notification_configuration( @@ -1587,69 +1538,8 @@ impl S3 for FS { #[instrument(level = "debug", skip(self))] async fn list_buckets(&self, req: S3Request) -> S3Result> { - // mc ls - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let mut req = req; - - if req.credentials.as_ref().is_none_or(|cred| cred.access_key.is_empty()) { - return Err(S3Error::with_message(S3ErrorCode::AccessDenied, "Access Denied")); - } - - let bucket_infos = if let Err(e) = authorize_request(&mut req, Action::S3Action(S3Action::ListAllMyBucketsAction)).await { - if e.code() != &S3ErrorCode::AccessDenied { - return Err(e); - } - - let mut list_bucket_infos = store.list_bucket(&BucketOptions::default()).await.map_err(ApiError::from)?; - - list_bucket_infos = futures::stream::iter(list_bucket_infos) - .filter_map(|info| async { - let mut req_clone = req.clone(); - let req_info = req_clone.extensions.get_mut::().expect("ReqInfo not found"); - req_info.bucket = Some(info.name.clone()); - - if authorize_request(&mut req_clone, Action::S3Action(S3Action::ListBucketAction)) - .await - .is_ok() - || authorize_request(&mut req_clone, Action::S3Action(S3Action::GetBucketLocationAction)) - .await - .is_ok() - { - Some(info) - } else { - None - } - }) - .collect() - .await; - - if list_bucket_infos.is_empty() { - return Err(S3Error::with_message(S3ErrorCode::AccessDenied, "Access Denied")); - } - list_bucket_infos - } else { - store.list_bucket(&BucketOptions::default()).await.map_err(ApiError::from)? - }; - - let buckets: Vec = bucket_infos - .iter() - .map(|v| Bucket { - creation_date: v.created.map(Timestamp::from), - name: Some(v.name.clone()), - ..Default::default() - }) - .collect(); - - let output = ListBucketsOutput { - buckets: Some(buckets), - owner: Some(RUSTFS_OWNER.to_owned()), - ..Default::default() - }; - Ok(S3Response::new(output)) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_list_buckets(req).await } async fn list_multipart_uploads( @@ -1943,67 +1833,8 @@ impl S3 for FS { } async fn put_bucket_acl(&self, req: S3Request) -> S3Result> { - let PutBucketAclInput { - bucket, - acl, - access_control_policy, - grant_full_control, - grant_read, - grant_read_acp, - grant_write, - grant_write_acp, - .. - } = req.input; - - // TODO:checkRequestAuthType - - 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 owner = default_owner(); - let mut stored_acl = access_control_policy - .as_ref() - .map(|policy| stored_acl_from_policy(policy, &owner)) - .transpose()?; - - if stored_acl.is_none() { - stored_acl = stored_acl_from_grant_headers( - &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_bucket(canned.as_str(), &owner)); - } - - let stored_acl = stored_acl.unwrap_or_else(|| stored_acl_from_canned_bucket(BucketCannedACL::PRIVATE, &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 data = serialize_acl(&stored_acl)?; - metadata_sys::update(&bucket, BUCKET_ACL_CONFIG, data) - .await - .map_err(ApiError::from)?; - - Ok(S3Response::new(PutBucketAclOutput::default())) + let usecase = DefaultBucketUsecase::from_global(); + usecase.execute_put_bucket_acl(req).await } #[instrument(level = "debug", skip(self))]