fix: blocking forever

Closes #102
Signed-off-by: bestgopher <84328409@qq.com>
This commit is contained in:
bestgopher
2024-10-29 20:48:33 +08:00
parent dd9abc1db0
commit 16d66a4206
4 changed files with 54 additions and 77 deletions
+8 -3
View File
@@ -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<ObjectApi>) -> LocalResult<Json<Vec<PoolStatus>>> {
pub async fn handler() -> LocalResult<Json<Vec<PoolStatus>>> {
// 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()接口获取每个池的数据
//
+1 -6
View File
@@ -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<T> = std::result::Result<T, ErrorCode>;
const API_VERSION: &str = "/v3";
pub fn register_admin_router(
ec_store: Option<ECStore>,
) -> impl Service<Request, Response = Response, Error: Into<BoxError>, Future: Send> + Clone {
pub fn register_admin_router() -> impl Service<Request, Response = Response, Error: Into<BoxError>, 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()
}
-20
View File
@@ -1,20 +0,0 @@
use std::ops::Deref;
use ecstore::store::ECStore;
#[derive(Clone)]
pub struct ObjectApi(Option<ECStore>);
impl Deref for ObjectApi {
type Target = Option<ECStore>;
fn deref(&self) -> &Self::Target {
&self.0
}
}
impl ObjectApi {
pub fn new(t: Option<ECStore>) -> Self {
Self(t)
}
}
+45 -48
View File
@@ -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 {