mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-24 13:16:28 +00:00
+1
-1
@@ -23,7 +23,7 @@ pub mod bucket;
|
|||||||
pub mod file_meta_inline;
|
pub mod file_meta_inline;
|
||||||
pub mod options;
|
pub mod options;
|
||||||
pub mod pools;
|
pub mod pools;
|
||||||
pub(crate) mod store_err;
|
pub mod store_err;
|
||||||
pub mod xhttp;
|
pub mod xhttp;
|
||||||
|
|
||||||
pub use global::is_legacy;
|
pub use global::is_legacy;
|
||||||
|
|||||||
+127
-97
@@ -15,9 +15,7 @@ use ecstore::bucket::policy_sys::PolicySys;
|
|||||||
use ecstore::bucket::tagging::decode_tags;
|
use ecstore::bucket::tagging::decode_tags;
|
||||||
use ecstore::bucket::tagging::encode_tags;
|
use ecstore::bucket::tagging::encode_tags;
|
||||||
use ecstore::bucket::versioning_sys::BucketVersioningSys;
|
use ecstore::bucket::versioning_sys::BucketVersioningSys;
|
||||||
use ecstore::disk::error::is_err_file_not_found;
|
|
||||||
use ecstore::disk::error::DiskError;
|
use ecstore::disk::error::DiskError;
|
||||||
use ecstore::error::Error as EcError;
|
|
||||||
use ecstore::new_object_layer_fn;
|
use ecstore::new_object_layer_fn;
|
||||||
use ecstore::options::extract_metadata;
|
use ecstore::options::extract_metadata;
|
||||||
use ecstore::options::put_opts;
|
use ecstore::options::put_opts;
|
||||||
@@ -54,6 +52,8 @@ use tracing::info;
|
|||||||
use transform_stream::AsyncTryStream;
|
use transform_stream::AsyncTryStream;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
|
use crate::storage::error::to_s3_error;
|
||||||
|
|
||||||
macro_rules! try_ {
|
macro_rules! try_ {
|
||||||
($result:expr) => {
|
($result:expr) => {
|
||||||
match $result {
|
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)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct FS {
|
pub struct FS {
|
||||||
// pub store: ECStore,
|
// pub store: ECStore,
|
||||||
@@ -112,18 +104,17 @@ impl S3 for FS {
|
|||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
||||||
};
|
};
|
||||||
|
|
||||||
try_!(
|
store
|
||||||
store
|
.make_bucket(
|
||||||
.make_bucket(
|
&bucket,
|
||||||
&bucket,
|
&MakeBucketOptions {
|
||||||
&MakeBucketOptions {
|
force_create: true,
|
||||||
force_create: true,
|
lock_enabled: object_lock_enabled_for_bucket.is_some_and(|v| v),
|
||||||
lock_enabled: object_lock_enabled_for_bucket.is_some_and(|v| v),
|
..Default::default()
|
||||||
..Default::default()
|
},
|
||||||
}
|
)
|
||||||
)
|
.await
|
||||||
.await
|
.map_err(to_s3_error)?;
|
||||||
);
|
|
||||||
|
|
||||||
let output = CreateBucketOutput::default();
|
let output = CreateBucketOutput::default();
|
||||||
Ok(S3Response::new(output))
|
Ok(S3Response::new(output))
|
||||||
@@ -151,17 +142,17 @@ impl S3 for FS {
|
|||||||
Some(s) => s,
|
Some(s) => s,
|
||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
||||||
};
|
};
|
||||||
try_!(
|
|
||||||
store
|
store
|
||||||
.delete_bucket(
|
.delete_bucket(
|
||||||
&input.bucket,
|
&input.bucket,
|
||||||
&DeleteBucketOptions {
|
&DeleteBucketOptions {
|
||||||
force: false,
|
force: false,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
}
|
},
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
);
|
.map_err(to_s3_error)?;
|
||||||
|
|
||||||
Ok(S3Response::new(DeleteBucketOutput {}))
|
Ok(S3Response::new(DeleteBucketOutput {}))
|
||||||
}
|
}
|
||||||
@@ -192,7 +183,10 @@ impl S3 for FS {
|
|||||||
Some(s) => s,
|
Some(s) => s,
|
||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
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;
|
// TODO: let errors;
|
||||||
|
|
||||||
@@ -260,7 +254,10 @@ impl S3 for FS {
|
|||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
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);
|
// info!("delete_objects res {:?} {:?}", &dobjs, errs);
|
||||||
|
|
||||||
let deleted = dobjs
|
let deleted = dobjs
|
||||||
@@ -335,7 +332,10 @@ impl S3 for FS {
|
|||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
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;
|
let info = reader.object_info;
|
||||||
|
|
||||||
@@ -363,7 +363,6 @@ impl S3 for FS {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
debug!("get_object response {:?}", output);
|
|
||||||
Ok(S3Response::new(output))
|
Ok(S3Response::new(output))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -406,7 +405,6 @@ impl S3 for FS {
|
|||||||
.get_object_info(&bucket, &key, &ObjectOptions::default())
|
.get_object_info(&bucket, &key, &ObjectOptions::default())
|
||||||
.await
|
.await
|
||||||
.map_err(to_s3_error)?;
|
.map_err(to_s3_error)?;
|
||||||
debug!("info {:?}", info);
|
|
||||||
|
|
||||||
let content_type = {
|
let content_type = {
|
||||||
if let Some(content_type) = info.content_type {
|
if let Some(content_type) = info.content_type {
|
||||||
@@ -445,7 +443,7 @@ impl S3 for FS {
|
|||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
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
|
let buckets: Vec<Bucket> = bucket_infos
|
||||||
.iter()
|
.iter()
|
||||||
@@ -503,19 +501,18 @@ impl S3 for FS {
|
|||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
||||||
};
|
};
|
||||||
|
|
||||||
let object_infos = try_!(
|
let object_infos = store
|
||||||
store
|
.list_objects_v2(
|
||||||
.list_objects_v2(
|
&bucket,
|
||||||
&bucket,
|
&prefix,
|
||||||
&prefix,
|
&continuation_token.unwrap_or_default(),
|
||||||
&continuation_token.unwrap_or_default(),
|
&delimiter,
|
||||||
&delimiter,
|
max_keys.unwrap_or_default(),
|
||||||
max_keys.unwrap_or_default(),
|
fetch_owner.unwrap_or_default(),
|
||||||
fetch_owner.unwrap_or_default(),
|
&start_after.unwrap_or_default(),
|
||||||
&start_after.unwrap_or_default()
|
)
|
||||||
)
|
.await
|
||||||
.await
|
.map_err(to_s3_error)?;
|
||||||
);
|
|
||||||
|
|
||||||
// warn!("object_infos {:?}", object_infos);
|
// warn!("object_infos {:?}", object_infos);
|
||||||
|
|
||||||
@@ -614,9 +611,14 @@ impl S3 for FS {
|
|||||||
metadata.insert(xhttp::AMZ_OBJECT_TAGGING.to_owned(), tags);
|
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;
|
let e_tag = obj_info.etag;
|
||||||
|
|
||||||
@@ -655,9 +657,12 @@ impl S3 for FS {
|
|||||||
metadata.insert(xhttp::AMZ_OBJECT_TAGGING.to_owned(), tags);
|
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 {
|
let output = CreateMultipartUploadOutput {
|
||||||
bucket: Some(bucket),
|
bucket: Some(bucket),
|
||||||
@@ -714,11 +719,10 @@ impl S3 for FS {
|
|||||||
|
|
||||||
// TODO: hash_reader
|
// TODO: hash_reader
|
||||||
|
|
||||||
let info = try_!(
|
let info = store
|
||||||
store
|
.put_object_part(&bucket, &key, &upload_id, part_id, &mut data, &opts)
|
||||||
.put_object_part(&bucket, &key, &upload_id, part_id, &mut data, &opts)
|
.await
|
||||||
.await
|
.map_err(to_s3_error)?;
|
||||||
);
|
|
||||||
|
|
||||||
let output = UploadPartOutput {
|
let output = UploadPartOutput {
|
||||||
e_tag: info.etag,
|
e_tag: info.etag,
|
||||||
@@ -784,11 +788,10 @@ impl S3 for FS {
|
|||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
|
||||||
};
|
};
|
||||||
|
|
||||||
try_!(
|
store
|
||||||
store
|
.complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, opts)
|
||||||
.complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, opts)
|
.await
|
||||||
.await
|
.map_err(to_s3_error)?;
|
||||||
);
|
|
||||||
|
|
||||||
let output = CompleteMultipartUploadOutput {
|
let output = CompleteMultipartUploadOutput {
|
||||||
bucket: Some(bucket),
|
bucket: Some(bucket),
|
||||||
@@ -815,11 +818,11 @@ impl S3 for FS {
|
|||||||
};
|
};
|
||||||
|
|
||||||
let opts = &ObjectOptions::default();
|
let opts = &ObjectOptions::default();
|
||||||
try_!(
|
|
||||||
store
|
store
|
||||||
.abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts)
|
.abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts)
|
||||||
.await
|
.await
|
||||||
);
|
.map_err(to_s3_error)?;
|
||||||
Ok(S3Response::new(AbortMultipartUploadOutput { ..Default::default() }))
|
Ok(S3Response::new(AbortMultipartUploadOutput { ..Default::default() }))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -849,7 +852,6 @@ impl S3 for FS {
|
|||||||
#[tracing::instrument(level = "debug", skip(self))]
|
#[tracing::instrument(level = "debug", skip(self))]
|
||||||
async fn put_bucket_tagging(&self, req: S3Request<PutBucketTaggingInput>) -> S3Result<S3Response<PutBucketTaggingOutput>> {
|
async fn put_bucket_tagging(&self, req: S3Request<PutBucketTaggingInput>) -> S3Result<S3Response<PutBucketTaggingOutput>> {
|
||||||
let PutBucketTaggingInput { bucket, tagging, .. } = req.input;
|
let PutBucketTaggingInput { bucket, tagging, .. } = req.input;
|
||||||
log::debug!("bucket: {bucket}, tagging: {tagging:?}");
|
|
||||||
|
|
||||||
let layer = new_object_layer_fn();
|
let layer = new_object_layer_fn();
|
||||||
let lock = layer.read().await;
|
let lock = layer.read().await;
|
||||||
@@ -868,7 +870,9 @@ impl S3 for FS {
|
|||||||
|
|
||||||
let data = try_!(xml::serialize(&tagging));
|
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()))
|
Ok(S3Response::new(Default::default()))
|
||||||
}
|
}
|
||||||
@@ -880,7 +884,9 @@ impl S3 for FS {
|
|||||||
) -> S3Result<S3Response<DeleteBucketTaggingOutput>> {
|
) -> S3Result<S3Response<DeleteBucketTaggingOutput>> {
|
||||||
let DeleteBucketTaggingInput { bucket, .. } = req.input;
|
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 {}))
|
Ok(S3Response::new(DeleteBucketTaggingOutput {}))
|
||||||
}
|
}
|
||||||
@@ -900,17 +906,15 @@ impl S3 for FS {
|
|||||||
.as_ref()
|
.as_ref()
|
||||||
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
|
.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);
|
let tags = encode_tags(tagging.tag_set);
|
||||||
|
|
||||||
// TODO: getOpts
|
// TODO: getOpts
|
||||||
// TODO: Replicate
|
// TODO: Replicate
|
||||||
try_!(
|
|
||||||
store
|
store
|
||||||
.put_object_tags(&bucket, &object, &tags, &ObjectOptions::default())
|
.put_object_tags(&bucket, &object, &tags, &ObjectOptions::default())
|
||||||
.await
|
.await
|
||||||
);
|
.map_err(to_s3_error)?;
|
||||||
|
|
||||||
Ok(S3Response::new(PutObjectTaggingOutput { version_id: None }))
|
Ok(S3Response::new(PutObjectTaggingOutput { version_id: None }))
|
||||||
}
|
}
|
||||||
@@ -926,7 +930,10 @@ impl S3 for FS {
|
|||||||
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
|
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
|
||||||
|
|
||||||
// TODO: version
|
// 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());
|
let tag_set = decode_tags(tags.as_str());
|
||||||
|
|
||||||
@@ -951,7 +958,10 @@ impl S3 for FS {
|
|||||||
|
|
||||||
// TODO: Replicate
|
// TODO: Replicate
|
||||||
// TODO: version
|
// 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 }))
|
Ok(S3Response::new(DeleteObjectTaggingOutput { version_id: None }))
|
||||||
}
|
}
|
||||||
@@ -977,7 +987,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 {
|
Ok(S3Response::new(GetBucketVersioningOutput {
|
||||||
status,
|
status,
|
||||||
@@ -1003,7 +1013,9 @@ impl S3 for FS {
|
|||||||
|
|
||||||
let data = try_!(xml::serialize(&versioning_configuration));
|
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
|
// TODO: globalSiteReplicationSys.BucketMetaHook
|
||||||
|
|
||||||
@@ -1068,7 +1080,7 @@ impl S3 for FS {
|
|||||||
|
|
||||||
// warn!("input policy {}", &policy);
|
// 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);
|
// warn!("parse policy {:?}", &cfg);
|
||||||
|
|
||||||
@@ -1077,9 +1089,11 @@ impl S3 for FS {
|
|||||||
return Err(s3_error!(InvalidPolicyDocument));
|
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 {}))
|
Ok(S3Response::new(PutBucketPolicyOutput {}))
|
||||||
}
|
}
|
||||||
@@ -1104,7 +1118,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 {}))
|
Ok(S3Response::new(DeleteBucketPolicyOutput {}))
|
||||||
}
|
}
|
||||||
@@ -1165,7 +1181,9 @@ impl S3 for FS {
|
|||||||
let Some(input_cfg) = lifecycle_configuration else { return Err(s3_error!(InvalidArgument)) };
|
let Some(input_cfg) = lifecycle_configuration else { return Err(s3_error!(InvalidArgument)) };
|
||||||
|
|
||||||
let data = try_!(xml::serialize(&input_cfg));
|
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()))
|
Ok(S3Response::new(PutBucketLifecycleConfigurationOutput::default()))
|
||||||
}
|
}
|
||||||
@@ -1191,7 +1209,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()))
|
Ok(S3Response::new(DeleteBucketLifecycleOutput::default()))
|
||||||
}
|
}
|
||||||
@@ -1261,7 +1281,9 @@ impl S3 for FS {
|
|||||||
// TODO: check kms
|
// TODO: check kms
|
||||||
|
|
||||||
let data = try_!(xml::serialize(&server_side_encryption_configuration));
|
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()))
|
Ok(S3Response::new(PutBucketEncryptionOutput::default()))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1284,7 +1306,7 @@ impl S3 for FS {
|
|||||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
|
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()))
|
Ok(S3Response::new(DeleteBucketEncryptionOutput::default()))
|
||||||
}
|
}
|
||||||
@@ -1340,7 +1362,9 @@ impl S3 for FS {
|
|||||||
|
|
||||||
let data = try_!(xml::serialize(&input_cfg));
|
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()))
|
Ok(S3Response::new(PutObjectLockConfigurationOutput::default()))
|
||||||
}
|
}
|
||||||
@@ -1405,7 +1429,9 @@ impl S3 for FS {
|
|||||||
// TODO: check enable, versioning enable
|
// TODO: check enable, versioning enable
|
||||||
let data = try_!(xml::serialize(&replication_configuration));
|
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()))
|
Ok(S3Response::new(PutBucketReplicationOutput::default()))
|
||||||
}
|
}
|
||||||
@@ -1429,7 +1455,9 @@ impl S3 for FS {
|
|||||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
|
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
|
// TODO: remove targets
|
||||||
|
|
||||||
@@ -1510,7 +1538,9 @@ impl S3 for FS {
|
|||||||
|
|
||||||
let data = try_!(xml::serialize(¬ification_configuration));
|
let data = try_!(xml::serialize(¬ification_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
|
// TODO: event notice add rule
|
||||||
|
|
||||||
|
|||||||
@@ -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,2 +1,3 @@
|
|||||||
pub mod acess;
|
pub mod acess;
|
||||||
pub mod ecfs;
|
pub mod ecfs;
|
||||||
|
pub mod error;
|
||||||
|
|||||||
+1
-1
@@ -7,7 +7,7 @@ mkdir -p ./target/volume/test{0..4}
|
|||||||
|
|
||||||
|
|
||||||
if [ -z "$RUST_LOG" ]; then
|
if [ -z "$RUST_LOG" ]; then
|
||||||
export RUST_LOG="rustfs=debug,ecstore=debug,s3s=debug,reader=debug,router=debug"
|
export RUST_LOG="rustfs=debug,ecstore=debug,s3s=info,reader=debug,router=debug"
|
||||||
fi
|
fi
|
||||||
|
|
||||||
# export RUSTFS_ERASURE_SET_DRIVE_COUNT=5
|
# export RUSTFS_ERASURE_SET_DRIVE_COUNT=5
|
||||||
|
|||||||
Reference in New Issue
Block a user