add to_s3_error

This commit is contained in:
weisd
2024-11-21 17:25:49 +08:00
parent f88126008b
commit 1bfba1e913
4 changed files with 201 additions and 95 deletions
+1 -1
View File
@@ -23,7 +23,7 @@ pub mod bucket;
pub mod file_meta_inline;
pub mod options;
pub mod pools;
pub(crate) mod store_err;
pub mod store_err;
pub mod xhttp;
pub use global::is_legacy;
+127 -94
View File
@@ -15,9 +15,7 @@ use ecstore::bucket::policy_sys::PolicySys;
use ecstore::bucket::tagging::decode_tags;
use ecstore::bucket::tagging::encode_tags;
use ecstore::bucket::versioning_sys::BucketVersioningSys;
use ecstore::disk::error::is_err_file_not_found;
use ecstore::disk::error::DiskError;
use ecstore::error::Error as EcError;
use ecstore::new_object_layer_fn;
use ecstore::options::extract_metadata;
use ecstore::options::put_opts;
@@ -54,6 +52,8 @@ use tracing::info;
use transform_stream::AsyncTryStream;
use uuid::Uuid;
use crate::storage::error::to_s3_error;
macro_rules! try_ {
($result:expr) => {
match $result {
@@ -72,14 +72,6 @@ lazy_static! {
};
}
fn to_s3_error(err: EcError) -> S3Error {
if is_err_file_not_found(&err) {
return S3Error::with_message(S3ErrorCode::NoSuchKey, format!(" ec err {}", err));
}
S3Error::with_message(S3ErrorCode::InternalError, format!(" ec err {}", err))
}
#[derive(Debug, Clone)]
pub struct FS {
// pub store: ECStore,
@@ -112,18 +104,17 @@ impl S3 for FS {
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
try_!(
store
.make_bucket(
&bucket,
&MakeBucketOptions {
force_create: true,
lock_enabled: object_lock_enabled_for_bucket.is_some_and(|v| v),
..Default::default()
}
)
.await
);
store
.make_bucket(
&bucket,
&MakeBucketOptions {
force_create: true,
lock_enabled: object_lock_enabled_for_bucket.is_some_and(|v| v),
..Default::default()
},
)
.await
.map_err(to_s3_error)?;
let output = CreateBucketOutput::default();
Ok(S3Response::new(output))
@@ -151,17 +142,17 @@ impl S3 for FS {
Some(s) => s,
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
try_!(
store
.delete_bucket(
&input.bucket,
&DeleteBucketOptions {
force: false,
..Default::default()
}
)
.await
);
store
.delete_bucket(
&input.bucket,
&DeleteBucketOptions {
force: false,
..Default::default()
},
)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(DeleteBucketOutput {}))
}
@@ -192,7 +183,10 @@ impl S3 for FS {
Some(s) => s,
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
let (dobjs, _errs) = try_!(store.delete_objects(&bucket, objects, ObjectOptions::default()).await);
let (dobjs, _errs) = store
.delete_objects(&bucket, objects, ObjectOptions::default())
.await
.map_err(to_s3_error)?;
// TODO: let errors;
@@ -260,7 +254,10 @@ impl S3 for FS {
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
let (dobjs, _errs) = try_!(store.delete_objects(&bucket, objects, ObjectOptions::default()).await);
let (dobjs, _errs) = store
.delete_objects(&bucket, objects, ObjectOptions::default())
.await
.map_err(to_s3_error)?;
// info!("delete_objects res {:?} {:?}", &dobjs, errs);
let deleted = dobjs
@@ -335,7 +332,10 @@ impl S3 for FS {
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
let reader = try_!(store.get_object_reader(bucket.as_str(), key.as_str(), range, h, opts).await);
let reader = store
.get_object_reader(bucket.as_str(), key.as_str(), range, h, opts)
.await
.map_err(to_s3_error)?;
let info = reader.object_info;
@@ -445,7 +445,7 @@ impl S3 for FS {
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
let bucket_infos = try_!(store.list_bucket(&BucketOptions::default()).await);
let bucket_infos = store.list_bucket(&BucketOptions::default()).await.map_err(to_s3_error)?;
let buckets: Vec<Bucket> = bucket_infos
.iter()
@@ -503,19 +503,18 @@ impl S3 for FS {
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
let object_infos = try_!(
store
.list_objects_v2(
&bucket,
&prefix,
&continuation_token.unwrap_or_default(),
&delimiter,
max_keys.unwrap_or_default(),
fetch_owner.unwrap_or_default(),
&start_after.unwrap_or_default()
)
.await
);
let object_infos = store
.list_objects_v2(
&bucket,
&prefix,
&continuation_token.unwrap_or_default(),
&delimiter,
max_keys.unwrap_or_default(),
fetch_owner.unwrap_or_default(),
&start_after.unwrap_or_default(),
)
.await
.map_err(to_s3_error)?;
// warn!("object_infos {:?}", object_infos);
@@ -614,9 +613,14 @@ impl S3 for FS {
metadata.insert(xhttp::AMZ_OBJECT_TAGGING.to_owned(), tags);
}
let opts: ObjectOptions = try_!(put_opts(&bucket, &key, None, &req.headers, Some(metadata)).await);
let opts: ObjectOptions = put_opts(&bucket, &key, None, &req.headers, Some(metadata))
.await
.map_err(to_s3_error)?;
let obj_info = try_!(store.put_object(&bucket, &key, &mut reader, &opts).await);
let obj_info = store
.put_object(&bucket, &key, &mut reader, &opts)
.await
.map_err(to_s3_error)?;
let e_tag = obj_info.etag;
@@ -655,9 +659,12 @@ impl S3 for FS {
metadata.insert(xhttp::AMZ_OBJECT_TAGGING.to_owned(), tags);
}
let opts: ObjectOptions = try_!(put_opts(&bucket, &key, None, &req.headers, Some(metadata)).await);
let opts: ObjectOptions = put_opts(&bucket, &key, None, &req.headers, Some(metadata))
.await
.map_err(to_s3_error)?;
let MultipartUploadResult { upload_id, .. } = try_!(store.new_multipart_upload(&bucket, &key, &opts).await);
let MultipartUploadResult { upload_id, .. } =
store.new_multipart_upload(&bucket, &key, &opts).await.map_err(to_s3_error)?;
let output = CreateMultipartUploadOutput {
bucket: Some(bucket),
@@ -714,11 +721,10 @@ impl S3 for FS {
// TODO: hash_reader
let info = try_!(
store
.put_object_part(&bucket, &key, &upload_id, part_id, &mut data, &opts)
.await
);
let info = store
.put_object_part(&bucket, &key, &upload_id, part_id, &mut data, &opts)
.await
.map_err(to_s3_error)?;
let output = UploadPartOutput {
e_tag: info.etag,
@@ -784,11 +790,10 @@ impl S3 for FS {
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
try_!(
store
.complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, opts)
.await
);
store
.complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, opts)
.await
.map_err(to_s3_error)?;
let output = CompleteMultipartUploadOutput {
bucket: Some(bucket),
@@ -815,11 +820,11 @@ impl S3 for FS {
};
let opts = &ObjectOptions::default();
try_!(
store
.abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts)
.await
);
store
.abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(AbortMultipartUploadOutput { ..Default::default() }))
}
@@ -868,7 +873,9 @@ impl S3 for FS {
let data = try_!(xml::serialize(&tagging));
try_!(metadata_sys::update(&bucket, BUCKET_TAGGING_CONFIG, data).await);
metadata_sys::update(&bucket, BUCKET_TAGGING_CONFIG, data)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(Default::default()))
}
@@ -880,7 +887,9 @@ impl S3 for FS {
) -> S3Result<S3Response<DeleteBucketTaggingOutput>> {
let DeleteBucketTaggingInput { bucket, .. } = req.input;
try_!(metadata_sys::delete(&bucket, BUCKET_TAGGING_CONFIG).await);
metadata_sys::delete(&bucket, BUCKET_TAGGING_CONFIG)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(DeleteBucketTaggingOutput {}))
}
@@ -900,17 +909,15 @@ impl S3 for FS {
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
// let mut object_info = try_!(store.get_object_info(&bucket, &object, &ObjectOptions::default()).await);
let tags = encode_tags(tagging.tag_set);
// TODO: getOpts
// TODO: Replicate
try_!(
store
.put_object_tags(&bucket, &object, &tags, &ObjectOptions::default())
.await
);
store
.put_object_tags(&bucket, &object, &tags, &ObjectOptions::default())
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(PutObjectTaggingOutput { version_id: None }))
}
@@ -926,7 +933,10 @@ impl S3 for FS {
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
// TODO: version
let tags = try_!(store.get_object_tags(&bucket, &object, &ObjectOptions::default()).await);
let tags = store
.get_object_tags(&bucket, &object, &ObjectOptions::default())
.await
.map_err(to_s3_error)?;
let tag_set = decode_tags(tags.as_str());
@@ -951,7 +961,10 @@ impl S3 for FS {
// TODO: Replicate
// TODO: version
try_!(store.delete_object_tags(&bucket, &object, &ObjectOptions::default()).await);
store
.delete_object_tags(&bucket, &object, &ObjectOptions::default())
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(DeleteObjectTaggingOutput { version_id: None }))
}
@@ -977,7 +990,7 @@ impl S3 for FS {
}
}
let VersioningConfiguration { status, .. } = try_!(BucketVersioningSys::get(&bucket).await);
let VersioningConfiguration { status, .. } = BucketVersioningSys::get(&bucket).await.map_err(to_s3_error)?;
Ok(S3Response::new(GetBucketVersioningOutput {
status,
@@ -1003,7 +1016,9 @@ impl S3 for FS {
let data = try_!(xml::serialize(&versioning_configuration));
try_!(metadata_sys::update(&bucket, BUCKET_VERSIONING_CONFIG, data).await);
metadata_sys::update(&bucket, BUCKET_VERSIONING_CONFIG, data)
.await
.map_err(to_s3_error)?;
// TODO: globalSiteReplicationSys.BucketMetaHook
@@ -1068,7 +1083,7 @@ impl S3 for FS {
// warn!("input policy {}", &policy);
let cfg = try_!(BucketPolicy::unmarshal(policy.as_bytes()));
let cfg = BucketPolicy::unmarshal(policy.as_bytes()).map_err(to_s3_error)?;
// warn!("parse policy {:?}", &cfg);
@@ -1077,9 +1092,11 @@ impl S3 for FS {
return Err(s3_error!(InvalidPolicyDocument));
}
let data = try_!(cfg.marshal_msg());
let data = cfg.marshal_msg().map_err(to_s3_error)?;
try_!(metadata_sys::update(&bucket, BUCKET_POLICY_CONFIG, data.into()).await);
metadata_sys::update(&bucket, BUCKET_POLICY_CONFIG, data.into())
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(PutBucketPolicyOutput {}))
}
@@ -1104,7 +1121,9 @@ impl S3 for FS {
}
}
try_!(metadata_sys::delete(&bucket, BUCKET_POLICY_CONFIG).await);
metadata_sys::delete(&bucket, BUCKET_POLICY_CONFIG)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(DeleteBucketPolicyOutput {}))
}
@@ -1165,7 +1184,9 @@ impl S3 for FS {
let Some(input_cfg) = lifecycle_configuration else { return Err(s3_error!(InvalidArgument)) };
let data = try_!(xml::serialize(&input_cfg));
try_!(metadata_sys::update(&bucket, BUCKET_LIFECYCLE_CONFIG, data).await);
metadata_sys::update(&bucket, BUCKET_LIFECYCLE_CONFIG, data)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(PutBucketLifecycleConfigurationOutput::default()))
}
@@ -1191,7 +1212,9 @@ impl S3 for FS {
}
}
try_!(metadata_sys::delete(&bucket, BUCKET_LIFECYCLE_CONFIG).await);
metadata_sys::delete(&bucket, BUCKET_LIFECYCLE_CONFIG)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(DeleteBucketLifecycleOutput::default()))
}
@@ -1261,7 +1284,9 @@ impl S3 for FS {
// TODO: check kms
let data = try_!(xml::serialize(&server_side_encryption_configuration));
try_!(metadata_sys::update(&bucket, BUCKET_SSECONFIG, data).await);
metadata_sys::update(&bucket, BUCKET_SSECONFIG, data)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(PutBucketEncryptionOutput::default()))
}
@@ -1284,7 +1309,7 @@ impl S3 for FS {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
try_!(metadata_sys::delete(&bucket, BUCKET_SSECONFIG).await);
metadata_sys::delete(&bucket, BUCKET_SSECONFIG).await.map_err(to_s3_error)?;
Ok(S3Response::new(DeleteBucketEncryptionOutput::default()))
}
@@ -1340,7 +1365,9 @@ impl S3 for FS {
let data = try_!(xml::serialize(&input_cfg));
try_!(metadata_sys::update(&bucket, OBJECT_LOCK_CONFIG, data).await);
metadata_sys::update(&bucket, OBJECT_LOCK_CONFIG, data)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(PutObjectLockConfigurationOutput::default()))
}
@@ -1405,7 +1432,9 @@ impl S3 for FS {
// TODO: check enable, versioning enable
let data = try_!(xml::serialize(&replication_configuration));
try_!(metadata_sys::update(&bucket, BUCKET_REPLICATION_CONFIG, data).await);
metadata_sys::update(&bucket, BUCKET_REPLICATION_CONFIG, data)
.await
.map_err(to_s3_error)?;
Ok(S3Response::new(PutBucketReplicationOutput::default()))
}
@@ -1429,7 +1458,9 @@ impl S3 for FS {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
try_!(metadata_sys::delete(&bucket, BUCKET_REPLICATION_CONFIG).await);
metadata_sys::delete(&bucket, BUCKET_REPLICATION_CONFIG)
.await
.map_err(to_s3_error)?;
// TODO: remove targets
@@ -1510,7 +1541,9 @@ impl S3 for FS {
let data = try_!(xml::serialize(&notification_configuration));
try_!(metadata_sys::update(&bucket, BUCKET_NOTIFICATION_CONFIG, data).await);
metadata_sys::update(&bucket, BUCKET_NOTIFICATION_CONFIG, data)
.await
.map_err(to_s3_error)?;
// TODO: event notice add rule
+72
View File
@@ -0,0 +1,72 @@
use ecstore::{disk::error::is_err_file_not_found, error::Error, store_err::StorageError};
use s3s::{s3_error, S3Error, S3ErrorCode};
pub fn to_s3_error(err: Error) -> S3Error {
if let Some(storage_err) = err.downcast_ref::<StorageError>() {
return match storage_err {
StorageError::NotImplemented => s3_error!(NotImplemented),
StorageError::InvalidArgument(bucket, object, version_id) => {
s3_error!(InvalidArgument, "Invalid arguments provided for {}/{}-{}", bucket, object, version_id)
}
StorageError::MethodNotAllowed => s3_error!(MethodNotAllowed),
StorageError::BucketNotFound(bucket) => {
s3_error!(InvalidArgument, "bucket not found {}", bucket)
}
StorageError::BucketNotEmpty(bucket) => s3_error!(BucketNotEmpty, "bucket not empty {}", bucket),
StorageError::BucketNameInvalid(bucket) => s3_error!(InvalidBucketName, "invalid bucket name {}", bucket),
StorageError::ObjectNameInvalid(bucket, object) => {
s3_error!(InvalidArgument, "invalid object name {}/{}", bucket, object)
}
StorageError::BucketExists(bucket) => s3_error!(BucketAlreadyExists, "{}", bucket),
StorageError::StorageFull => s3_error!(ServiceUnavailable, "Storage reached its minimum free drive threshold."),
StorageError::SlowDown => s3_error!(SlowDown, "Please reduce your request rate"),
StorageError::PrefixAccessDenied(bucket, object) => {
s3_error!(AccessDenied, "PrefixAccessDenied {}/{}", bucket, object)
}
StorageError::InvalidUploadIDKeyCombination(bucket, object) => {
s3_error!(InvalidArgument, "Invalid UploadID KeyCombination: {}/{}", bucket, object)
}
StorageError::MalformedUploadID(bucket) => s3_error!(InvalidArgument, "Malformed UploadID: {}", bucket),
StorageError::ObjectNameTooLong(bucket, object) => {
s3_error!(InvalidArgument, "Object name too long: {}/{}", bucket, object)
}
StorageError::ObjectNamePrefixAsSlash(bucket, object) => {
s3_error!(InvalidArgument, "Object name contains forward slash as prefix: {}/{}", bucket, object)
}
StorageError::ObjectNotFound(bucket, object) => s3_error!(NoSuchKey, "{}/{}", bucket, object),
StorageError::VersionNotFound(bucket, object, version_id) => {
s3_error!(NoSuchVersion, "{}/{}/{}", bucket, object, version_id)
}
StorageError::InvalidUploadID(bucket, object, version_id) => {
s3_error!(InvalidPart, "Invalid upload id: {}/{}-{}", bucket, object, version_id)
}
StorageError::InvalidVersionID(bucket, object, version_id) => {
s3_error!(InvalidArgument, "Invalid version id: {}/{}-{}", bucket, object, version_id)
}
// extended
StorageError::DataMovementOverwriteErr(bucket, object, version_id) => s3_error!(
InvalidArgument,
"invalid data movement operation, source and destination pool are the same for : {}/{}-{}",
bucket,
object,
version_id
),
// extended
StorageError::ObjectExistsAsDirectory(bucket, object) => {
s3_error!(InvalidArgument, "Object exists on :{} as directory {}", bucket, object)
}
StorageError::InsufficientReadQuorum => {
s3_error!(SlowDown, "Storage resources are insufficient for the read operation")
}
StorageError::InsufficientWriteQuorum => {
s3_error!(SlowDown, "Storage resources are insufficient for the write operation")
}
};
}
if is_err_file_not_found(&err) {
return S3Error::with_message(S3ErrorCode::NoSuchKey, format!(" ec err {}", err));
}
S3Error::with_message(S3ErrorCode::InternalError, format!(" ec err {}", err))
}
+1
View File
@@ -1,2 +1,3 @@
pub mod acess;
pub mod ecfs;
pub mod error;