From 9d74f56f5741ae13bc9d8c1f96a4bc0a0412c0ac Mon Sep 17 00:00:00 2001 From: weisd Date: Tue, 1 Jul 2025 15:23:53 +0800 Subject: [PATCH] feat: add ExportBucketMetadata handler --- ecstore/src/bucket/metadata_sys.rs | 26 +- ecstore/src/bucket/policy_sys.rs | 2 +- rustfs/src/admin/handlers.rs | 1 + rustfs/src/admin/handlers/bucket_meta.rs | 346 +++++++++++++++++++++++ rustfs/src/admin/mod.rs | 140 ++++----- 5 files changed, 446 insertions(+), 69 deletions(-) create mode 100644 rustfs/src/admin/handlers/bucket_meta.rs diff --git a/ecstore/src/bucket/metadata_sys.rs b/ecstore/src/bucket/metadata_sys.rs index 5e3cf190e..333f5900f 100644 --- a/ecstore/src/bucket/metadata_sys.rs +++ b/ecstore/src/bucket/metadata_sys.rs @@ -19,7 +19,7 @@ use std::{collections::HashMap, sync::Arc}; use time::OffsetDateTime; use tokio::sync::RwLock; use tokio::time::sleep; -use tracing::{error, warn}; +use tracing::error; use super::metadata::{BucketMetadata, load_bucket_metadata}; use super::quota::BucketQuota; @@ -76,6 +76,27 @@ pub async fn delete(bucket: &str, config_file: &str) -> Result { bucket_meta_sys.delete(bucket, config_file).await } +pub async fn get_bucket_policy(bucket: &str) -> Result<(BucketPolicy, OffsetDateTime)> { + let bucket_meta_sys_lock = get_bucket_metadata_sys()?; + let bucket_meta_sys = bucket_meta_sys_lock.read().await; + + bucket_meta_sys.get_bucket_policy(bucket).await +} + +pub async fn get_quota_config(bucket: &str) -> Result<(BucketQuota, OffsetDateTime)> { + let bucket_meta_sys_lock = get_bucket_metadata_sys()?; + let bucket_meta_sys = bucket_meta_sys_lock.read().await; + + bucket_meta_sys.get_quota_config(bucket).await +} + +pub async fn get_bucket_targets_config(bucket: &str) -> Result { + let bucket_meta_sys_lock = get_bucket_metadata_sys()?; + let bucket_meta_sys = bucket_meta_sys_lock.read().await; + + bucket_meta_sys.get_bucket_targets_config(bucket).await +} + pub async fn get_tagging_config(bucket: &str) -> Result<(Tagging, OffsetDateTime)> { let bucket_meta_sys_lock = get_bucket_metadata_sys()?; let bucket_meta_sys = bucket_meta_sys_lock.read().await; @@ -340,7 +361,6 @@ impl BucketMetadataSys { } pub async fn get_config_from_disk(&self, bucket: &str) -> Result { - println!("load data from disk"); if is_meta_bucketname(bucket) { return Err(Error::other("errInvalidArgument")); } @@ -381,7 +401,6 @@ impl BucketMetadataSys { let bm = match self.get_config(bucket).await { Ok((res, _)) => res, Err(err) => { - warn!("get_versioning_config err {:?}", &err); return if err == Error::ConfigNotFound { Ok((VersioningConfiguration::default(), OffsetDateTime::UNIX_EPOCH)) } else { @@ -445,7 +464,6 @@ impl BucketMetadataSys { let bm = match self.get_config(bucket).await { Ok((bm, _)) => bm.notification_config.clone(), Err(err) => { - warn!("get_notification_config err {:?}", &err); if err == Error::ConfigNotFound { None } else { diff --git a/ecstore/src/bucket/policy_sys.rs b/ecstore/src/bucket/policy_sys.rs index 37c0c1a32..5f6e15ed9 100644 --- a/ecstore/src/bucket/policy_sys.rs +++ b/ecstore/src/bucket/policy_sys.rs @@ -21,7 +21,7 @@ impl PolicySys { } pub async fn get(bucket: &str) -> Result { let bucket_meta_sys_lock = get_bucket_metadata_sys()?; - let bucket_meta_sys = bucket_meta_sys_lock.write().await; + let bucket_meta_sys = bucket_meta_sys_lock.read().await; let (cfg, _) = bucket_meta_sys.get_bucket_policy(bucket).await?; diff --git a/rustfs/src/admin/handlers.rs b/rustfs/src/admin/handlers.rs index 5993438bb..b0f9acca0 100644 --- a/rustfs/src/admin/handlers.rs +++ b/rustfs/src/admin/handlers.rs @@ -56,6 +56,7 @@ use tokio_stream::wrappers::ReceiverStream; use tracing::{error, info, warn}; // use url::UrlQuery; +pub mod bucket_meta; pub mod event; pub mod group; pub mod policys; diff --git a/rustfs/src/admin/handlers/bucket_meta.rs b/rustfs/src/admin/handlers/bucket_meta.rs new file mode 100644 index 000000000..6fc0f4dea --- /dev/null +++ b/rustfs/src/admin/handlers/bucket_meta.rs @@ -0,0 +1,346 @@ +use std::io::{Cursor, Write as _}; + +use crate::{ + admin::router::Operation, + auth::{check_key_valid, get_session_token}, +}; +use ecstore::bucket::utils::serialize; +use ecstore::{ + StorageAPI, + bucket::{ + metadata::{ + BUCKET_LIFECYCLE_CONFIG, BUCKET_NOTIFICATION_CONFIG, BUCKET_POLICY_CONFIG, BUCKET_QUOTA_CONFIG_FILE, + BUCKET_REPLICATION_CONFIG, BUCKET_SSECONFIG, BUCKET_TAGGING_CONFIG, BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, + OBJECT_LOCK_CONFIG, + }, + metadata_sys, + }, + error::StorageError, + new_object_layer_fn, + store_api::BucketOptions, +}; +use http::{HeaderMap, StatusCode}; +use matchit::Params; +use rustfs_utils::path::path_join_buf; +use s3s::{ + Body, S3Request, S3Response, S3Result, + header::{CONTENT_DISPOSITION, CONTENT_LENGTH, CONTENT_TYPE}, + s3_error, +}; +use serde::Deserialize; +use serde_urlencoded::from_bytes; +use zip::{ZipArchive, ZipWriter, result::ZipError, write::SimpleFileOptions}; + +#[derive(Debug, Default, serde::Deserialize)] +pub struct ExportBucketMetadataQuery { + pub bucket: String, +} + +pub struct ExportBucketMetadata {} + +#[async_trait::async_trait] +impl Operation for ExportBucketMetadata { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let query = { + if let Some(query) = req.uri.query() { + let input: ExportBucketMetadataQuery = + from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get query failed"))?; + input + } else { + ExportBucketMetadataQuery::default() + } + }; + + let Some(input_cred) = req.credentials else { + return Err(s3_error!(InvalidRequest, "get cred failed")); + }; + + let (_cred, _owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; + + let Some(store) = new_object_layer_fn() else { + return Err(s3_error!(InvalidRequest, "object store not init")); + }; + + let buckets = if query.bucket.is_empty() { + store + .list_bucket(&BucketOptions::default()) + .await + .map_err(|e| s3_error!(InternalError, "list buckets failed: {e}"))? + } else { + let bucket = store + .get_bucket_info(&query.bucket, &BucketOptions::default()) + .await + .map_err(|e| s3_error!(InternalError, "get bucket failed: {e}"))?; + vec![bucket] + }; + + let mut zip_writer = ZipWriter::new(Cursor::new(Vec::new())); + + let confs = [ + BUCKET_POLICY_CONFIG, + BUCKET_NOTIFICATION_CONFIG, + BUCKET_LIFECYCLE_CONFIG, + BUCKET_SSECONFIG, + BUCKET_TAGGING_CONFIG, + BUCKET_QUOTA_CONFIG_FILE, + OBJECT_LOCK_CONFIG, + BUCKET_VERSIONING_CONFIG, + BUCKET_REPLICATION_CONFIG, + BUCKET_TARGETS_FILE, + ]; + + for bucket in buckets { + for &conf in confs.iter() { + let conf_path = path_join_buf(&[bucket.name.as_str(), conf]); + match conf { + BUCKET_POLICY_CONFIG => { + let config = match metadata_sys::get_bucket_policy(&bucket.name).await { + Ok((res, _)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + let config_json = + serde_json::to_vec(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_json) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + BUCKET_NOTIFICATION_CONFIG => { + let config = match metadata_sys::get_notification_config(&bucket.name).await { + Ok(Some(res)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + Ok(None) => continue, + }; + + let config_xml = + serialize(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_xml) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + BUCKET_LIFECYCLE_CONFIG => { + let config = match metadata_sys::get_lifecycle_config(&bucket.name).await { + Ok((res, _)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + let config_xml = + serialize(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_xml) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + BUCKET_TAGGING_CONFIG => { + let config = match metadata_sys::get_tagging_config(&bucket.name).await { + Ok((res, _)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + let config_xml = + serialize(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_xml) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + BUCKET_QUOTA_CONFIG_FILE => { + let config = match metadata_sys::get_quota_config(&bucket.name).await { + Ok((res, _)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + let config_json = + serde_json::to_vec(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_json) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + OBJECT_LOCK_CONFIG => { + let config = match metadata_sys::get_object_lock_config(&bucket.name).await { + Ok((res, _)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + let config_xml = + serialize(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_xml) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + BUCKET_SSECONFIG => { + let config = match metadata_sys::get_sse_config(&bucket.name).await { + Ok((res, _)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + let config_xml = + serialize(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_xml) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + BUCKET_VERSIONING_CONFIG => { + let config = match metadata_sys::get_versioning_config(&bucket.name).await { + Ok((res, _)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + let config_xml = + serialize(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_xml) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + BUCKET_REPLICATION_CONFIG => { + let config = match metadata_sys::get_replication_config(&bucket.name).await { + Ok((res, _)) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + let config_xml = + serialize(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_xml) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + BUCKET_TARGETS_FILE => { + let config = match metadata_sys::get_bucket_targets_config(&bucket.name).await { + Ok(res) => res, + Err(e) => { + if e == StorageError::ConfigNotFound { + continue; + } + return Err(s3_error!(InternalError, "get bucket metadata failed: {e}")); + } + }; + + let config_json = + serde_json::to_vec(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?; + + zip_writer + .start_file(conf_path, SimpleFileOptions::default()) + .map_err(|e| s3_error!(InternalError, "start file failed: {e}"))?; + zip_writer + .write_all(&config_json) + .map_err(|e| s3_error!(InternalError, "write file failed: {e}"))?; + } + _ => {} + } + } + } + + let zip_bytes = zip_writer + .finish() + .map_err(|e| s3_error!(InternalError, "finish zip failed: {e}"))?; + let mut header = HeaderMap::new(); + header.insert(CONTENT_TYPE, "application/zip".parse().unwrap()); + header.insert(CONTENT_DISPOSITION, "attachment; filename=bucket-meta.zip".parse().unwrap()); + header.insert(CONTENT_LENGTH, zip_bytes.get_ref().len().to_string().parse().unwrap()); + Ok(S3Response::with_headers((StatusCode::OK, Body::from(zip_bytes.into_inner())), header)) + } +} + +#[derive(Debug, Default, Deserialize)] +pub struct ImportBucketMetadataQuery { + pub bucket: String, +} + +pub struct ImportBucketMetadata {} + +#[async_trait::async_trait] +impl Operation for ImportBucketMetadata { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + let _query = { + if let Some(query) = req.uri.query() { + let input: ImportBucketMetadataQuery = + from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get query failed"))?; + input + } else { + ImportBucketMetadataQuery::default() + } + }; + + let Some(input_cred) = req.credentials else { + return Err(s3_error!(InvalidRequest, "get cred failed")); + }; + + let (_cred, _owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; + + let mut header = HeaderMap::new(); + header.insert(CONTENT_TYPE, "application/json".parse().unwrap()); + Ok(S3Response::with_headers((StatusCode::OK, Body::empty()), header)) + } +} diff --git a/rustfs/src/admin/mod.rs b/rustfs/src/admin/mod.rs index 38c3030a6..17046522d 100644 --- a/rustfs/src/admin/mod.rs +++ b/rustfs/src/admin/mod.rs @@ -5,7 +5,7 @@ pub mod utils; // use ecstore::global::{is_dist_erasure, is_erasure}; use handlers::{ - group, policys, pools, rebalance, + bucket_meta, group, policys, pools, rebalance, service_account::{AddServiceAccount, DeleteServiceAccount, InfoServiceAccount, ListServiceAccount, UpdateServiceAccount}, sts, tier, user, }; @@ -119,11 +119,85 @@ pub fn make_admin_route() -> std::io::Result { format!("{}{}", ADMIN_PREFIX, "/v3/background-heal/status").as_str(), AdminOperation(&handlers::BackgroundHealStatusHandler {}), )?; - // } + + // ? + r.insert( + Method::GET, + format!("{}{}", ADMIN_PREFIX, "/v3/tier").as_str(), + AdminOperation(&tier::ListTiers {}), + )?; + // ? + r.insert( + Method::GET, + format!("{}{}", ADMIN_PREFIX, "/v3/tier-stats").as_str(), + AdminOperation(&tier::GetTierInfo {}), + )?; + // ?force=xxx + r.insert( + Method::DELETE, + format!("{}{}", ADMIN_PREFIX, "/v3/tier/{tiername}").as_str(), + AdminOperation(&tier::RemoveTier {}), + )?; + // ?force=xxx + // body: AddOrUpdateTierReq + r.insert( + Method::PUT, + format!("{}{}", ADMIN_PREFIX, "/v3/tier").as_str(), + AdminOperation(&tier::AddTier {}), + )?; + // ? + // body: AddOrUpdateTierReq + r.insert( + Method::POST, + format!("{}{}", ADMIN_PREFIX, "/v3/tier/{tiername}").as_str(), + AdminOperation(&tier::EditTier {}), + )?; + r.insert( + Method::POST, + format!("{}{}", ADMIN_PREFIX, "/v3/tier/clear").as_str(), + AdminOperation(&tier::ClearTier {}), + )?; + + r.insert( + Method::GET, + format!("{}{}", ADMIN_PREFIX, "/export-bucket-metadata").as_str(), + AdminOperation(&bucket_meta::ExportBucketMetadata {}), + )?; + + r.insert( + Method::PUT, + format!("{}{}", ADMIN_PREFIX, "/import-bucket-metadata").as_str(), + AdminOperation(&bucket_meta::ImportBucketMetadata {}), + )?; + + r.insert( + Method::GET, + format!("{}{}", ADMIN_PREFIX, "/v3/list-remote-targets").as_str(), + AdminOperation(&ListRemoteTargetHandler {}), + )?; + + r.insert( + Method::GET, + format!("{}{}", ADMIN_PREFIX, "/v3/replicationmetrics").as_str(), + AdminOperation(&GetReplicationMetricsHandler {}), + )?; + + r.insert( + Method::PUT, + format!("{}{}", ADMIN_PREFIX, "/v3/set-remote-target").as_str(), + AdminOperation(&SetRemoteTargetHandler {}), + )?; + + r.insert( + Method::DELETE, + format!("{}{}", ADMIN_PREFIX, "/v3/remove-remote-target").as_str(), + AdminOperation(&RemoveRemoteTargetHandler {}), + )?; Ok(r) } +/// user router fn register_user_route(r: &mut S3Router) -> std::io::Result<()> { // 1 r.insert( @@ -240,30 +314,6 @@ fn register_user_route(r: &mut S3Router) -> std::io::Result<()> AdminOperation(&user::ImportIam {}), )?; - r.insert( - Method::GET, - format!("{}{}", ADMIN_PREFIX, "/v3/list-remote-targets").as_str(), - AdminOperation(&ListRemoteTargetHandler {}), - )?; - - r.insert( - Method::GET, - format!("{}{}", ADMIN_PREFIX, "/v3/replicationmetrics").as_str(), - AdminOperation(&GetReplicationMetricsHandler {}), - )?; - - r.insert( - Method::PUT, - format!("{}{}", ADMIN_PREFIX, "/v3/set-remote-target").as_str(), - AdminOperation(&SetRemoteTargetHandler {}), - )?; - - r.insert( - Method::DELETE, - format!("{}{}", ADMIN_PREFIX, "/v3/remove-remote-target").as_str(), - AdminOperation(&RemoveRemoteTargetHandler {}), - )?; - // list-canned-policies?bucket=xxx r.insert( Method::GET, @@ -299,43 +349,5 @@ fn register_user_route(r: &mut S3Router) -> std::io::Result<()> AdminOperation(&policys::SetPolicyForUserOrGroup {}), )?; - // ? - r.insert( - Method::GET, - format!("{}{}", ADMIN_PREFIX, "/v3/tier").as_str(), - AdminOperation(&tier::ListTiers {}), - )?; - // ? - r.insert( - Method::GET, - format!("{}{}", ADMIN_PREFIX, "/v3/tier-stats").as_str(), - AdminOperation(&tier::GetTierInfo {}), - )?; - // ?force=xxx - r.insert( - Method::DELETE, - format!("{}{}", ADMIN_PREFIX, "/v3/tier/{tiername}").as_str(), - AdminOperation(&tier::RemoveTier {}), - )?; - // ?force=xxx - // body: AddOrUpdateTierReq - r.insert( - Method::PUT, - format!("{}{}", ADMIN_PREFIX, "/v3/tier").as_str(), - AdminOperation(&tier::AddTier {}), - )?; - // ? - // body: AddOrUpdateTierReq - r.insert( - Method::POST, - format!("{}{}", ADMIN_PREFIX, "/v3/tier/{tiername}").as_str(), - AdminOperation(&tier::EditTier {}), - )?; - r.insert( - Method::POST, - format!("{}{}", ADMIN_PREFIX, "/v3/tier/clear").as_str(), - AdminOperation(&tier::ClearTier {}), - )?; - Ok(()) }