mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-27 15:37:02 +00:00
refactor(app): migrate bucket config flows to usecase (#1919)
This commit is contained in:
@@ -26,12 +26,14 @@ use crate::storage::ecfs::{
|
|||||||
};
|
};
|
||||||
use crate::storage::helper::OperationHelper;
|
use crate::storage::helper::OperationHelper;
|
||||||
use crate::storage::*;
|
use crate::storage::*;
|
||||||
|
use http::StatusCode;
|
||||||
use metrics::counter;
|
use metrics::counter;
|
||||||
use rustfs_ecstore::bucket::{
|
use rustfs_ecstore::bucket::{
|
||||||
lifecycle::bucket_lifecycle_ops::validate_transition_tier,
|
lifecycle::bucket_lifecycle_ops::validate_transition_tier,
|
||||||
metadata::{
|
metadata::{
|
||||||
BUCKET_ACL_CONFIG, BUCKET_LIFECYCLE_CONFIG, BUCKET_NOTIFICATION_CONFIG, BUCKET_POLICY_CONFIG, BUCKET_SSECONFIG,
|
BUCKET_ACL_CONFIG, BUCKET_CORS_CONFIG, BUCKET_LIFECYCLE_CONFIG, BUCKET_NOTIFICATION_CONFIG, BUCKET_POLICY_CONFIG,
|
||||||
BUCKET_TAGGING_CONFIG, BUCKET_VERSIONING_CONFIG,
|
BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, BUCKET_REPLICATION_CONFIG, BUCKET_SSECONFIG, BUCKET_TAGGING_CONFIG,
|
||||||
|
BUCKET_VERSIONING_CONFIG,
|
||||||
},
|
},
|
||||||
metadata_sys,
|
metadata_sys,
|
||||||
policy_sys::PolicySys,
|
policy_sys::PolicySys,
|
||||||
@@ -57,7 +59,7 @@ use s3s::dto::*;
|
|||||||
use s3s::xml;
|
use s3s::xml;
|
||||||
use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
|
use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
|
||||||
use std::{fmt::Display, sync::Arc};
|
use std::{fmt::Display, sync::Arc};
|
||||||
use tracing::{debug, info, instrument, warn};
|
use tracing::{debug, error, info, instrument, warn};
|
||||||
use urlencoding::encode;
|
use urlencoding::encode;
|
||||||
|
|
||||||
pub type BucketUsecaseResult<T> = Result<T, ApiError>;
|
pub type BucketUsecaseResult<T> = Result<T, ApiError>;
|
||||||
@@ -325,6 +327,33 @@ impl DefaultBucketUsecase {
|
|||||||
Ok(S3Response::new(DeleteBucketEncryptionOutput::default()))
|
Ok(S3Response::new(DeleteBucketEncryptionOutput::default()))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[instrument(level = "debug", skip(self))]
|
||||||
|
pub async fn execute_delete_bucket_cors(
|
||||||
|
&self,
|
||||||
|
req: S3Request<DeleteBucketCorsInput>,
|
||||||
|
) -> S3Result<S3Response<DeleteBucketCorsOutput>> {
|
||||||
|
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))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
pub async fn execute_delete_bucket_lifecycle(
|
pub async fn execute_delete_bucket_lifecycle(
|
||||||
&self,
|
&self,
|
||||||
@@ -378,6 +407,34 @@ impl DefaultBucketUsecase {
|
|||||||
Ok(S3Response::new(DeleteBucketPolicyOutput {}))
|
Ok(S3Response::new(DeleteBucketPolicyOutput {}))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn execute_delete_bucket_replication(
|
||||||
|
&self,
|
||||||
|
req: S3Request<DeleteBucketReplicationInput>,
|
||||||
|
) -> S3Result<S3Response<DeleteBucketReplicationOutput>> {
|
||||||
|
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))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
pub async fn execute_delete_bucket_tagging(
|
pub async fn execute_delete_bucket_tagging(
|
||||||
&self,
|
&self,
|
||||||
@@ -396,6 +453,33 @@ impl DefaultBucketUsecase {
|
|||||||
Ok(S3Response::new(DeleteBucketTaggingOutput {}))
|
Ok(S3Response::new(DeleteBucketTaggingOutput {}))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[instrument(level = "debug", skip(self))]
|
||||||
|
pub async fn execute_delete_public_access_block(
|
||||||
|
&self,
|
||||||
|
req: S3Request<DeletePublicAccessBlockInput>,
|
||||||
|
) -> S3Result<S3Response<DeletePublicAccessBlockOutput>> {
|
||||||
|
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(
|
pub async fn execute_get_bucket_encryption(
|
||||||
&self,
|
&self,
|
||||||
req: S3Request<GetBucketEncryptionInput>,
|
req: S3Request<GetBucketEncryptionInput>,
|
||||||
@@ -428,6 +512,42 @@ impl DefaultBucketUsecase {
|
|||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[instrument(level = "debug", skip(self))]
|
||||||
|
pub async fn execute_get_bucket_cors(&self, req: S3Request<GetBucketCorsInput>) -> S3Result<S3Response<GetBucketCorsOutput>> {
|
||||||
|
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))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
pub async fn execute_get_bucket_lifecycle_configuration(
|
pub async fn execute_get_bucket_lifecycle_configuration(
|
||||||
&self,
|
&self,
|
||||||
@@ -632,6 +752,44 @@ impl DefaultBucketUsecase {
|
|||||||
Ok(S3Response::new(output))
|
Ok(S3Response::new(output))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn execute_get_bucket_replication(
|
||||||
|
&self,
|
||||||
|
req: S3Request<GetBucketReplicationInput>,
|
||||||
|
) -> S3Result<S3Response<GetBucketReplicationOutput>> {
|
||||||
|
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))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
pub async fn execute_get_bucket_tagging(
|
pub async fn execute_get_bucket_tagging(
|
||||||
&self,
|
&self,
|
||||||
@@ -666,6 +824,44 @@ impl DefaultBucketUsecase {
|
|||||||
Ok(S3Response::new(GetBucketTaggingOutput { tag_set }))
|
Ok(S3Response::new(GetBucketTaggingOutput { tag_set }))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[instrument(level = "debug", skip(self))]
|
||||||
|
pub async fn execute_get_public_access_block(
|
||||||
|
&self,
|
||||||
|
req: S3Request<GetPublicAccessBlockInput>,
|
||||||
|
) -> S3Result<S3Response<GetPublicAccessBlockOutput>> {
|
||||||
|
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))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
pub async fn execute_get_bucket_versioning(
|
pub async fn execute_get_bucket_versioning(
|
||||||
&self,
|
&self,
|
||||||
@@ -889,6 +1085,100 @@ impl DefaultBucketUsecase {
|
|||||||
Ok(S3Response::new(PutBucketPolicyOutput {}))
|
Ok(S3Response::new(PutBucketPolicyOutput {}))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[instrument(level = "debug", skip(self))]
|
||||||
|
pub async fn execute_put_bucket_cors(&self, req: S3Request<PutBucketCorsInput>) -> S3Result<S3Response<PutBucketCorsOutput>> {
|
||||||
|
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<PutBucketReplicationInput>,
|
||||||
|
) -> S3Result<S3Response<PutBucketReplicationOutput>> {
|
||||||
|
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<PutPublicAccessBlockInput>,
|
||||||
|
) -> S3Result<S3Response<PutPublicAccessBlockOutput>> {
|
||||||
|
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))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
pub async fn execute_put_bucket_tagging(
|
pub async fn execute_put_bucket_tagging(
|
||||||
&self,
|
&self,
|
||||||
@@ -1197,6 +1487,34 @@ mod tests {
|
|||||||
assert_eq!(err.code(), &S3ErrorCode::InternalError);
|
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]
|
#[tokio::test]
|
||||||
async fn execute_head_bucket_returns_internal_error_when_store_uninitialized() {
|
async fn execute_head_bucket_returns_internal_error_when_store_uninitialized() {
|
||||||
let input = HeadBucketInput::builder().bucket("test-bucket".to_string()).build().unwrap();
|
let input = HeadBucketInput::builder().bucket("test-bucket".to_string()).build().unwrap();
|
||||||
@@ -1222,6 +1540,62 @@ mod tests {
|
|||||||
assert_eq!(err.code(), &S3ErrorCode::InternalError);
|
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]
|
#[tokio::test]
|
||||||
async fn execute_get_bucket_versioning_returns_internal_error_when_store_uninitialized() {
|
async fn execute_get_bucket_versioning_returns_internal_error_when_store_uninitialized() {
|
||||||
let input = GetBucketVersioningInput::builder()
|
let input = GetBucketVersioningInput::builder()
|
||||||
@@ -1265,6 +1639,54 @@ mod tests {
|
|||||||
assert_eq!(err.code(), &S3ErrorCode::InternalError);
|
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]
|
#[tokio::test]
|
||||||
async fn execute_list_objects_v2_rejects_negative_max_keys() {
|
async fn execute_list_objects_v2_rejects_negative_max_keys() {
|
||||||
let input = ListObjectsV2Input::builder()
|
let input = ListObjectsV2Input::builder()
|
||||||
|
|||||||
+20
-224
@@ -28,16 +28,12 @@ use crate::storage::{
|
|||||||
parse_object_lock_legal_hold, parse_object_lock_retention, validate_bucket_object_lock_enabled,
|
parse_object_lock_legal_hold, parse_object_lock_retention, validate_bucket_object_lock_enabled,
|
||||||
};
|
};
|
||||||
use futures::StreamExt;
|
use futures::StreamExt;
|
||||||
use http::{HeaderMap, StatusCode};
|
use http::HeaderMap;
|
||||||
use metrics::{counter, histogram};
|
use metrics::{counter, histogram};
|
||||||
use rustfs_ecstore::{
|
use rustfs_ecstore::{
|
||||||
bucket::{
|
bucket::{
|
||||||
metadata::{
|
metadata::{BUCKET_ACL_CONFIG, BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG},
|
||||||
BUCKET_ACL_CONFIG, BUCKET_CORS_CONFIG, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, BUCKET_REPLICATION_CONFIG,
|
|
||||||
BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG,
|
|
||||||
},
|
|
||||||
metadata_sys,
|
metadata_sys,
|
||||||
metadata_sys::get_replication_config,
|
|
||||||
object_lock::objectlock_sys::{check_object_lock_for_deletion, check_retention_for_modification},
|
object_lock::objectlock_sys::{check_object_lock_for_deletion, check_retention_for_modification},
|
||||||
replication::{DeletedObjectReplicationInfo, check_replicate_delete, schedule_replication_delete},
|
replication::{DeletedObjectReplicationInfo, check_replicate_delete, schedule_replication_delete},
|
||||||
tagging::{decode_tags, encode_tags},
|
tagging::{decode_tags, encode_tags},
|
||||||
@@ -1013,22 +1009,8 @@ impl S3 for FS {
|
|||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
async fn delete_bucket_cors(&self, req: S3Request<DeleteBucketCorsInput>) -> S3Result<S3Response<DeleteBucketCorsOutput>> {
|
async fn delete_bucket_cors(&self, req: S3Request<DeleteBucketCorsInput>) -> S3Result<S3Response<DeleteBucketCorsOutput>> {
|
||||||
let DeleteBucketCorsInput { bucket, .. } = req.input;
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
|
usecase.execute_delete_bucket_cors(req).await
|
||||||
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 {}))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn delete_bucket_encryption(
|
async fn delete_bucket_encryption(
|
||||||
@@ -1060,24 +1042,8 @@ impl S3 for FS {
|
|||||||
&self,
|
&self,
|
||||||
req: S3Request<DeleteBucketReplicationInput>,
|
req: S3Request<DeleteBucketReplicationInput>,
|
||||||
) -> S3Result<S3Response<DeleteBucketReplicationOutput>> {
|
) -> S3Result<S3Response<DeleteBucketReplicationOutput>> {
|
||||||
let DeleteBucketReplicationInput { bucket, .. } = req.input;
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
|
usecase.execute_delete_bucket_replication(req).await
|
||||||
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()))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
@@ -1094,22 +1060,8 @@ impl S3 for FS {
|
|||||||
&self,
|
&self,
|
||||||
req: S3Request<DeletePublicAccessBlockInput>,
|
req: S3Request<DeletePublicAccessBlockInput>,
|
||||||
) -> S3Result<S3Response<DeletePublicAccessBlockOutput>> {
|
) -> S3Result<S3Response<DeletePublicAccessBlockOutput>> {
|
||||||
let DeletePublicAccessBlockInput { bucket, .. } = req.input;
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
|
usecase.execute_delete_public_access_block(req).await
|
||||||
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))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Delete an object
|
/// Delete an object
|
||||||
@@ -1534,32 +1486,8 @@ impl S3 for FS {
|
|||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
async fn get_bucket_cors(&self, req: S3Request<GetBucketCorsInput>) -> S3Result<S3Response<GetBucketCorsOutput>> {
|
async fn get_bucket_cors(&self, req: S3Request<GetBucketCorsInput>) -> S3Result<S3Response<GetBucketCorsOutput>> {
|
||||||
let bucket = req.input.bucket.clone();
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
// check bucket exists.
|
usecase.execute_get_bucket_cors(req).await
|
||||||
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),
|
|
||||||
}))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn get_bucket_encryption(
|
async fn get_bucket_encryption(
|
||||||
@@ -1629,55 +1557,8 @@ impl S3 for FS {
|
|||||||
&self,
|
&self,
|
||||||
req: S3Request<GetBucketReplicationInput>,
|
req: S3Request<GetBucketReplicationInput>,
|
||||||
) -> S3Result<S3Response<GetBucketReplicationOutput>> {
|
) -> S3Result<S3Response<GetBucketReplicationOutput>> {
|
||||||
let GetBucketReplicationInput { bucket, .. } = req.input;
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
|
usecase.execute_get_bucket_replication(req).await
|
||||||
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),
|
|
||||||
}))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
@@ -1691,31 +1572,8 @@ impl S3 for FS {
|
|||||||
&self,
|
&self,
|
||||||
req: S3Request<GetPublicAccessBlockInput>,
|
req: S3Request<GetPublicAccessBlockInput>,
|
||||||
) -> S3Result<S3Response<GetPublicAccessBlockOutput>> {
|
) -> S3Result<S3Response<GetPublicAccessBlockOutput>> {
|
||||||
let bucket = req.input.bucket.clone();
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
|
usecase.execute_get_public_access_block(req).await
|
||||||
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),
|
|
||||||
}))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
@@ -2432,28 +2290,8 @@ impl S3 for FS {
|
|||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
async fn put_bucket_cors(&self, req: S3Request<PutBucketCorsInput>) -> S3Result<S3Response<PutBucketCorsOutput>> {
|
async fn put_bucket_cors(&self, req: S3Request<PutBucketCorsInput>) -> S3Result<S3Response<PutBucketCorsOutput>> {
|
||||||
let PutBucketCorsInput {
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
bucket,
|
usecase.execute_put_bucket_cors(req).await
|
||||||
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()))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn put_bucket_encryption(
|
async fn put_bucket_encryption(
|
||||||
@@ -2490,30 +2328,8 @@ impl S3 for FS {
|
|||||||
&self,
|
&self,
|
||||||
req: S3Request<PutBucketReplicationInput>,
|
req: S3Request<PutBucketReplicationInput>,
|
||||||
) -> S3Result<S3Response<PutBucketReplicationOutput>> {
|
) -> S3Result<S3Response<PutBucketReplicationOutput>> {
|
||||||
let PutBucketReplicationInput {
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
bucket,
|
usecase.execute_put_bucket_replication(req).await
|
||||||
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()))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
@@ -2521,28 +2337,8 @@ impl S3 for FS {
|
|||||||
&self,
|
&self,
|
||||||
req: S3Request<PutPublicAccessBlockInput>,
|
req: S3Request<PutPublicAccessBlockInput>,
|
||||||
) -> S3Result<S3Response<PutPublicAccessBlockOutput>> {
|
) -> S3Result<S3Response<PutPublicAccessBlockOutput>> {
|
||||||
let PutPublicAccessBlockInput {
|
let usecase = DefaultBucketUsecase::from_global();
|
||||||
bucket,
|
usecase.execute_put_public_access_block(req).await
|
||||||
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()))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
|
|||||||
Reference in New Issue
Block a user