fix:#38 implement the basic storage functions of bucketmeta config use s3s struct define

This commit is contained in:
weisd
2024-10-14 14:58:32 +08:00
parent 0f39887cdb
commit 12114dce16
40 changed files with 920 additions and 1234 deletions
+1 -1
View File
@@ -59,7 +59,7 @@ futures.workspace = true
futures-util.workspace = true
# uuid = { version = "1.8.0", features = ["v4", "fast-rng", "serde"] }
ecstore = { path = "../ecstore" }
s3s = "0.10.0"
s3s.workspace = true
clap = { version = "4.5.20", features = ["derive"] }
tracing-subscriber = { version = "0.3.18", features = ["env-filter", "time"] }
hyper-util = { version = "0.1.9", features = [
+6 -6
View File
@@ -6,7 +6,7 @@ mod storage;
use clap::Parser;
use common::error::{Error, Result};
use ecstore::{
bucket::init_bucket_metadata_sys,
bucket::metadata_sys::init_bucket_metadata_sys,
endpoints::EndpointServerPools,
set_global_endpoints,
store::{init_local_disks, ECStore},
@@ -128,11 +128,11 @@ async fn run(opt: config::Opt) -> Result<()> {
info!("authentication is enabled {}, {}", &access_key, &secret_key);
b.set_auth(SimpleAuth::from_single(access_key, secret_key));
// Enable parsing virtual-hosted-style requests
if let Some(dm) = opt.domain_name {
info!("virtual-hosted-style requests are enabled use domain_name {}", &dm);
b.set_base_domain(dm);
}
// // Enable parsing virtual-hosted-style requests
// if let Some(dm) = opt.domain_name {
// info!("virtual-hosted-style requests are enabled use domain_name {}", &dm);
// b.set_base_domain(dm);
// }
// if domain_name.is_some() {
// info!(
+399 -145
View File
@@ -1,14 +1,17 @@
use bytes::Bytes;
use ecstore::bucket::error::BucketMetadataError;
use ecstore::bucket::get_bucket_metadata_sys;
use ecstore::bucket::metadata;
use ecstore::bucket::metadata::BUCKET_LIFECYCLE_CONFIG;
use ecstore::bucket::metadata::BUCKET_NOTIFICATION_CONFIG;
use ecstore::bucket::metadata::BUCKET_POLICY_CONFIG;
use ecstore::bucket::metadata::BUCKET_REPLICATION_CONFIG;
use ecstore::bucket::metadata::BUCKET_SSECONFIG;
use ecstore::bucket::metadata::BUCKET_TAGGING_CONFIG;
use ecstore::bucket::metadata::BUCKET_VERSIONING_CONFIG;
use ecstore::bucket::metadata::OBJECT_LOCK_CONFIG;
use ecstore::bucket::metadata_sys;
use ecstore::bucket::policy::bucket_policy::BucketPolicy;
use ecstore::bucket::policy_sys::PolicySys;
use ecstore::bucket::tags::Tags;
use ecstore::bucket::versioning::State as VersioningState;
use ecstore::bucket::versioning::Versioning;
use ecstore::bucket::versioning_sys::BucketVersioningSys;
use ecstore::disk::error::DiskError;
use ecstore::new_object_layer_fn;
@@ -34,9 +37,9 @@ use s3s::S3ErrorCode;
use s3s::S3Result;
use s3s::S3;
use s3s::{S3Request, S3Response};
use std::collections::HashMap;
use std::fmt::Debug;
use std::str::FromStr;
use tracing::info;
use transform_stream::AsyncTryStream;
use uuid::Uuid;
@@ -73,7 +76,11 @@ impl S3 for FS {
fields(start_time=?time::OffsetDateTime::now_utc())
)]
async fn create_bucket(&self, req: S3Request<CreateBucketInput>) -> S3Result<S3Response<CreateBucketOutput>> {
let input = req.input;
let CreateBucketInput {
bucket,
object_lock_enabled_for_bucket,
..
} = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
@@ -84,7 +91,14 @@ impl S3 for FS {
try_!(
store
.make_bucket(&input.bucket, &MakeBucketOptions { force_create: true })
.make_bucket(
&bucket,
&MakeBucketOptions {
force_create: true,
lock_enabled: object_lock_enabled_for_bucket.is_some_and(|v| v),
..Default::default()
}
)
.await
);
@@ -690,36 +704,6 @@ impl S3 for FS {
Ok(S3Response::new(AbortMultipartUploadOutput { ..Default::default() }))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn put_bucket_tagging(&self, req: S3Request<PutBucketTaggingInput>) -> S3Result<S3Response<PutBucketTaggingOutput>> {
let PutBucketTaggingInput { bucket, tagging, .. } = req.input;
log::debug!("bucket: {bucket}, tagging: {tagging:?}");
// check bucket exists.
let _bucket = self
.head_bucket(S3Request::new(HeadBucketInput {
bucket: bucket.clone(),
expected_bucket_owner: None,
}))
.await?;
let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
let mut bucket_meta_sys = bucket_meta_sys_lock.write().await;
let mut tag_map = HashMap::new();
for tag in tagging.tag_set.iter() {
tag_map.insert(tag.key.clone(), tag.value.clone());
}
let tags = Tags::new(tag_map, false);
let data = try_!(tags.marshal_msg());
let _updated = try_!(bucket_meta_sys.update(&bucket, BUCKET_TAGGING_CONFIG, data).await);
Ok(S3Response::new(Default::default()))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn get_bucket_tagging(&self, req: S3Request<GetBucketTaggingInput>) -> S3Result<S3Response<GetBucketTaggingOutput>> {
let GetBucketTaggingInput { bucket, .. } = req.input;
@@ -731,49 +715,53 @@ impl S3 for FS {
}))
.await?;
let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
let bucket_meta_sys = bucket_meta_sys_lock.read().await;
let tag_set: Vec<Tag> = match bucket_meta_sys.get_tagging_config(&bucket).await {
Ok((tags, _)) => tags
.tag_set
.tag_map
.into_iter()
.map(|(key, value)| Tag { key, value })
.collect(),
let Tagging { tag_set } = match metadata_sys::get_tagging_config(&bucket).await {
Ok((tags, _)) => tags,
Err(err) => {
warn!("get_tagging_config err {:?}", &err);
// TODO: check not found
Vec::new()
Tagging::default()
}
};
Ok(S3Response::new(GetBucketTaggingOutput { tag_set }))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn put_bucket_tagging(&self, req: S3Request<PutBucketTaggingInput>) -> S3Result<S3Response<PutBucketTaggingOutput>> {
let PutBucketTaggingInput { bucket, tagging, .. } = req.input;
log::debug!("bucket: {bucket}, tagging: {tagging:?}");
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = match lock.as_ref() {
Some(s) => s,
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())),
};
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
let data = try_!(metadata::serialize(&tagging));
try_!(metadata_sys::update(&bucket, BUCKET_TAGGING_CONFIG, data).await);
Ok(S3Response::new(Default::default()))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn delete_bucket_tagging(
&self,
req: S3Request<DeleteBucketTaggingInput>,
) -> S3Result<S3Response<DeleteBucketTaggingOutput>> {
let DeleteBucketTaggingInput { bucket, .. } = req.input;
// check bucket exists.
let _bucket = self
.head_bucket(S3Request::new(HeadBucketInput {
bucket: bucket.clone(),
expected_bucket_owner: None,
}))
.await?;
let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
let mut bucket_meta_sys = bucket_meta_sys_lock.write().await;
let tag_map = HashMap::new();
let tags = Tags::new(tag_map, false);
let data = try_!(tags.marshal_msg());
let _updated = try_!(bucket_meta_sys.update(&bucket, BUCKET_TAGGING_CONFIG, data).await);
try_!(metadata_sys::delete(&bucket, BUCKET_TAGGING_CONFIG).await);
Ok(S3Response::new(DeleteBucketTaggingOutput {}))
}
@@ -872,16 +860,11 @@ impl S3 for FS {
}
}
let cfg = try_!(BucketVersioningSys::get(&bucket).await);
let status = match cfg.status {
VersioningState::Enabled => Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)),
VersioningState::Suspended => Some(BucketVersioningStatus::from_static(BucketVersioningStatus::SUSPENDED)),
};
let VersioningConfiguration { status, .. } = try_!(BucketVersioningSys::get(&bucket).await);
Ok(S3Response::new(GetBucketVersioningOutput {
mfa_delete: None,
status,
..Default::default()
}))
}
@@ -901,28 +884,9 @@ impl S3 for FS {
// check bucket object lock enable
// check replication suspended
let mut cfg = match BucketVersioningSys::get(&bucket).await {
Ok(res) => res,
Err(err) => {
warn!("BucketVersioningSys::get err {:?}", err);
Versioning::default()
}
};
let data = try_!(metadata::serialize(&versioning_configuration));
if let Some(verstatus) = versioning_configuration.status {
cfg.status = match verstatus.as_str() {
BucketVersioningStatus::ENABLED => VersioningState::Enabled,
BucketVersioningStatus::SUSPENDED => VersioningState::Suspended,
_ => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init")),
}
}
let data = try_!(cfg.marshal_msg());
let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
let mut bucket_meta_sys = bucket_meta_sys_lock.write().await;
try_!(bucket_meta_sys.update(&bucket, BUCKET_VERSIONING_CONFIG, data).await);
try_!(metadata_sys::update(&bucket, BUCKET_VERSIONING_CONFIG, data).await);
// TODO: globalSiteReplicationSys.BucketMetaHook
@@ -963,7 +927,7 @@ impl S3 for FS {
}
};
let policys = try_!(cfg.marshal_msg());
let policys = try_!(serde_json::to_string(&cfg));
Ok(S3Response::new(GetBucketPolicyOutput { policy: Some(policys) }))
}
@@ -986,20 +950,15 @@ impl S3 for FS {
}
let cfg = try_!(BucketPolicy::unmarshal(policy.as_bytes()));
warn!("put_bucket_policy struct {:?}", &cfg);
if let Err(err) = cfg.validate(&bucket) {
warn!("put_bucket_policy input {:?}", &policy);
warn!("cfg.validate err {:?}", err);
if let Err(err) = cfg.validate(&bucket) {
warn!("put_bucket_policy err input {:?}, {:?}", &policy, err);
return Err(s3_error!(InvalidPolicyDocument));
}
let data = try_!(cfg.marshal_msg());
let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
let mut bucket_meta_sys = bucket_meta_sys_lock.write().await;
try_!(bucket_meta_sys.update(&bucket, BUCKET_POLICY_CONFIG, data.into()).await);
try_!(metadata_sys::update(&bucket, BUCKET_POLICY_CONFIG, data.into()).await);
Ok(S3Response::new(PutBucketPolicyOutput {}))
}
@@ -1024,10 +983,7 @@ impl S3 for FS {
}
}
let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
let mut bucket_meta_sys = bucket_meta_sys_lock.write().await;
try_!(bucket_meta_sys.delete(&bucket, BUCKET_POLICY_CONFIG).await);
try_!(metadata_sys::delete(&bucket, BUCKET_POLICY_CONFIG).await);
Ok(S3Response::new(DeleteBucketPolicyOutput {}))
}
@@ -1035,99 +991,409 @@ impl S3 for FS {
#[tracing::instrument(level = "debug", skip(self))]
async fn get_bucket_lifecycle_configuration(
&self,
_req: S3Request<GetBucketLifecycleConfigurationInput>,
req: S3Request<GetBucketLifecycleConfigurationInput>,
) -> S3Result<S3Response<GetBucketLifecycleConfigurationOutput>> {
Err(s3_error!(NotImplemented, "GetBucketLifecycleConfiguration is not implemented yet"))
let GetBucketLifecycleConfigurationInput { bucket, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
let rules = match metadata_sys::get_lifecycle_config(&bucket).await {
Ok((cfg, _)) => Some(cfg.rules),
Err(_err) => {
// if BucketMetadataError::BucketLifecycleNotFound.is(&err) {
// return Err(s3_error!(NoSuchLifecycleConfiguration));
// }
// warn!("get_lifecycle_config err {:?}", err);
None
}
};
Ok(S3Response::new(GetBucketLifecycleConfigurationOutput {
rules,
..Default::default()
}))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn put_bucket_lifecycle_configuration(
&self,
_req: S3Request<PutBucketLifecycleConfigurationInput>,
req: S3Request<PutBucketLifecycleConfigurationInput>,
) -> S3Result<S3Response<PutBucketLifecycleConfigurationOutput>> {
Err(s3_error!(NotImplemented, "PutBucketLifecycleConfiguration is not implemented yet"))
let PutBucketLifecycleConfigurationInput {
bucket,
lifecycle_configuration,
..
} = req.input;
warn!("lifecycle_configuration {:?}", &lifecycle_configuration);
// TODO: objcetLock
let Some(input_cfg) = lifecycle_configuration else { return Err(s3_error!(InvalidArgument)) };
let data = try_!(metadata::serialize(&input_cfg));
try_!(metadata_sys::update(&bucket, BUCKET_LIFECYCLE_CONFIG, data).await);
Ok(S3Response::new(PutBucketLifecycleConfigurationOutput::default()))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn delete_bucket_lifecycle(
&self,
_req: S3Request<DeleteBucketLifecycleInput>,
req: S3Request<DeleteBucketLifecycleInput>,
) -> S3Result<S3Response<DeleteBucketLifecycleOutput>> {
Err(s3_error!(NotImplemented, "DeleteBucketLifecycle is not implemented yet"))
let DeleteBucketLifecycleInput { bucket, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
try_!(metadata_sys::delete(&bucket, BUCKET_LIFECYCLE_CONFIG).await);
Ok(S3Response::new(DeleteBucketLifecycleOutput::default()))
}
async fn get_bucket_encryption(
&self,
_req: S3Request<GetBucketEncryptionInput>,
req: S3Request<GetBucketEncryptionInput>,
) -> S3Result<S3Response<GetBucketEncryptionOutput>> {
Err(s3_error!(NotImplemented, "GetBucketEncryption is not implemented yet"))
let GetBucketEncryptionInput { bucket, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
let server_side_encryption_configuration = match metadata_sys::get_sse_config(&bucket).await {
Ok((cfg, _)) => Some(cfg),
Err(err) => {
// if BucketMetadataError::BucketLifecycleNotFound.is(&err) {
// return Err(s3_error!(ErrNoSuchBucketSSEConfig));
// }
warn!("get_sse_config err {:?}", err);
None
}
};
Ok(S3Response::new(GetBucketEncryptionOutput {
server_side_encryption_configuration,
}))
}
async fn put_bucket_encryption(
&self,
_req: S3Request<PutBucketEncryptionInput>,
req: S3Request<PutBucketEncryptionInput>,
) -> S3Result<S3Response<PutBucketEncryptionOutput>> {
Err(s3_error!(NotImplemented, "PutBucketEncryption is not implemented yet"))
let PutBucketEncryptionInput {
bucket,
server_side_encryption_configuration,
..
} = req.input;
info!("sse_config {:?}", &server_side_encryption_configuration);
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
// TODO: check kms
let data = try_!(metadata::serialize(&server_side_encryption_configuration));
try_!(metadata_sys::update(&bucket, BUCKET_SSECONFIG, data).await);
Ok(S3Response::new(PutBucketEncryptionOutput::default()))
}
async fn delete_bucket_encryption(
&self,
_req: S3Request<DeleteBucketEncryptionInput>,
req: S3Request<DeleteBucketEncryptionInput>,
) -> S3Result<S3Response<DeleteBucketEncryptionOutput>> {
Err(s3_error!(NotImplemented, "DeleteBucketEncryption is not implemented yet"))
let DeleteBucketEncryptionInput { bucket, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
try_!(metadata_sys::delete(&bucket, BUCKET_SSECONFIG).await);
Ok(S3Response::new(DeleteBucketEncryptionOutput::default()))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn get_object_lock_configuration(
&self,
_req: S3Request<GetObjectLockConfigurationInput>,
req: S3Request<GetObjectLockConfigurationInput>,
) -> S3Result<S3Response<GetObjectLockConfigurationOutput>> {
// mc cp step 1
let output = GetObjectLockConfigurationOutput::default();
Ok(S3Response::new(output))
let GetObjectLockConfigurationInput { bucket, .. } = req.input;
let object_lock_configuration = match metadata_sys::get_object_lock_config(&bucket).await {
Ok((cfg, _created)) => Some(cfg),
Err(err) => {
warn!("get_object_lock_config err {:?}", err);
None
}
};
warn!("object_lock_configuration {:?}", &object_lock_configuration);
Ok(S3Response::new(GetObjectLockConfigurationOutput {
object_lock_configuration,
}))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn put_object_lock_configuration(
&self,
_req: S3Request<PutObjectLockConfigurationInput>,
req: S3Request<PutObjectLockConfigurationInput>,
) -> S3Result<S3Response<PutObjectLockConfigurationOutput>> {
Err(s3_error!(NotImplemented, "PutObjectLockConfiguration is not implemented yet"))
let PutObjectLockConfigurationInput {
bucket,
object_lock_configuration,
..
} = req.input;
let Some(input_cfg) = object_lock_configuration else { return Err(s3_error!(InvalidArgument)) };
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
let data = try_!(metadata::serialize(&input_cfg));
try_!(metadata_sys::update(&bucket, OBJECT_LOCK_CONFIG, data).await);
Ok(S3Response::new(PutObjectLockConfigurationOutput::default()))
}
async fn get_bucket_replication(
&self,
_req: S3Request<GetBucketReplicationInput>,
req: S3Request<GetBucketReplicationInput>,
) -> S3Result<S3Response<GetBucketReplicationOutput>> {
Err(s3_error!(NotImplemented, "GetBucketReplication is not implemented yet"))
let GetBucketReplicationInput { bucket, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
let replication_configuration = match metadata_sys::get_replication_config(&bucket).await {
Ok((cfg, _created)) => Some(cfg),
Err(err) => {
warn!("get_object_lock_config err {:?}", err);
None
}
};
Ok(S3Response::new(GetBucketReplicationOutput {
replication_configuration,
}))
}
async fn put_bucket_replication(
&self,
_req: S3Request<PutBucketReplicationInput>,
req: S3Request<PutBucketReplicationInput>,
) -> S3Result<S3Response<PutBucketReplicationOutput>> {
Err(s3_error!(NotImplemented, "PutBucketReplication is not implemented yet"))
let PutBucketReplicationInput {
bucket,
replication_configuration,
..
} = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
// TODO: check enable, versioning enable
let data = try_!(metadata::serialize(&replication_configuration));
try_!(metadata_sys::update(&bucket, BUCKET_REPLICATION_CONFIG, data).await);
Ok(S3Response::new(PutBucketReplicationOutput::default()))
}
async fn delete_bucket_replication(
&self,
_req: S3Request<DeleteBucketReplicationInput>,
req: S3Request<DeleteBucketReplicationInput>,
) -> S3Result<S3Response<DeleteBucketReplicationOutput>> {
Err(s3_error!(NotImplemented, "DeleteBucketReplication is not implemented yet"))
let DeleteBucketReplicationInput { bucket, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
try_!(metadata_sys::delete(&bucket, BUCKET_REPLICATION_CONFIG).await);
// TODO: remove targets
Ok(S3Response::new(DeleteBucketReplicationOutput::default()))
}
async fn get_bucket_notification_configuration(
&self,
_req: S3Request<GetBucketNotificationConfigurationInput>,
req: S3Request<GetBucketNotificationConfigurationInput>,
) -> S3Result<S3Response<GetBucketNotificationConfigurationOutput>> {
Err(s3_error!(NotImplemented, "GetBucketNotificationConfiguration is not implemented yet"))
let GetBucketNotificationConfigurationInput { bucket, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
let has_notification_config = match metadata_sys::get_notification_config(&bucket).await {
Ok(cfg) => cfg,
Err(err) => {
warn!("get_notification_config err {:?}", err);
None
}
};
// TODO: valid target list
if let Some(NotificationConfiguration {
event_bridge_configuration,
lambda_function_configurations,
queue_configurations,
topic_configurations,
}) = has_notification_config
{
Ok(S3Response::new(GetBucketNotificationConfigurationOutput {
event_bridge_configuration,
lambda_function_configurations,
queue_configurations,
topic_configurations,
}))
} else {
Ok(S3Response::new(GetBucketNotificationConfigurationOutput::default()))
}
}
async fn put_bucket_notification_configuration(
&self,
_req: S3Request<PutBucketNotificationConfigurationInput>,
req: S3Request<PutBucketNotificationConfigurationInput>,
) -> S3Result<S3Response<PutBucketNotificationConfigurationOutput>> {
Err(s3_error!(NotImplemented, "PutBucketNotificationConfiguration is not implemented yet"))
let PutBucketNotificationConfigurationInput {
bucket,
notification_configuration,
..
} = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket));
} else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));
}
}
let data = try_!(metadata::serialize(&notification_configuration));
try_!(metadata_sys::update(&bucket, BUCKET_NOTIFICATION_CONFIG, data).await);
// TODO: event notice add rule
Ok(S3Response::new(PutBucketNotificationConfigurationOutput::default()))
}
}
@@ -1151,15 +1417,3 @@ where
Ok(())
})
}
// Consumes this body object to return a bytes stream.
// pub fn into_bytes_stream(mut body: StreamingBlob) -> impl Stream<Item = Result<Bytes, std::io::Error>> + Send + 'static {
// futures_util::stream::poll_fn(move |ctx| loop {
// match Pin::new(&mut body).poll_next(ctx) {
// Poll::Ready(Some(Ok(data))) => return Poll::Ready(Some(Ok(data))),
// Poll::Ready(Some(Err(err))) => return Poll::Ready(Some(Err(std::io::Error::new(std::io::ErrorKind::Other, err)))),
// Poll::Ready(None) => return Poll::Ready(None),
// Poll::Pending => return Poll::Pending,
// }
// })
// }