From 1092b2696a6f154af7286df3c6acf0539074ce10 Mon Sep 17 00:00:00 2001 From: weisd Date: Sat, 5 Oct 2024 01:41:19 +0800 Subject: [PATCH] update bucket tagging op use bucketmetadata_sys --- ecstore/src/bucket/metadata.rs | 1 - ecstore/src/bucket/metadata_sys.rs | 30 +++- ecstore/src/bucket/mod.rs | 2 +- ecstore/src/bucket/tags/mod.rs | 17 +-- ecstore/src/config/common.rs | 1 + ecstore/src/store.rs | 2 +- rustfs/src/storage/ecfs.rs | 218 +++++++++++++++++------------ 7 files changed, 168 insertions(+), 103 deletions(-) diff --git a/ecstore/src/bucket/metadata.rs b/ecstore/src/bucket/metadata.rs index 8f349a983..5f6dd3884 100644 --- a/ecstore/src/bucket/metadata.rs +++ b/ecstore/src/bucket/metadata.rs @@ -265,7 +265,6 @@ impl BucketMetadata { fn parse_all_configs(&mut self, _api: &ECStore) -> Result<()> { 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)?); } diff --git a/ecstore/src/bucket/metadata_sys.rs b/ecstore/src/bucket/metadata_sys.rs index a73c08c61..2e0711cec 100644 --- a/ecstore/src/bucket/metadata_sys.rs +++ b/ecstore/src/bucket/metadata_sys.rs @@ -31,7 +31,7 @@ pub async fn get_bucket_metadata_sys() -> Arc> { 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; sys.set(bucket, bm).await } @@ -135,10 +135,10 @@ impl BucketMetadataSys { } } - pub async fn set(&self, bucket: &str, bm: BucketMetadata) { - if !is_meta_bucketname(bucket) { + pub async fn set(&self, bucket: String, bm: BucketMetadata) { + if !is_meta_bucketname(&bucket) { 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)?; - bm.save(store).await?; + self.save(&mut bm).await?; 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)> { if let Some(api) = self.api.as_ref() { let has_bm = { @@ -221,6 +240,7 @@ impl BucketMetadataSys { let bm = match self.get_config(bucket).await { Ok((res, _)) => res, Err(err) => { + warn!("get_tagging_config err {:?}", &err); if config::error::is_not_found(&err) { return Err(Error::new(BucketMetadataError::TaggingNotFound)); } else { diff --git a/ecstore/src/bucket/mod.rs b/ecstore/src/bucket/mod.rs index b1f0a6d8a..a905fba13 100644 --- a/ecstore/src/bucket/mod.rs +++ b/ecstore/src/bucket/mod.rs @@ -8,7 +8,7 @@ mod objectlock; mod policy; mod quota; mod replication; -mod tags; +pub mod tags; mod target; pub mod utils; mod versioning; diff --git a/ecstore/src/bucket/tags/mod.rs b/ecstore/src/bucket/tags/mod.rs index 62cfb9848..34c0ca793 100644 --- a/ecstore/src/bucket/tags/mod.rs +++ b/ecstore/src/bucket/tags/mod.rs @@ -1,4 +1,4 @@ -use crate::error::{Error, Result}; +use crate::error::Result; use rmp_serde::Serializer as rmpSerializer; use serde::{Deserialize, Serialize}; use std::collections::HashMap; @@ -6,21 +6,22 @@ use std::collections::HashMap; // 定义tagSet结构体 #[derive(Debug, Deserialize, Serialize, Default, Clone)] pub struct TagSet { - #[serde(rename = "Tag")] - tag_map: HashMap, - is_object: bool, + pub tag_map: HashMap, + pub is_object: bool, } // 定义tagging结构体 #[derive(Debug, Deserialize, Serialize, Default, Clone)] pub struct Tags { - #[serde(rename = "Tagging")] - xml_name: String, - #[serde(rename = "TagSet")] - tag_set: Option, + pub tag_set: TagSet, } impl Tags { + pub fn new(tag_map: HashMap, is_object: bool) -> Self { + Self { + tag_set: TagSet { tag_map, is_object }, + } + } pub fn marshal_msg(&self) -> Result> { let mut buf = Vec::new(); diff --git a/ecstore/src/config/common.rs b/ecstore/src/config/common.rs index 6dedfbed8..62e53252a 100644 --- a/ecstore/src/config/common.rs +++ b/ecstore/src/config/common.rs @@ -5,6 +5,7 @@ use crate::store_api::{HTTPRangeSpec, ObjectIO, ObjectInfo, ObjectOptions, PutOb use http::HeaderMap; use s3s::dto::StreamingBlob; use s3s::Body; +use tracing::warn; use super::error::ConfigError; diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 321740e3a..41bcba04d 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -510,7 +510,7 @@ impl StorageAPI for ECStore { let mut meta = BucketMetadata::new(bucket); meta.save(self).await?; - bucket_metadata_sys_set(bucket, meta).await; + bucket_metadata_sys_set(bucket.to_string(), meta).await; // TODO: toObjectErr diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index dd1855f38..53607001a 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -1,5 +1,8 @@ use bytes::BufMut; 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::disk::error::DiskError; use ecstore::disk::RUSTFS_META_BUCKET; @@ -18,6 +21,7 @@ use ecstore::store_api::StorageAPI; use futures::pin_mut; use futures::{Stream, StreamExt}; use http::HeaderMap; +use log::warn; use s3s::dto::*; use s3s::s3_error; use s3s::Body; @@ -26,6 +30,7 @@ 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 transform_stream::AsyncTryStream; @@ -702,53 +707,67 @@ impl S3 for FS { })) .await?; - 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"))?; + let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let mut bucket_meta_sys = bucket_meta_sys_lock.write().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 tag_map = HashMap::new(); + for tag in tagging.tag_set.iter() { + tag_map.insert(tag.key.clone(), tag.value.clone()); } - 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 tags = Tags::new(tag_map, false); - 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 - ); + let data = try_!(tags.marshal_msg()); + + let _updated = try_!(bucket_meta_sys.update(&bucket, BUCKET_TAGGING_CONFIG, data).await); + + // 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"))?; + + // 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())) } @@ -764,15 +783,29 @@ impl S3 for FS { })) .await?; - 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"))?; + // 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"))?; - // 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 = 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_!( // store @@ -842,48 +875,59 @@ impl S3 for FS { })) .await?; - 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"))?; + let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let mut bucket_meta_sys = bucket_meta_sys_lock.write().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 tag_map = HashMap::new(); - let stream = meta_obj.stream; + let tags = Tags::new(tag_map, false); - let mut data = vec![]; - pin_mut!(stream); + let data = try_!(tags.marshal_msg()); - while let Some(x) = stream.next().await { - let x = try_!(x); - data.put_slice(&x[..]); - } + let _updated = try_!(bucket_meta_sys.update(&bucket, BUCKET_TAGGING_CONFIG, data).await); - 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 - ); + // 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"))?; + + // 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[..])); + // 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 {})) }