From 16d66a420646a91b07e4d75259b348be7806fdd6 Mon Sep 17 00:00:00 2001 From: bestgopher <84328409@qq.com> Date: Tue, 29 Oct 2024 20:48:33 +0800 Subject: [PATCH] fix: blocking forever Closes #102 Signed-off-by: bestgopher <84328409@qq.com> --- api/admin/src/handlers/list_pools.rs | 11 +++- api/admin/src/lib.rs | 7 +-- api/admin/src/object_api.rs | 20 ------ rustfs/src/main.rs | 93 ++++++++++++++-------------- 4 files changed, 54 insertions(+), 77 deletions(-) delete mode 100644 api/admin/src/object_api.rs diff --git a/api/admin/src/handlers/list_pools.rs b/api/admin/src/handlers/list_pools.rs index b74e2d25d..faa9142e9 100644 --- a/api/admin/src/handlers/list_pools.rs +++ b/api/admin/src/handlers/list_pools.rs @@ -1,5 +1,5 @@ +use crate::error::ErrorCode; use crate::Result as LocalResult; -use crate::{error::ErrorCode, object_api::ObjectApi}; use axum::{extract::State, Json}; use serde::Serialize; @@ -39,12 +39,17 @@ struct PoolDecommissionInfo { bytes_failed: i64, } -pub async fn handler(State(ec_store): State) -> LocalResult>> { +pub async fn handler() -> LocalResult>> { // if ecstore::is_legacy().await { // return Err(ErrorCode::ErrNotImplemented); // } + // + // - let pools = (*ec_store).as_ref().ok_or(ErrorCode::ErrNotImplemented)?; + // todo 实用oncelock作为全局变量 + let layer = ecstore::new_object_layer_fn(); + let lock = layer.read().await; + let pools = lock.as_ref().ok_or(ErrorCode::ErrNotImplemented)?; // todo, 调用pool.status()接口获取每个池的数据 // diff --git a/api/admin/src/lib.rs b/api/admin/src/lib.rs index b04621748..a2270fbcd 100644 --- a/api/admin/src/lib.rs +++ b/api/admin/src/lib.rs @@ -1,26 +1,21 @@ pub mod error; pub mod handlers; -pub mod object_api; use axum::{extract::Request, response::Response, routing::get, BoxError, Router}; use ecstore::store::ECStore; use error::ErrorCode; use handlers::list_pools; -use object_api::ObjectApi; use tower::Service; pub type Result = std::result::Result; const API_VERSION: &str = "/v3"; -pub fn register_admin_router( - ec_store: Option, -) -> impl Service, Future: Send> + Clone { +pub fn register_admin_router() -> impl Service, Future: Send> + Clone { Router::new() .nest( "/rustfs/admin", Router::new().nest(API_VERSION, Router::new().route("/pools/list", get(list_pools::handler))), ) - .with_state::<()>(ObjectApi::new(ec_store)) .into_service() } diff --git a/api/admin/src/object_api.rs b/api/admin/src/object_api.rs deleted file mode 100644 index 2395df404..000000000 --- a/api/admin/src/object_api.rs +++ /dev/null @@ -1,20 +0,0 @@ -use std::ops::Deref; - -use ecstore::store::ECStore; - -#[derive(Clone)] -pub struct ObjectApi(Option); - -impl Deref for ObjectApi { - type Target = Option; - - fn deref(&self) -> &Self::Target { - &self.0 - } -} - -impl ObjectApi { - pub fn new(t: Option) -> Self { - Self(t) - } -} diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index feadc029e..edd96c4e9 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -146,59 +146,56 @@ async fn run(opt: config::Opt) -> Result<()> { }; let rpc_service = NodeServiceServer::with_interceptor(make_server(), check_auth); + info!(" init store success!"); + + tokio::spawn(async move { + let hyper_service = service.into_shared(); + let adm_service = admin::register_admin_router(); + + let hybrid_service = TowerToHyperService::new(hybrid(hyper_service, rpc_service, adm_service)); + + let http_server = ConnBuilder::new(TokioExecutor::new()); + let mut ctrl_c = std::pin::pin!(tokio::signal::ctrl_c()); + let graceful = hyper_util::server::graceful::GracefulShutdown::new(); + info!("server is running at http://{local_addr}"); + + loop { + let (socket, _) = tokio::select! { + res = listener.accept() => { + match res { + Ok(conn) => conn, + Err(err) => { + tracing::error!("error accepting connection: {err}"); + continue; + } + } + } + _ = ctrl_c.as_mut() => { + break; + } + }; + + let conn = http_server.serve_connection(TokioIo::new(socket), hybrid_service.clone()); + let conn = graceful.watch(conn.into_owned()); + tokio::spawn(async move { + let _ = conn.await; + }); + } + + tokio::select! { + () = graceful.shutdown() => { + tracing::debug!("Gracefully shutdown!"); + }, + () = tokio::time::sleep(std::time::Duration::from_secs(10)) => { + tracing::debug!("Waited 10 seconds for graceful shutdown, aborting..."); + } + } + }); // init store let store = ECStore::new(opt.address.clone(), endpoint_pools.clone()) .await .map_err(|err| Error::from_string(err.to_string()))?; - info!(" init store success!"); - - tokio::spawn({ - let store = store.clone(); - async move { - let hyper_service = service.into_shared(); - let adm_service = admin::register_admin_router(Some(store)); - - let hybrid_service = TowerToHyperService::new(hybrid(hyper_service, rpc_service, adm_service)); - - let http_server = ConnBuilder::new(TokioExecutor::new()); - let mut ctrl_c = std::pin::pin!(tokio::signal::ctrl_c()); - let graceful = hyper_util::server::graceful::GracefulShutdown::new(); - info!("server is running at http://{local_addr}"); - - loop { - let (socket, _) = tokio::select! { - res = listener.accept() => { - match res { - Ok(conn) => conn, - Err(err) => { - tracing::error!("error accepting connection: {err}"); - continue; - } - } - } - _ = ctrl_c.as_mut() => { - break; - } - }; - - let conn = http_server.serve_connection(TokioIo::new(socket), hybrid_service.clone()); - let conn = graceful.watch(conn.into_owned()); - tokio::spawn(async move { - let _ = conn.await; - }); - } - - tokio::select! { - () = graceful.shutdown() => { - tracing::debug!("Gracefully shutdown!"); - }, - () = tokio::time::sleep(std::time::Duration::from_secs(10)) => { - tracing::debug!("Waited 10 seconds for graceful shutdown, aborting..."); - } - } - } - }); let buckets_list = store .list_bucket(&BucketOptions {