update bucket tagging op use bucketmetadata_sys

This commit is contained in:
weisd
2024-10-05 01:41:19 +08:00
parent c29199515c
commit 1092b2696a
7 changed files with 168 additions and 103 deletions
-1
View File
@@ -265,7 +265,6 @@ impl BucketMetadata {
fn parse_all_configs(&mut self, _api: &ECStore) -> Result<()> { fn parse_all_configs(&mut self, _api: &ECStore) -> Result<()> {
if !self.tagging_config_xml.is_empty() { if !self.tagging_config_xml.is_empty() {
warn!("self.tagging_config_xml {:?}", &self.tagging_config_xml);
self.tagging_config = Some(tags::Tags::unmarshal(&self.tagging_config_xml)?); self.tagging_config = Some(tags::Tags::unmarshal(&self.tagging_config_xml)?);
} }
+25 -5
View File
@@ -31,7 +31,7 @@ pub async fn get_bucket_metadata_sys() -> Arc<RwLock<BucketMetadataSys>> {
GLOBAL_BucketMetadataSys.clone() GLOBAL_BucketMetadataSys.clone()
} }
pub async fn bucket_metadata_sys_set(bucket: &str, bm: BucketMetadata) { pub async fn bucket_metadata_sys_set(bucket: String, bm: BucketMetadata) {
let sys = GLOBAL_BucketMetadataSys.write().await; let sys = GLOBAL_BucketMetadataSys.write().await;
sys.set(bucket, bm).await sys.set(bucket, bm).await
} }
@@ -135,10 +135,10 @@ impl BucketMetadataSys {
} }
} }
pub async fn set(&self, bucket: &str, bm: BucketMetadata) { pub async fn set(&self, bucket: String, bm: BucketMetadata) {
if !is_meta_bucketname(bucket) { if !is_meta_bucketname(&bucket) {
let mut map = self.metadata_map.write().await; let mut map = self.metadata_map.write().await;
map.insert(bucket.to_string(), bm); map.insert(bucket, bm);
} }
} }
@@ -176,11 +176,30 @@ impl BucketMetadataSys {
let updated = bm.update_config(config_file, data)?; let updated = bm.update_config(config_file, data)?;
bm.save(store).await?; self.save(&mut bm).await?;
Ok(updated) Ok(updated)
} }
async fn save(&self, bm: &mut BucketMetadata) -> Result<()> {
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = match lock.as_ref() {
Some(s) => s,
None => return Err(Error::msg("errServerNotInitialized")),
};
if is_meta_bucketname(&bm.name) {
return Err(Error::msg("errInvalidArgument"));
}
bm.save(store).await?;
self.set(bm.name.clone(), bm.clone()).await;
Ok(())
}
pub async fn get_config(&self, bucket: &str) -> Result<(BucketMetadata, bool)> { pub async fn get_config(&self, bucket: &str) -> Result<(BucketMetadata, bool)> {
if let Some(api) = self.api.as_ref() { if let Some(api) = self.api.as_ref() {
let has_bm = { let has_bm = {
@@ -221,6 +240,7 @@ impl BucketMetadataSys {
let bm = match self.get_config(bucket).await { let bm = match self.get_config(bucket).await {
Ok((res, _)) => res, Ok((res, _)) => res,
Err(err) => { Err(err) => {
warn!("get_tagging_config err {:?}", &err);
if config::error::is_not_found(&err) { if config::error::is_not_found(&err) {
return Err(Error::new(BucketMetadataError::TaggingNotFound)); return Err(Error::new(BucketMetadataError::TaggingNotFound));
} else { } else {
+1 -1
View File
@@ -8,7 +8,7 @@ mod objectlock;
mod policy; mod policy;
mod quota; mod quota;
mod replication; mod replication;
mod tags; pub mod tags;
mod target; mod target;
pub mod utils; pub mod utils;
mod versioning; mod versioning;
+9 -8
View File
@@ -1,4 +1,4 @@
use crate::error::{Error, Result}; use crate::error::Result;
use rmp_serde::Serializer as rmpSerializer; use rmp_serde::Serializer as rmpSerializer;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::collections::HashMap; use std::collections::HashMap;
@@ -6,21 +6,22 @@ use std::collections::HashMap;
// 定义tagSet结构体 // 定义tagSet结构体
#[derive(Debug, Deserialize, Serialize, Default, Clone)] #[derive(Debug, Deserialize, Serialize, Default, Clone)]
pub struct TagSet { pub struct TagSet {
#[serde(rename = "Tag")] pub tag_map: HashMap<String, String>,
tag_map: HashMap<String, String>, pub is_object: bool,
is_object: bool,
} }
// 定义tagging结构体 // 定义tagging结构体
#[derive(Debug, Deserialize, Serialize, Default, Clone)] #[derive(Debug, Deserialize, Serialize, Default, Clone)]
pub struct Tags { pub struct Tags {
#[serde(rename = "Tagging")] pub tag_set: TagSet,
xml_name: String,
#[serde(rename = "TagSet")]
tag_set: Option<TagSet>,
} }
impl Tags { impl Tags {
pub fn new(tag_map: HashMap<String, String>, is_object: bool) -> Self {
Self {
tag_set: TagSet { tag_map, is_object },
}
}
pub fn marshal_msg(&self) -> Result<Vec<u8>> { pub fn marshal_msg(&self) -> Result<Vec<u8>> {
let mut buf = Vec::new(); let mut buf = Vec::new();
+1
View File
@@ -5,6 +5,7 @@ use crate::store_api::{HTTPRangeSpec, ObjectIO, ObjectInfo, ObjectOptions, PutOb
use http::HeaderMap; use http::HeaderMap;
use s3s::dto::StreamingBlob; use s3s::dto::StreamingBlob;
use s3s::Body; use s3s::Body;
use tracing::warn;
use super::error::ConfigError; use super::error::ConfigError;
+1 -1
View File
@@ -510,7 +510,7 @@ impl StorageAPI for ECStore {
let mut meta = BucketMetadata::new(bucket); let mut meta = BucketMetadata::new(bucket);
meta.save(self).await?; meta.save(self).await?;
bucket_metadata_sys_set(bucket, meta).await; bucket_metadata_sys_set(bucket.to_string(), meta).await;
// TODO: toObjectErr // TODO: toObjectErr
+131 -87
View File
@@ -1,5 +1,8 @@
use bytes::BufMut; use bytes::BufMut;
use bytes::Bytes; use bytes::Bytes;
use ecstore::bucket::get_bucket_metadata_sys;
use ecstore::bucket::metadata::BUCKET_TAGGING_CONFIG;
use ecstore::bucket::tags::Tags;
use ecstore::bucket_meta::BucketMetadata; use ecstore::bucket_meta::BucketMetadata;
use ecstore::disk::error::DiskError; use ecstore::disk::error::DiskError;
use ecstore::disk::RUSTFS_META_BUCKET; use ecstore::disk::RUSTFS_META_BUCKET;
@@ -18,6 +21,7 @@ use ecstore::store_api::StorageAPI;
use futures::pin_mut; use futures::pin_mut;
use futures::{Stream, StreamExt}; use futures::{Stream, StreamExt};
use http::HeaderMap; use http::HeaderMap;
use log::warn;
use s3s::dto::*; use s3s::dto::*;
use s3s::s3_error; use s3s::s3_error;
use s3s::Body; use s3s::Body;
@@ -26,6 +30,7 @@ use s3s::S3ErrorCode;
use s3s::S3Result; use s3s::S3Result;
use s3s::S3; use s3s::S3;
use s3s::{S3Request, S3Response}; use s3s::{S3Request, S3Response};
use std::collections::HashMap;
use std::fmt::Debug; use std::fmt::Debug;
use std::str::FromStr; use std::str::FromStr;
use transform_stream::AsyncTryStream; use transform_stream::AsyncTryStream;
@@ -702,53 +707,67 @@ impl S3 for FS {
})) }))
.await?; .await?;
let layer = new_object_layer_fn(); let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
let lock = layer.read().await; let mut bucket_meta_sys = bucket_meta_sys_lock.write().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let meta_obj = try_!( let mut tag_map = HashMap::new();
store for tag in tagging.tag_set.iter() {
.get_object_reader( tag_map.insert(tag.key.clone(), tag.value.clone());
RUSTFS_META_BUCKET,
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
HTTPRangeSpec::nil(),
Default::default(),
&ObjectOptions::default(),
)
.await
);
let stream = meta_obj.stream;
let mut data = vec![];
pin_mut!(stream);
while let Some(x) = stream.next().await {
let x = try_!(x);
data.put_slice(&x[..]);
} }
let mut meta = try_!(BucketMetadata::unmarshal_from(&data[..])); let tags = Tags::new(tag_map, false);
if tagging.tag_set.is_empty() {
meta.tagging = None;
} else {
meta.tagging = Some(tagging.tag_set.into_iter().map(|x| (x.key, x.value)).collect())
}
let data = try_!(meta.marshal_msg()); let data = try_!(tags.marshal_msg());
let len = data.len();
try_!( let _updated = try_!(bucket_meta_sys.update(&bucket, BUCKET_TAGGING_CONFIG, data).await);
store
.put_object( // let layer = new_object_layer_fn();
RUSTFS_META_BUCKET, // let lock = layer.read().await;
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(), // let store = lock
PutObjReader::new(StreamingBlob::from(Body::from(data)), len), // .as_ref()
&ObjectOptions::default(), // .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
)
.await // let meta_obj = try_!(
); // store
// .get_object_reader(
// RUSTFS_META_BUCKET,
// BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
// HTTPRangeSpec::nil(),
// Default::default(),
// &ObjectOptions::default(),
// )
// .await
// );
// let stream = meta_obj.stream;
// let mut data = vec![];
// pin_mut!(stream);
// while let Some(x) = stream.next().await {
// let x = try_!(x);
// data.put_slice(&x[..]);
// }
// let mut meta = try_!(BucketMetadata::unmarshal_from(&data[..]));
// if tagging.tag_set.is_empty() {
// meta.tagging = None;
// } else {
// meta.tagging = Some(tagging.tag_set.into_iter().map(|x| (x.key, x.value)).collect())
// }
// let data = try_!(meta.marshal_msg());
// let len = data.len();
// try_!(
// store
// .put_object(
// RUSTFS_META_BUCKET,
// BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
// PutObjReader::new(StreamingBlob::from(Body::from(data)), len),
// &ObjectOptions::default(),
// )
// .await
// );
Ok(S3Response::new(Default::default())) Ok(S3Response::new(Default::default()))
} }
@@ -764,15 +783,29 @@ impl S3 for FS {
})) }))
.await?; .await?;
let layer = new_object_layer_fn(); // let layer = new_object_layer_fn();
let lock = layer.read().await; // let lock = layer.read().await;
let store = lock // let store = lock
.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 a = get_bucket_metadata_sys(); 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(),
Err(err) => {
warn!("get_tagging_config err {:?}", &err);
// TODO: check not found
Vec::new()
}
};
Ok(S3Response::new(GetBucketTaggingOutput { ..Default::default() })) Ok(S3Response::new(GetBucketTaggingOutput { tag_set }))
// let meta_obj = try_!( // let meta_obj = try_!(
// store // store
@@ -842,48 +875,59 @@ impl S3 for FS {
})) }))
.await?; .await?;
let layer = new_object_layer_fn(); let bucket_meta_sys_lock = get_bucket_metadata_sys().await;
let lock = layer.read().await; let mut bucket_meta_sys = bucket_meta_sys_lock.write().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let meta_obj = try_!( let tag_map = HashMap::new();
store
.get_object_reader(
RUSTFS_META_BUCKET,
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
HTTPRangeSpec::nil(),
Default::default(),
&ObjectOptions::default(),
)
.await
);
let stream = meta_obj.stream; let tags = Tags::new(tag_map, false);
let mut data = vec![]; let data = try_!(tags.marshal_msg());
pin_mut!(stream);
while let Some(x) = stream.next().await { let _updated = try_!(bucket_meta_sys.update(&bucket, BUCKET_TAGGING_CONFIG, data).await);
let x = try_!(x);
data.put_slice(&x[..]);
}
let mut meta = try_!(BucketMetadata::unmarshal_from(&data[..])); // let layer = new_object_layer_fn();
meta.tagging = None; // let lock = layer.read().await;
let data = try_!(meta.marshal_msg()); // let store = lock
let len = data.len(); // .as_ref()
try_!( // .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
store
.put_object( // let meta_obj = try_!(
RUSTFS_META_BUCKET, // store
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(), // .get_object_reader(
PutObjReader::new(StreamingBlob::from(Body::from(data)), len), // RUSTFS_META_BUCKET,
&ObjectOptions::default(), // BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
) // HTTPRangeSpec::nil(),
.await // Default::default(),
); // &ObjectOptions::default(),
// )
// .await
// );
// let stream = meta_obj.stream;
// let mut data = vec![];
// pin_mut!(stream);
// while let Some(x) = stream.next().await {
// let x = try_!(x);
// data.put_slice(&x[..]);
// }
// let mut meta = try_!(BucketMetadata::unmarshal_from(&data[..]));
// meta.tagging = None;
// let data = try_!(meta.marshal_msg());
// let len = data.len();
// try_!(
// store
// .put_object(
// RUSTFS_META_BUCKET,
// BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
// PutObjReader::new(StreamingBlob::from(Body::from(data)), len),
// &ObjectOptions::default(),
// )
// .await
// );
Ok(S3Response::new(DeleteBucketTaggingOutput {})) Ok(S3Response::new(DeleteBucketTaggingOutput {}))
} }