mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-10 07:06:53 +00:00
refactor: align app runtime facade helpers (#3947)
This commit is contained in:
@@ -25,7 +25,7 @@ use super::storage_api::admin_usecase::data_usage::{
|
||||
};
|
||||
use super::storage_api::admin_usecase::{ECStore, EndpointServerPools};
|
||||
use crate::app::runtime_sources::{
|
||||
AppContext, current_app_context, resolve_endpoints_handle, resolve_object_store_handle_for_context,
|
||||
AppContext, current_app_context, current_endpoints_handle, current_object_store_handle_for_context,
|
||||
};
|
||||
use crate::capacity::resolve_admin_used_capacity;
|
||||
use crate::cluster_snapshot::{
|
||||
@@ -196,7 +196,7 @@ impl DefaultAdminUsecase {
|
||||
}
|
||||
|
||||
fn object_store(&self) -> Option<Arc<ECStore>> {
|
||||
resolve_object_store_handle_for_context(self.context.as_deref())
|
||||
current_object_store_handle_for_context(self.context.as_deref())
|
||||
}
|
||||
|
||||
fn app_error(code: S3ErrorCode, message: impl Into<String>) -> ApiError {
|
||||
@@ -583,7 +583,7 @@ impl DefaultAdminUsecase {
|
||||
|
||||
pub async fn execute_collect_dependency_readiness(&self) -> DependencyReadiness {
|
||||
let report = collect_runtime_dependency_readiness_report().await;
|
||||
if let Some(endpoint_pools) = resolve_endpoints_handle() {
|
||||
if let Some(endpoint_pools) = current_endpoints_handle() {
|
||||
let runtime_status = ClusterRuntimeStatusSnapshot::from_readiness_report(report);
|
||||
return cluster_read_only_snapshot_from_endpoint_pools(&endpoint_pools, runtime_status)
|
||||
.runtime_status
|
||||
@@ -593,7 +593,7 @@ impl DefaultAdminUsecase {
|
||||
}
|
||||
|
||||
pub async fn execute_collect_cluster_read_only_snapshot(&self) -> Option<ClusterReadOnlySnapshot> {
|
||||
let endpoint_pools = resolve_endpoints_handle()?;
|
||||
let endpoint_pools = current_endpoints_handle()?;
|
||||
collect_cluster_read_only_snapshot(&endpoint_pools).await
|
||||
}
|
||||
|
||||
|
||||
@@ -60,8 +60,8 @@ use crate::admin::handlers::site_replication::{
|
||||
site_replication_bucket_meta_hook, site_replication_delete_bucket_hook, site_replication_make_bucket_hook,
|
||||
};
|
||||
use crate::app::runtime_sources::{
|
||||
AppContext, current_app_context, resolve_encryption_service, resolve_notification_system,
|
||||
resolve_notify_interface_for_context, resolve_object_store_handle_for_context,
|
||||
AppContext, current_app_context, current_encryption_service, current_notification_system,
|
||||
current_notify_interface_for_context, current_object_store_handle_for_context,
|
||||
};
|
||||
use crate::auth::get_condition_values_with_client_info;
|
||||
use crate::error::ApiError;
|
||||
@@ -238,7 +238,7 @@ fn notify_bucket_metadata_reload(
|
||||
request_context: Option<request_context::RequestContext>,
|
||||
) {
|
||||
spawn_background_with_context(request_context, async move {
|
||||
if let Some(notification_sys) = resolve_notification_system()
|
||||
if let Some(notification_sys) = current_notification_system()
|
||||
&& let Err(err) = notification_sys.load_bucket_metadata(&bucket).await
|
||||
{
|
||||
warn!(bucket = %bucket, error = %err, "failed to notify peers after {operation}");
|
||||
@@ -778,7 +778,7 @@ impl DefaultBucketUsecase {
|
||||
}
|
||||
|
||||
fn object_store(&self) -> Option<Arc<ECStore>> {
|
||||
resolve_object_store_handle_for_context(self.context.as_deref())
|
||||
current_object_store_handle_for_context(self.context.as_deref())
|
||||
}
|
||||
|
||||
#[instrument(
|
||||
@@ -1594,7 +1594,7 @@ impl DefaultBucketUsecase {
|
||||
&& by_default.sse_algorithm.as_str() == ServerSideEncryption::AWS_KMS
|
||||
&& by_default.kms_master_key_id.as_deref().is_none_or(str::is_empty)
|
||||
{
|
||||
let service = resolve_encryption_service()
|
||||
let service = current_encryption_service()
|
||||
.await
|
||||
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "KMS service not initialized".to_string()))?;
|
||||
let default_key = service
|
||||
@@ -1746,7 +1746,7 @@ impl DefaultBucketUsecase {
|
||||
.map_err(ApiError::from)?;
|
||||
|
||||
let region = resolve_notification_region(self.global_region(), request_region);
|
||||
let notify = resolve_notify_interface_for_context(self.context.as_deref());
|
||||
let notify = current_notify_interface_for_context(self.context.as_deref());
|
||||
let clear_rules = notify.clear_bucket_notification_rules(&bucket);
|
||||
let parse_rules = async {
|
||||
let mut event_rules = Vec::new();
|
||||
|
||||
@@ -33,7 +33,7 @@ use super::storage_api::test::{
|
||||
};
|
||||
use super::{multipart_usecase::DefaultMultipartUsecase, object_usecase::DefaultObjectUsecase};
|
||||
use crate::app::bucket_usecase::DefaultBucketUsecase;
|
||||
use crate::app::runtime_sources::resolve_tier_config_handle;
|
||||
use crate::app::runtime_sources::current_tier_config_handle;
|
||||
use bytes::Bytes;
|
||||
use futures::FutureExt;
|
||||
use futures::stream;
|
||||
@@ -343,7 +343,7 @@ impl AppWarmBackend for MockWarmBackend {
|
||||
|
||||
async fn register_mock_tier(tier_name: &str) -> MockWarmBackend {
|
||||
let backend = MockWarmBackend::default();
|
||||
let tier_config_mgr_handle = resolve_tier_config_handle();
|
||||
let tier_config_mgr_handle = current_tier_config_handle();
|
||||
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
|
||||
tier_config_mgr.tiers.insert(
|
||||
tier_name.to_string(),
|
||||
|
||||
@@ -56,7 +56,7 @@ use super::storage_api::multipart_usecase::sse::{
|
||||
};
|
||||
use super::storage_api::multipart_usecase::{StorageObjectOptions as ObjectOptions, StoragePutObjReader as PutObjReader};
|
||||
use crate::app::object_usecase::{build_put_like_object_lock_metadata, validate_existing_object_lock_for_write};
|
||||
use crate::app::runtime_sources::{AppContext, current_app_context, resolve_object_store_handle_for_context};
|
||||
use crate::app::runtime_sources::{AppContext, current_app_context, current_object_store_handle_for_context};
|
||||
use crate::capacity::record_capacity_write;
|
||||
use crate::error::ApiError;
|
||||
use crate::table_catalog;
|
||||
@@ -295,7 +295,7 @@ impl DefaultMultipartUsecase {
|
||||
}
|
||||
|
||||
fn object_store(&self) -> Option<Arc<ECStore>> {
|
||||
resolve_object_store_handle_for_context(self.context.as_deref())
|
||||
current_object_store_handle_for_context(self.context.as_deref())
|
||||
}
|
||||
|
||||
#[instrument(level = "debug", skip(self))]
|
||||
|
||||
@@ -90,8 +90,8 @@ use super::storage_api::object_usecase::{
|
||||
validate_sse_headers_for_write, validate_ssec_for_read, wrap_response_with_cors,
|
||||
};
|
||||
use crate::app::runtime_sources::{
|
||||
AppContext, current_app_context, resolve_expiry_state_handle, resolve_notify_interface_for_context,
|
||||
resolve_object_store_handle_for_context,
|
||||
AppContext, current_app_context, current_expiry_state_handle, current_notify_interface_for_context,
|
||||
current_object_store_handle_for_context,
|
||||
};
|
||||
use crate::config::RustFSBufferConfig;
|
||||
use crate::delete_tail_activity::{DeleteTailActivityGuard, DeleteTailStage};
|
||||
@@ -376,7 +376,7 @@ async fn enqueue_transitioned_delete_cleanup(
|
||||
|
||||
tier_delete_journal::persist_tier_delete_journal_entry(store, &je).await?;
|
||||
|
||||
let expiry_state = resolve_expiry_state_handle();
|
||||
let expiry_state = current_expiry_state_handle();
|
||||
let mut expiry_state = expiry_state.write().await;
|
||||
if let Err(err) = expiry_state.enqueue_tier_journal_entry(&je).await {
|
||||
warn!(
|
||||
@@ -1845,7 +1845,7 @@ impl DefaultObjectUsecase {
|
||||
}
|
||||
|
||||
fn object_store(&self) -> Option<Arc<ECStore>> {
|
||||
resolve_object_store_handle_for_context(self.context.as_deref())
|
||||
current_object_store_handle_for_context(self.context.as_deref())
|
||||
}
|
||||
|
||||
fn base_buffer_size(&self) -> usize {
|
||||
@@ -4313,7 +4313,7 @@ impl DefaultObjectUsecase {
|
||||
}
|
||||
|
||||
let req_headers = req.headers.clone();
|
||||
let notify = resolve_notify_interface_for_context(self.context.as_deref());
|
||||
let notify = current_notify_interface_for_context(self.context.as_deref());
|
||||
let request_context = req.extensions.get::<request_context::RequestContext>().cloned();
|
||||
let deleted_any = delete_results.iter().any(|result| result.delete_object.is_some());
|
||||
let notify_bucket = bucket.clone();
|
||||
@@ -5318,7 +5318,7 @@ impl DefaultObjectUsecase {
|
||||
None => String::new(),
|
||||
};
|
||||
|
||||
let notify = resolve_notify_interface_for_context(self.context.as_deref());
|
||||
let notify = current_notify_interface_for_context(self.context.as_deref());
|
||||
let req_params = extract_params_header(&req.headers);
|
||||
let host = get_request_host(&req.headers);
|
||||
let port = get_request_port(&req.headers);
|
||||
|
||||
@@ -13,13 +13,16 @@
|
||||
// limitations under the License.
|
||||
|
||||
pub(crate) use crate::runtime_sources::{
|
||||
AppContext, resolve_encryption_service, resolve_endpoints_handle, resolve_expiry_state_handle, resolve_notification_system,
|
||||
resolve_notify_interface_for_context, resolve_object_store_handle_for_context, resolve_s3select_db,
|
||||
AppContext, resolve_encryption_service as current_encryption_service, resolve_endpoints_handle as current_endpoints_handle,
|
||||
resolve_expiry_state_handle as current_expiry_state_handle, resolve_notification_system as current_notification_system,
|
||||
resolve_notify_interface_for_context as current_notify_interface_for_context,
|
||||
resolve_object_store_handle_for_context as current_object_store_handle_for_context,
|
||||
resolve_s3select_db as current_s3select_db,
|
||||
};
|
||||
use std::sync::Arc;
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) use crate::runtime_sources::resolve_tier_config_handle;
|
||||
pub(crate) use crate::runtime_sources::resolve_tier_config_handle as current_tier_config_handle;
|
||||
|
||||
pub(crate) fn current_app_context() -> Option<Arc<AppContext>> {
|
||||
crate::runtime_sources::current_app_context()
|
||||
|
||||
@@ -2,7 +2,7 @@ use super::storage_api::select_object::contract::object::ObjectOperations as _;
|
||||
use super::storage_api::select_object::options::get_opts;
|
||||
use super::storage_api::select_object::request_context::spawn_traced;
|
||||
use super::storage_api::select_object::{get_validated_store, validate_sse_headers_for_read, validate_ssec_for_read};
|
||||
use crate::app::runtime_sources::resolve_s3select_db;
|
||||
use crate::app::runtime_sources::current_s3select_db;
|
||||
use crate::error::ApiError;
|
||||
use bytes::Bytes;
|
||||
use datafusion::arrow::{
|
||||
@@ -60,7 +60,7 @@ pub async fn execute_select_object_content(
|
||||
validate_scan_range_for_object_size(&input.request, metadata.size)?;
|
||||
|
||||
let input = Arc::new(input);
|
||||
let db = resolve_s3select_db((*input).clone(), false)
|
||||
let db = current_s3select_db((*input).clone(), false)
|
||||
.await
|
||||
.map_err(map_query_error_to_s3)?;
|
||||
let query = Query::new(Context { input: input.clone() }, input.request.expression.clone());
|
||||
|
||||
Reference in New Issue
Block a user