From 35f4cdd2970cac175c539118abf09d549e8ae4f9 Mon Sep 17 00:00:00 2001 From: weisd Date: Wed, 26 Mar 2025 16:22:09 +0800 Subject: [PATCH] move admin handlers --- rustfs/src/admin/handlers.rs | 247 +--------------------------- rustfs/src/admin/handlers/pools.rs | 252 +++++++++++++++++++++++++++++ rustfs/src/admin/mod.rs | 10 +- 3 files changed, 258 insertions(+), 251 deletions(-) create mode 100644 rustfs/src/admin/handlers/pools.rs diff --git a/rustfs/src/admin/handlers.rs b/rustfs/src/admin/handlers.rs index 3a53de3c9..886731c87 100644 --- a/rustfs/src/admin/handlers.rs +++ b/rustfs/src/admin/handlers.rs @@ -1,5 +1,4 @@ use super::router::Operation; -use crate::storage::error::to_s3_error; use ::policy::policy::action::{Action, S3Action}; use ::policy::policy::resource::Resource; use ::policy::policy::statement::BPStatement; @@ -17,7 +16,6 @@ use ecstore::peer::is_reserved_or_invalid_bucket; use ecstore::store::is_valid_object_prefix; use ecstore::store_api::StorageAPI; use ecstore::utils::path::path_join; -use ecstore::GLOBAL_Endpoints; use futures::{Stream, StreamExt}; use http::{HeaderMap, Uri}; use hyper::StatusCode; @@ -29,7 +27,6 @@ use s3s::stream::{ByteStream, DynByteStream}; use s3s::{s3_error, Body, S3Error, S3Request, S3Response, S3Result}; use s3s::{S3ErrorCode, StdError}; use serde::{Deserialize, Serialize}; -use serde_urlencoded::from_bytes; use std::collections::{HashMap, HashSet}; use std::path::PathBuf; use std::pin::Pin; @@ -44,6 +41,7 @@ use tracing::{error, info, warn}; pub mod group; pub mod policy; +pub mod pools; pub mod service_account; pub mod sts; pub mod trace; @@ -593,249 +591,6 @@ impl Operation for BackgroundHealStatusHandler { } } -pub struct ListPools {} - -#[async_trait::async_trait] -impl Operation for ListPools { - // GET //pools/list - async fn call(&self, _req: S3Request, _params: Params<'_, '_>) -> S3Result> { - warn!("handle ListPools"); - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let Some(endpoints) = GLOBAL_Endpoints.get() else { - return Err(s3_error!(NotImplemented)); - }; - - if endpoints.legacy() { - return Err(s3_error!(NotImplemented)); - } - - let mut pools_status = Vec::new(); - - for (idx, _) in endpoints.as_ref().iter().enumerate() { - let state = store.status(idx).await.map_err(to_s3_error)?; - - pools_status.push(state); - } - - let data = serde_json::to_vec(&pools_status) - .map_err(|_e| S3Error::with_message(S3ErrorCode::InternalError, "parse accountInfo failed"))?; - - let mut header = HeaderMap::new(); - header.insert(CONTENT_TYPE, "application/json".parse().unwrap()); - - Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), header)) - } -} - -#[derive(Debug, Deserialize, Default)] -#[serde(default)] -pub struct StatusPoolQuery { - pub pool: String, - #[serde(rename = "by-id")] - pub by_id: String, -} - -pub struct StatusPool {} - -#[async_trait::async_trait] -impl Operation for StatusPool { - // GET //pools/status?pool=http://server{1...4}/disk{1...4} - async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { - warn!("handle StatusPool"); - - let Some(endpoints) = GLOBAL_Endpoints.get() else { - return Err(s3_error!(NotImplemented)); - }; - - if endpoints.legacy() { - return Err(s3_error!(NotImplemented)); - } - - let query = { - if let Some(query) = req.uri.query() { - let input: StatusPoolQuery = - from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get body failed"))?; - input - } else { - StatusPoolQuery::default() - } - }; - - let is_byid = query.by_id.as_str() == "true"; - - let has_idx = { - if is_byid { - let a = query.pool.parse::().unwrap_or_default(); - if a < endpoints.as_ref().len() { - Some(a) - } else { - None - } - } else { - endpoints.get_pool_idx(&query.pool) - } - }; - - let Some(idx) = has_idx else { - warn!("specified pool {} not found, please specify a valid pool", &query.pool); - return Err(s3_error!(InvalidArgument)); - }; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let pools_status = store.status(idx).await.map_err(to_s3_error)?; - - let data = serde_json::to_vec(&pools_status) - .map_err(|_e| S3Error::with_message(S3ErrorCode::InternalError, "parse accountInfo failed"))?; - - let mut header = HeaderMap::new(); - header.insert(CONTENT_TYPE, "application/json".parse().unwrap()); - - Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), header)) - } -} - -pub struct StartDecommission {} - -#[async_trait::async_trait] -impl Operation for StartDecommission { - // POST //pools/decommission?pool=http://server{1...4}/disk{1...4} - async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { - warn!("handle StartDecommission"); - - let Some(endpoints) = GLOBAL_Endpoints.get() else { - return Err(s3_error!(NotImplemented)); - }; - - if endpoints.legacy() { - return Err(s3_error!(NotImplemented)); - } - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - if store.is_decommission_running().await { - return Err(S3Error::with_message( - S3ErrorCode::InvalidRequest, - "DecommissionAlreadyRunning".to_string(), - )); - } - - // TODO: check IsRebalanceStarted - - let query = { - if let Some(query) = req.uri.query() { - let input: StatusPoolQuery = - from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get body failed"))?; - input - } else { - StatusPoolQuery::default() - } - }; - let is_byid = query.by_id.as_str() == "true"; - - let pools: Vec<&str> = query.pool.split(",").collect(); - let mut pools_indices = Vec::with_capacity(pools.len()); - - for pool in pools.iter() { - let idx = { - if is_byid { - pool.parse::() - .map_err(|_e| s3_error!(InvalidArgument, "pool parse failed"))? - } else { - let Some(idx) = endpoints.get_pool_idx(pool) else { - return Err(s3_error!(InvalidArgument, "pool parse failed")); - }; - idx - } - }; - - let mut has_found = None; - for (i, pool) in store.pools.iter().enumerate() { - if i == idx { - has_found = Some(pool.clone()); - break; - } - } - - let Some(_p) = has_found else { - return Err(s3_error!(InvalidArgument)); - }; - - pools_indices.push(idx); - } - - if !pools_indices.is_empty() { - store.decommission(pools_indices).await.map_err(to_s3_error)?; - } - - Ok(S3Response::new((StatusCode::OK, Body::default()))) - } -} - -pub struct CancelDecommission {} - -#[async_trait::async_trait] -impl Operation for CancelDecommission { - // POST //pools/cancel?pool=http://server{1...4}/disk{1...4} - async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { - warn!("handle CancelDecommission"); - - let Some(endpoints) = GLOBAL_Endpoints.get() else { - return Err(s3_error!(NotImplemented)); - }; - - if endpoints.legacy() { - return Err(s3_error!(NotImplemented)); - } - - let query = { - if let Some(query) = req.uri.query() { - let input: StatusPoolQuery = - from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get body failed"))?; - input - } else { - StatusPoolQuery::default() - } - }; - - let is_byid = query.by_id.as_str() == "true"; - - let has_idx = { - if is_byid { - let a = query.pool.parse::().unwrap_or_default(); - if a < endpoints.as_ref().len() { - Some(a) - } else { - None - } - } else { - endpoints.get_pool_idx(&query.pool) - } - }; - - let Some(idx) = has_idx else { - warn!("specified pool {} not found, please specify a valid pool", &query.pool); - return Err(s3_error!(InvalidArgument)); - }; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - store.decommission_cancel(idx).await.map_err(to_s3_error)?; - - Ok(S3Response::new((StatusCode::OK, Body::default()))) - } -} - pub struct RebalanceStart {} #[async_trait::async_trait] diff --git a/rustfs/src/admin/handlers/pools.rs b/rustfs/src/admin/handlers/pools.rs new file mode 100644 index 000000000..413cdb4dd --- /dev/null +++ b/rustfs/src/admin/handlers/pools.rs @@ -0,0 +1,252 @@ +use ecstore::{new_object_layer_fn, GLOBAL_Endpoints}; +use http::{HeaderMap, StatusCode}; +use matchit::Params; +use s3s::{header::CONTENT_TYPE, s3_error, Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result}; +use serde::Deserialize; +use serde_urlencoded::from_bytes; +use tracing::warn; + +use crate::{admin::router::Operation, storage::error::to_s3_error}; + +pub struct ListPools {} + +#[async_trait::async_trait] +impl Operation for ListPools { + // GET //pools/list + async fn call(&self, _req: S3Request, _params: Params<'_, '_>) -> S3Result> { + warn!("handle ListPools"); + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let Some(endpoints) = GLOBAL_Endpoints.get() else { + return Err(s3_error!(NotImplemented)); + }; + + if endpoints.legacy() { + return Err(s3_error!(NotImplemented)); + } + + let mut pools_status = Vec::new(); + + for (idx, _) in endpoints.as_ref().iter().enumerate() { + let state = store.status(idx).await.map_err(to_s3_error)?; + + pools_status.push(state); + } + + let data = serde_json::to_vec(&pools_status) + .map_err(|_e| S3Error::with_message(S3ErrorCode::InternalError, "parse accountInfo failed"))?; + + let mut header = HeaderMap::new(); + header.insert(CONTENT_TYPE, "application/json".parse().unwrap()); + + Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), header)) + } +} + +#[derive(Debug, Deserialize, Default)] +#[serde(default)] +pub struct StatusPoolQuery { + pub pool: String, + #[serde(rename = "by-id")] + pub by_id: String, +} + +pub struct StatusPool {} + +#[async_trait::async_trait] +impl Operation for StatusPool { + // GET //pools/status?pool=http://server{1...4}/disk{1...4} + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + warn!("handle StatusPool"); + + let Some(endpoints) = GLOBAL_Endpoints.get() else { + return Err(s3_error!(NotImplemented)); + }; + + if endpoints.legacy() { + return Err(s3_error!(NotImplemented)); + } + + let query = { + if let Some(query) = req.uri.query() { + let input: StatusPoolQuery = + from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get body failed"))?; + input + } else { + StatusPoolQuery::default() + } + }; + + let is_byid = query.by_id.as_str() == "true"; + + let has_idx = { + if is_byid { + let a = query.pool.parse::().unwrap_or_default(); + if a < endpoints.as_ref().len() { + Some(a) + } else { + None + } + } else { + endpoints.get_pool_idx(&query.pool) + } + }; + + let Some(idx) = has_idx else { + warn!("specified pool {} not found, please specify a valid pool", &query.pool); + return Err(s3_error!(InvalidArgument)); + }; + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let pools_status = store.status(idx).await.map_err(to_s3_error)?; + + let data = serde_json::to_vec(&pools_status) + .map_err(|_e| S3Error::with_message(S3ErrorCode::InternalError, "parse accountInfo failed"))?; + + let mut header = HeaderMap::new(); + header.insert(CONTENT_TYPE, "application/json".parse().unwrap()); + + Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), header)) + } +} + +pub struct StartDecommission {} + +#[async_trait::async_trait] +impl Operation for StartDecommission { + // POST //pools/decommission?pool=http://server{1...4}/disk{1...4} + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + warn!("handle StartDecommission"); + + let Some(endpoints) = GLOBAL_Endpoints.get() else { + return Err(s3_error!(NotImplemented)); + }; + + if endpoints.legacy() { + return Err(s3_error!(NotImplemented)); + } + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + if store.is_decommission_running().await { + return Err(S3Error::with_message( + S3ErrorCode::InvalidRequest, + "DecommissionAlreadyRunning".to_string(), + )); + } + + // TODO: check IsRebalanceStarted + + let query = { + if let Some(query) = req.uri.query() { + let input: StatusPoolQuery = + from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get body failed"))?; + input + } else { + StatusPoolQuery::default() + } + }; + let is_byid = query.by_id.as_str() == "true"; + + let pools: Vec<&str> = query.pool.split(",").collect(); + let mut pools_indices = Vec::with_capacity(pools.len()); + + for pool in pools.iter() { + let idx = { + if is_byid { + pool.parse::() + .map_err(|_e| s3_error!(InvalidArgument, "pool parse failed"))? + } else { + let Some(idx) = endpoints.get_pool_idx(pool) else { + return Err(s3_error!(InvalidArgument, "pool parse failed")); + }; + idx + } + }; + + let mut has_found = None; + for (i, pool) in store.pools.iter().enumerate() { + if i == idx { + has_found = Some(pool.clone()); + break; + } + } + + let Some(_p) = has_found else { + return Err(s3_error!(InvalidArgument)); + }; + + pools_indices.push(idx); + } + + if !pools_indices.is_empty() { + store.decommission(pools_indices).await.map_err(to_s3_error)?; + } + + Ok(S3Response::new((StatusCode::OK, Body::default()))) + } +} + +pub struct CancelDecommission {} + +#[async_trait::async_trait] +impl Operation for CancelDecommission { + // POST //pools/cancel?pool=http://server{1...4}/disk{1...4} + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + warn!("handle CancelDecommission"); + + let Some(endpoints) = GLOBAL_Endpoints.get() else { + return Err(s3_error!(NotImplemented)); + }; + + if endpoints.legacy() { + return Err(s3_error!(NotImplemented)); + } + + let query = { + if let Some(query) = req.uri.query() { + let input: StatusPoolQuery = + from_bytes(query.as_bytes()).map_err(|_e| s3_error!(InvalidArgument, "get body failed"))?; + input + } else { + StatusPoolQuery::default() + } + }; + + let is_byid = query.by_id.as_str() == "true"; + + let has_idx = { + if is_byid { + let a = query.pool.parse::().unwrap_or_default(); + if a < endpoints.as_ref().len() { + Some(a) + } else { + None + } + } else { + endpoints.get_pool_idx(&query.pool) + } + }; + + let Some(idx) = has_idx else { + warn!("specified pool {} not found, please specify a valid pool", &query.pool); + return Err(s3_error!(InvalidArgument)); + }; + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + store.decommission_cancel(idx).await.map_err(to_s3_error)?; + + Ok(S3Response::new((StatusCode::OK, Body::default()))) + } +} diff --git a/rustfs/src/admin/mod.rs b/rustfs/src/admin/mod.rs index 4b3981327..5870b284d 100644 --- a/rustfs/src/admin/mod.rs +++ b/rustfs/src/admin/mod.rs @@ -6,7 +6,7 @@ pub mod utils; use common::error::Result; // use ecstore::global::{is_dist_erasure, is_erasure}; use handlers::{ - group, policy, + group, policy, pools, service_account::{AddServiceAccount, DeleteServiceAccount, InfoServiceAccount, ListServiceAccount, UpdateServiceAccount}, sts, user, }; @@ -69,25 +69,25 @@ pub fn make_admin_route() -> Result { r.insert( Method::GET, format!("{}{}", ADMIN_PREFIX, "/v3/pools/list").as_str(), - AdminOperation(&handlers::ListPools {}), + AdminOperation(&pools::ListPools {}), )?; // 1 r.insert( Method::GET, format!("{}{}", ADMIN_PREFIX, "/v3/pools/status").as_str(), - AdminOperation(&handlers::StatusPool {}), + AdminOperation(&pools::StatusPool {}), )?; // todo r.insert( Method::POST, format!("{}{}", ADMIN_PREFIX, "/v3/pools/decommission").as_str(), - AdminOperation(&handlers::StartDecommission {}), + AdminOperation(&pools::StartDecommission {}), )?; // todo r.insert( Method::POST, format!("{}{}", ADMIN_PREFIX, "/v3/pools/cancel").as_str(), - AdminOperation(&handlers::CancelDecommission {}), + AdminOperation(&pools::CancelDecommission {}), )?; r.insert(