refactor(app): centralize context resolvers for admin/server paths (#1975)

This commit is contained in:
安正超
2026-02-26 20:41:11 +08:00
committed by GitHub
parent dafb31d208
commit 2c85721654
8 changed files with 55 additions and 37 deletions
+2 -2
View File
@@ -27,7 +27,7 @@ use crate::{
router::{AdminOperation, Operation, S3Router}, router::{AdminOperation, Operation, S3Router},
}, },
app::admin_usecase::{DefaultAdminUsecase, QueryPoolStatusRequest}, app::admin_usecase::{DefaultAdminUsecase, QueryPoolStatusRequest},
app::context::get_global_app_context, app::context::resolve_endpoints_handle,
auth::{check_key_valid, get_session_token}, auth::{check_key_valid, get_session_token},
error::ApiError, error::ApiError,
server::{ADMIN_PREFIX, RemoteAddr}, server::{ADMIN_PREFIX, RemoteAddr},
@@ -36,7 +36,7 @@ use hyper::Method;
use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::new_object_layer_fn;
fn endpoints_from_context() -> Option<rustfs_ecstore::endpoints::EndpointServerPools> { fn endpoints_from_context() -> Option<rustfs_ecstore::endpoints::EndpointServerPools> {
get_global_app_context().and_then(|context| context.endpoints().handle()) resolve_endpoints_handle()
} }
pub fn register_pool_route(r: &mut S3Router<AdminOperation>) -> std::io::Result<()> { pub fn register_pool_route(r: &mut S3Router<AdminOperation>) -> std::io::Result<()> {
+3 -3
View File
@@ -16,7 +16,7 @@
use crate::admin::auth::{validate_admin_request, validate_admin_request_with_bucket}; use crate::admin::auth::{validate_admin_request, validate_admin_request_with_bucket};
use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::router::{AdminOperation, Operation, S3Router};
use crate::app::context::get_global_app_context; use crate::app::context::{resolve_bucket_metadata_handle, resolve_object_store_handle};
use crate::auth::{check_key_valid, get_session_token}; use crate::auth::{check_key_valid, get_session_token};
use crate::server::ADMIN_PREFIX; use crate::server::ADMIN_PREFIX;
use hyper::{Method, StatusCode}; use hyper::{Method, StatusCode};
@@ -86,11 +86,11 @@ pub struct GetBucketQuotaStatsHandler;
pub struct CheckBucketQuotaHandler; pub struct CheckBucketQuotaHandler;
fn bucket_metadata_from_context() -> Option<Arc<RwLock<BucketMetadataSys>>> { fn bucket_metadata_from_context() -> Option<Arc<RwLock<BucketMetadataSys>>> {
get_global_app_context().and_then(|context| context.bucket_metadata().handle()) resolve_bucket_metadata_handle()
} }
async fn current_usage_from_context(bucket: &str) -> u64 { async fn current_usage_from_context(bucket: &str) -> u64 {
let Some(store) = get_global_app_context().map(|context| context.object_store()) else { let Some(store) = resolve_object_store_handle() else {
return 0; return 0;
}; };
+8 -16
View File
@@ -18,7 +18,7 @@ use crate::{
auth::validate_admin_request, auth::validate_admin_request,
router::{AdminOperation, Operation, S3Router}, router::{AdminOperation, Operation, S3Router},
}, },
app::context::{default_tier_config_interface, get_global_app_context}, app::context::resolve_tier_config_handle,
auth::{check_key_valid, get_session_token}, auth::{check_key_valid, get_session_token},
server::{ADMIN_PREFIX, RemoteAddr}, server::{ADMIN_PREFIX, RemoteAddr},
}; };
@@ -45,17 +45,9 @@ use s3s::{
s3_error, s3_error,
}; };
use serde_urlencoded::from_bytes; use serde_urlencoded::from_bytes;
use std::sync::Arc;
use time::OffsetDateTime; use time::OffsetDateTime;
use tokio::sync::RwLock;
use tracing::{debug, warn}; use tracing::{debug, warn};
fn tier_config_mgr_from_context() -> Arc<RwLock<rustfs_ecstore::tier::tier::TierConfigMgr>> {
get_global_app_context()
.map(|context| context.tier_config().handle())
.unwrap_or_else(|| default_tier_config_interface().handle())
}
#[derive(Debug, Clone, serde::Deserialize, Default)] #[derive(Debug, Clone, serde::Deserialize, Default)]
pub struct AddTierQuery { pub struct AddTierQuery {
#[serde(rename = "accessKey")] #[serde(rename = "accessKey")]
@@ -214,7 +206,7 @@ impl Operation for AddTier {
&_ => (), &_ => (),
} }
let tier_config_mgr_handle = tier_config_mgr_from_context(); let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await; let mut tier_config_mgr = tier_config_mgr_handle.write().await;
//tier_config_mgr.reload(api); //tier_config_mgr.reload(api);
if let Err(err) = tier_config_mgr.add(args, force).await { if let Err(err) = tier_config_mgr.add(args, force).await {
@@ -307,7 +299,7 @@ impl Operation for EditTier {
let tier_name = params.get("tiername").map(|s| s.to_string()).unwrap_or_default(); let tier_name = params.get("tiername").map(|s| s.to_string()).unwrap_or_default();
let tier_config_mgr_handle = tier_config_mgr_from_context(); let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await; let mut tier_config_mgr = tier_config_mgr_handle.write().await;
//tier_config_mgr.reload(api); //tier_config_mgr.reload(api);
if let Err(err) = tier_config_mgr.edit(&tier_name, creds).await { if let Err(err) = tier_config_mgr.edit(&tier_name, creds).await {
@@ -375,7 +367,7 @@ impl Operation for ListTiers {
) )
.await?; .await?;
let tier_config_mgr_handle = tier_config_mgr_from_context(); let tier_config_mgr_handle = resolve_tier_config_handle();
let tier_config_mgr = tier_config_mgr_handle.read().await; let tier_config_mgr = tier_config_mgr_handle.read().await;
let tiers = tier_config_mgr.list_tiers(); let tiers = tier_config_mgr.list_tiers();
@@ -431,7 +423,7 @@ impl Operation for RemoveTier {
let tier_name = params.get("tiername").map(|s| s.to_string()).unwrap_or_default(); let tier_name = params.get("tiername").map(|s| s.to_string()).unwrap_or_default();
let tier_config_mgr_handle = tier_config_mgr_from_context(); let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await; let mut tier_config_mgr = tier_config_mgr_handle.write().await;
//tier_config_mgr.reload(api); //tier_config_mgr.reload(api);
if let Err(err) = tier_config_mgr.remove(&tier_name, force).await { if let Err(err) = tier_config_mgr.remove(&tier_name, force).await {
@@ -492,7 +484,7 @@ impl Operation for VerifyTier {
) )
.await?; .await?;
let tier_config_mgr_handle = tier_config_mgr_from_context(); let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await; let mut tier_config_mgr = tier_config_mgr_handle.write().await;
tier_config_mgr.verify(&query.tier.unwrap()).await; tier_config_mgr.verify(&query.tier.unwrap()).await;
@@ -534,7 +526,7 @@ impl Operation for GetTierInfo {
} }
}; };
let tier_config_mgr_handle = tier_config_mgr_from_context(); let tier_config_mgr_handle = resolve_tier_config_handle();
let tier_config_mgr = tier_config_mgr_handle.read().await; let tier_config_mgr = tier_config_mgr_handle.read().await;
let info = tier_config_mgr.get(&query.tier.unwrap()); let info = tier_config_mgr.get(&query.tier.unwrap());
@@ -601,7 +593,7 @@ impl Operation for ClearTier {
return Err(s3_error!(InvalidRequest, "get rand failed")); return Err(s3_error!(InvalidRequest, "get rand failed"));
}; };
let tier_config_mgr_handle = tier_config_mgr_from_context(); let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await; let mut tier_config_mgr = tier_config_mgr_handle.write().await;
//tier_config_mgr.reload(api); //tier_config_mgr.reload(api);
if let Err(err) = tier_config_mgr.clear_tier(force).await { if let Err(err) = tier_config_mgr.clear_tier(force).await {
+2 -2
View File
@@ -13,7 +13,7 @@
// limitations under the License. // limitations under the License.
use crate::admin::router::Operation; use crate::admin::router::Operation;
use crate::app::context::get_global_app_context; use crate::app::context::resolve_endpoints_handle;
use http::StatusCode; use http::StatusCode;
use hyper::Uri; use hyper::Uri;
use matchit::Params; use matchit::Params;
@@ -43,7 +43,7 @@ impl Operation for Trace {
let _trace_opts = extract_trace_options(&req.uri)?; let _trace_opts = extract_trace_options(&req.uri)?;
// let (tx, rx) = mpsc::channel(10000); // let (tx, rx) = mpsc::channel(10000);
let _peers = match get_global_app_context().and_then(|context| context.endpoints().handle()) { let _peers = match resolve_endpoints_handle() {
Some(ep) => PeerRestClient::new_clients(ep.clone()).await, Some(ep) => PeerRestClient::new_clients(ep.clone()).await,
None => (Vec::new(), Vec::new()), None => (Vec::new(), Vec::new()),
}; };
+34
View File
@@ -326,6 +326,40 @@ pub fn resolve_kms_runtime_service_manager() -> Option<Arc<KmsServiceManager>> {
.or_else(|| default_kms_runtime_interface().service_manager()) .or_else(|| default_kms_runtime_interface().service_manager())
} }
/// Resolve bucket metadata handle using AppContext-first precedence.
pub fn resolve_bucket_metadata_handle() -> Option<Arc<RwLock<BucketMetadataSys>>> {
get_global_app_context()
.and_then(|context| context.bucket_metadata().handle())
.or_else(|| default_bucket_metadata_interface().handle())
}
/// Resolve object store handle from AppContext.
pub fn resolve_object_store_handle() -> Option<Arc<ECStore>> {
get_global_app_context().map(|context| context.object_store())
}
/// Resolve endpoints using AppContext-first precedence.
pub fn resolve_endpoints_handle() -> Option<EndpointServerPools> {
get_global_app_context()
.and_then(|context| context.endpoints().handle())
.or_else(|| default_endpoints_interface().handle())
}
/// Resolve tier config handle using AppContext-first precedence.
pub fn resolve_tier_config_handle() -> Arc<RwLock<TierConfigMgr>> {
get_global_app_context()
.map(|context| context.tier_config().handle())
.unwrap_or_else(|| default_tier_config_interface().handle())
}
/// Resolve server config using AppContext-first precedence.
pub fn resolve_server_config() -> Option<Config> {
match get_global_app_context() {
Some(context) => context.server_config().get(),
None => default_server_config_interface().get(),
}
}
pub fn default_bucket_metadata_interface() -> Arc<dyn BucketMetadataInterface> { pub fn default_bucket_metadata_interface() -> Arc<dyn BucketMetadataInterface> {
Arc::new(BucketMetadataHandle) Arc::new(BucketMetadataHandle)
} }
+2 -5
View File
@@ -12,16 +12,13 @@
// See the License for the specific language governing permissions and // See the License for the specific language governing permissions and
// limitations under the License. // limitations under the License.
use crate::app::context::{default_server_config_interface, get_global_app_context}; use crate::app::context::resolve_server_config;
use rustfs_audit::{AuditError, AuditResult, audit_system, init_audit_system, system::AuditSystemState}; use rustfs_audit::{AuditError, AuditResult, audit_system, init_audit_system, system::AuditSystemState};
use rustfs_config::DEFAULT_DELIMITER; use rustfs_config::DEFAULT_DELIMITER;
use tracing::{info, warn}; use tracing::{info, warn};
fn server_config_from_context() -> Option<rustfs_ecstore::config::Config> { fn server_config_from_context() -> Option<rustfs_ecstore::config::Config> {
match get_global_app_context() { resolve_server_config()
Some(context) => context.server_config().get(),
None => default_server_config_interface().get(),
}
} }
/// Start the audit system. /// Start the audit system.
+2 -5
View File
@@ -12,15 +12,12 @@
// See the License for the specific language governing permissions and // See the License for the specific language governing permissions and
// limitations under the License. // limitations under the License.
use crate::app::context::{default_server_config_interface, get_global_app_context}; use crate::app::context::resolve_server_config;
use rustfs_config::DEFAULT_DELIMITER; use rustfs_config::DEFAULT_DELIMITER;
use tracing::{error, info, instrument, warn}; use tracing::{error, info, instrument, warn};
fn server_config_from_context() -> Option<rustfs_ecstore::config::Config> { fn server_config_from_context() -> Option<rustfs_ecstore::config::Config> {
match get_global_app_context() { resolve_server_config()
Some(context) => context.server_config().get(),
None => default_server_config_interface().get(),
}
} }
/// Shuts down the event notifier system gracefully /// Shuts down the event notifier system gracefully
+2 -4
View File
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and // See the License for the specific language governing permissions and
// limitations under the License. // limitations under the License.
use crate::app::context::{default_bucket_metadata_interface, get_global_app_context}; use crate::app::context::resolve_bucket_metadata_handle;
use crate::error::ApiError; use crate::error::ApiError;
use crate::storage::concurrency::get_concurrency_manager; use crate::storage::concurrency::get_concurrency_manager;
use crate::storage::helper::OperationHelper; use crate::storage::helper::OperationHelper;
@@ -53,9 +53,7 @@ use tokio_util::io::StreamReader;
use tracing::{debug, error, instrument, warn}; use tracing::{debug, error, instrument, warn};
fn bucket_metadata_for_quota() -> Option<Arc<RwLock<metadata_sys::BucketMetadataSys>>> { fn bucket_metadata_for_quota() -> Option<Arc<RwLock<metadata_sys::BucketMetadataSys>>> {
get_global_app_context() resolve_bucket_metadata_handle()
.and_then(|context| context.bucket_metadata().handle())
.or_else(|| default_bucket_metadata_interface().handle())
} }
impl Objects { impl Objects {