mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-20 19:42:17 +00:00
get/put BucketVersionin done
This commit is contained in:
@@ -1,4 +1,7 @@
|
|||||||
use crate::error::{Error, Result};
|
use crate::{
|
||||||
|
error::{Error, Result},
|
||||||
|
utils,
|
||||||
|
};
|
||||||
use rmp_serde::Serializer as rmpSerializer;
|
use rmp_serde::Serializer as rmpSerializer;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
@@ -69,7 +72,6 @@ impl Versioning {
|
|||||||
return Err(Error::new(VersioningErr::TooManyExcludedPrefixes));
|
return Err(Error::new(VersioningErr::TooManyExcludedPrefixes));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
_ => return Err(Error::msg(format!("unsupported versioning status {}", self.status))),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
@@ -97,7 +99,7 @@ impl Versioning {
|
|||||||
|
|
||||||
for sprefix in self.excluded_prefixes.iter() {
|
for sprefix in self.excluded_prefixes.iter() {
|
||||||
let full_prefix = format!("{}*", sprefix.prefix);
|
let full_prefix = format!("{}*", sprefix.prefix);
|
||||||
if utils::wildcard::match_simple(full_prefix, prefix) {
|
if utils::wildcard::match_simple(&full_prefix, prefix) {
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -124,7 +126,7 @@ impl Versioning {
|
|||||||
|
|
||||||
for sprefix in self.excluded_prefixes.iter() {
|
for sprefix in self.excluded_prefixes.iter() {
|
||||||
let full_prefix = format!("{}*", sprefix.prefix);
|
let full_prefix = format!("{}*", sprefix.prefix);
|
||||||
if utils::wildcard::match_simple(full_prefix, prefix) {
|
if utils::wildcard::match_simple(&full_prefix, prefix) {
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
use super::get_bucket_metadata_sys;
|
use super::get_bucket_metadata_sys;
|
||||||
use super::versioning::Versioning;
|
use super::versioning::Versioning;
|
||||||
use crate::disk::RUSTFS_META_BUCKET;
|
use crate::disk::RUSTFS_META_BUCKET;
|
||||||
|
use crate::error::Result;
|
||||||
use tracing::warn;
|
use tracing::warn;
|
||||||
|
|
||||||
pub struct BucketVersioningSys {}
|
pub struct BucketVersioningSys {}
|
||||||
@@ -52,7 +53,7 @@ impl BucketVersioningSys {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
|
let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
|
||||||
let mut bucket_meta_sys = bucket_meta_sys_lock.write().await;
|
let bucket_meta_sys = bucket_meta_sys_lock.write().await;
|
||||||
|
|
||||||
let (cfg, _) = bucket_meta_sys.get_versioning_config(bucket).await?;
|
let (cfg, _) = bucket_meta_sys.get_versioning_config(bucket).await?;
|
||||||
|
|
||||||
|
|||||||
@@ -4,4 +4,4 @@ pub mod fs;
|
|||||||
pub mod hash;
|
pub mod hash;
|
||||||
pub mod net;
|
pub mod net;
|
||||||
pub mod path;
|
pub mod path;
|
||||||
mod wildcard;
|
pub mod wildcard;
|
||||||
|
|||||||
@@ -6,7 +6,7 @@ pub fn match_simple(pattern: &str, name: &str) -> bool {
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
// Do an extended wildcard '*' and '?' match.
|
// Do an extended wildcard '*' and '?' match.
|
||||||
deep_match_rune(name, pattern, true)
|
deep_match_rune(name.as_bytes(), pattern.as_bytes(), true)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn match_pattern(pattern: &str, name: &str) -> bool {
|
pub fn match_pattern(pattern: &str, name: &str) -> bool {
|
||||||
@@ -17,11 +17,11 @@ pub fn match_pattern(pattern: &str, name: &str) -> bool {
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
// Do an extended wildcard '*' and '?' match.
|
// Do an extended wildcard '*' and '?' match.
|
||||||
deep_match_rune(name, pattern, false)
|
deep_match_rune(name.as_bytes(), pattern.as_bytes(), false)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn deep_match_rune(str_: &str, pattern: &str, simple: bool) -> bool {
|
fn deep_match_rune(str_: &[u8], pattern: &[u8], simple: bool) -> bool {
|
||||||
let (mut str_, mut pattern) = (str_.as_bytes(), pattern.as_bytes());
|
let (mut str_, mut pattern) = (str_, pattern);
|
||||||
while !pattern.is_empty() {
|
while !pattern.is_empty() {
|
||||||
match pattern[0] as char {
|
match pattern[0] as char {
|
||||||
'*' => {
|
'*' => {
|
||||||
|
|||||||
@@ -1,8 +1,10 @@
|
|||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use ecstore::bucket::get_bucket_metadata_sys;
|
use ecstore::bucket::get_bucket_metadata_sys;
|
||||||
use ecstore::bucket::metadata::BUCKET_TAGGING_CONFIG;
|
use ecstore::bucket::metadata::BUCKET_TAGGING_CONFIG;
|
||||||
|
use ecstore::bucket::metadata::BUCKET_VERSIONING_CONFIG;
|
||||||
use ecstore::bucket::tags::Tags;
|
use ecstore::bucket::tags::Tags;
|
||||||
use ecstore::bucket::versioning::State as VersioningState;
|
use ecstore::bucket::versioning::State as VersioningState;
|
||||||
|
use ecstore::bucket::versioning::Versioning;
|
||||||
use ecstore::bucket::versioning_sys::BucketVersioningSys;
|
use ecstore::bucket::versioning_sys::BucketVersioningSys;
|
||||||
use ecstore::disk::error::DiskError;
|
use ecstore::disk::error::DiskError;
|
||||||
use ecstore::new_object_layer_fn;
|
use ecstore::new_object_layer_fn;
|
||||||
@@ -19,6 +21,7 @@ use ecstore::store_api::PutObjReader;
|
|||||||
use ecstore::store_api::StorageAPI;
|
use ecstore::store_api::StorageAPI;
|
||||||
use futures::pin_mut;
|
use futures::pin_mut;
|
||||||
use futures::{Stream, StreamExt};
|
use futures::{Stream, StreamExt};
|
||||||
|
use http::status;
|
||||||
use http::HeaderMap;
|
use http::HeaderMap;
|
||||||
use log::warn;
|
use log::warn;
|
||||||
use s3s::dto::*;
|
use s3s::dto::*;
|
||||||
@@ -858,7 +861,7 @@ impl S3 for FS {
|
|||||||
&self,
|
&self,
|
||||||
req: S3Request<GetBucketVersioningInput>,
|
req: S3Request<GetBucketVersioningInput>,
|
||||||
) -> S3Result<S3Response<GetBucketVersioningOutput>> {
|
) -> S3Result<S3Response<GetBucketVersioningOutput>> {
|
||||||
let GetBucketVersioningInput { bucket, .. } = req;
|
let GetBucketVersioningInput { bucket, .. } = req.input;
|
||||||
let layer = new_object_layer_fn();
|
let layer = new_object_layer_fn();
|
||||||
let lock = layer.read().await;
|
let lock = layer.read().await;
|
||||||
let store = match lock.as_ref() {
|
let store = match lock.as_ref() {
|
||||||
@@ -866,7 +869,7 @@ impl S3 for FS {
|
|||||||
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))),
|
None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))),
|
||||||
};
|
};
|
||||||
|
|
||||||
if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions::default()).await {
|
if let Err(e) = store.get_bucket_info(&bucket, &BucketOptions::default()).await {
|
||||||
if DiskError::VolumeNotFound.is(&e) {
|
if DiskError::VolumeNotFound.is(&e) {
|
||||||
return Err(s3_error!(NoSuchBucket));
|
return Err(s3_error!(NoSuchBucket));
|
||||||
} else {
|
} else {
|
||||||
@@ -877,8 +880,8 @@ impl S3 for FS {
|
|||||||
let cfg = try_!(BucketVersioningSys::get(&bucket).await);
|
let cfg = try_!(BucketVersioningSys::get(&bucket).await);
|
||||||
|
|
||||||
let status = match cfg.status {
|
let status = match cfg.status {
|
||||||
VersioningState::Enabled => Some(BucketVersioningStatus::ENABLED),
|
VersioningState::Enabled => Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)),
|
||||||
VersioningState::Suspended => Some(BucketVersioningStatus::SUSPENDED),
|
VersioningState::Suspended => Some(BucketVersioningStatus::from_static(BucketVersioningStatus::SUSPENDED)),
|
||||||
};
|
};
|
||||||
|
|
||||||
Ok(S3Response::new(GetBucketVersioningOutput {
|
Ok(S3Response::new(GetBucketVersioningOutput {
|
||||||
@@ -892,12 +895,43 @@ impl S3 for FS {
|
|||||||
&self,
|
&self,
|
||||||
req: S3Request<PutBucketVersioningInput>,
|
req: S3Request<PutBucketVersioningInput>,
|
||||||
) -> S3Result<S3Response<PutBucketVersioningOutput>> {
|
) -> S3Result<S3Response<PutBucketVersioningOutput>> {
|
||||||
let PutBucketVersioningInput { bucket, .. } = req;
|
let PutBucketVersioningInput {
|
||||||
|
bucket,
|
||||||
|
versioning_configuration,
|
||||||
|
..
|
||||||
|
} = req.input;
|
||||||
|
|
||||||
|
// TODO: check other sys
|
||||||
// check site replication enable
|
// check site replication enable
|
||||||
// check bucket object lock enable
|
// check bucket object lock enable
|
||||||
// check replication suspended
|
// check replication suspended
|
||||||
Err(s3_error!(NotImplemented, "PutBucketVersioning is not implemented yet"))
|
|
||||||
|
let mut cfg = match BucketVersioningSys::get(&bucket).await {
|
||||||
|
Ok(res) => res,
|
||||||
|
Err(err) => {
|
||||||
|
warn!("BucketVersioningSys::get err {:?}", err);
|
||||||
|
Versioning::default()
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
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);
|
||||||
|
|
||||||
|
// TODO: globalSiteReplicationSys.BucketMetaHook
|
||||||
|
|
||||||
|
Ok(S3Response::new(PutBucketVersioningOutput {}))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user