diff --git a/docs/architecture/migration-progress.md b/docs/architecture/migration-progress.md index 2bcf6533c..d3d639038 100644 --- a/docs/architecture/migration-progress.md +++ b/docs/architecture/migration-progress.md @@ -5,16 +5,16 @@ Status values: `[ ]` not started, `[~]` in progress, `[x]` complete, `[!]` block ## Current Context - Issue: [`rustfs/backlog#660`](https://github.com/rustfs/backlog/issues/660) -- Branch: `overtrue/arch-app-usecase-object-store-context` -- Baseline: `origin/main` at `bed875d3e0d114398da8b2f94d4473293cbc45a8` +- Branch: `overtrue/arch-admin-object-store-context` +- Baseline: `upstream/main` at `6fcd62ba5648f90f6a78a0819f54e4a881cc2a87` - PR type for this branch: `consumer-migration` -- Runtime behavior changes: no external behavior change expected; app usecases - prefer the AppContext-owned object store and keep the existing global - fallback. -- Rust code changes: migrate admin, bucket, multipart, and object usecase - object-store lookups to AppContext-first resolution. +- Runtime behavior changes: no external behavior change expected; admin + object-store lookups prefer the AppContext-owned object store and keep the + existing global fallback. +- Rust code changes: add global fallback to the shared object-store resolver and + migrate admin handlers, services, and router helpers to that resolver. - CI/script changes: none. -- Docs changes: record `CTX-004` consumer migration scope and verification. +- Docs changes: record `CTX-005` consumer migration scope and verification. ## Phase 0 Tasks @@ -167,6 +167,13 @@ Status values: `[ ]` not started, `[~]` in progress, `[x]` complete, `[!]` block preserve the existing global object-layer fallback when absent. - Verification: formatting, compile check, migration guards, diff hygiene, Rust risk scan, and full `make pre-commit`. +- [x] `CTX-005` Migrate admin object-store consumers. + - Do: migrate admin handlers, admin services, and admin router helpers to the + shared object-store resolver. + - Acceptance: admin object-store lookups use AppContext when present and + preserve the existing global object-layer fallback when absent. + - Verification: focused resolver test, formatting, compile check, migration + guards, diff hygiene, Rust risk scan, and full `make pre-commit`. ## Phase 1 Security Governance Tasks @@ -473,25 +480,26 @@ Status values: `[ ]` not started, `[~]` in progress, `[x]` complete, `[!]` block ## Next PRs -1. `pure-move`: start `R-009` boot wrapper with the IAM degraded readiness +1. `consumer-migration`: migrate the next low-risk non-admin AppContext + consumer group while preserving global fallback. +2. `pure-move`: start `R-009` boot wrapper with the IAM degraded readiness contract covered. -2. `consumer-migration`: migrate the next low-risk AppContext consumer group - while preserving global fallback. ## Pre-Push Review Log | Expert | Status | Notes | |---|---|---| -| Quality/architecture | pass | Single `consumer-migration` slice; helpers stay private to each usecase and no service container, registry, or cross-layer dependency is added. | -| Migration preservation | pass | AppContext object-store resolution is preferred when present, existing `new_object_layer_fn` fallback and `Not init` error paths are preserved, and method bodies keep their storage behavior. | -| Testing/verification | pass | Formatting, compile check, diff hygiene, migration guards, added-line Rust risk scan, and full `make pre-commit` passed. | +| Quality/architecture | pass | Single `consumer-migration` slice; the shared object-store resolver owns fallback behavior and admin callers do not add service containers, registries, or cross-layer dependencies. | +| Migration preservation | pass | AppContext object-store resolution is preferred when present, existing global object-layer fallback and admin error paths are preserved, and method bodies keep their storage behavior. | +| Testing/verification | pass | Focused resolver test, formatting, compile check, diff hygiene, migration guards, added-line Rust risk scan, and full `make pre-commit` passed. | ## Verification Notes -Passed on `bed875d3e0d114398da8b2f94d4473293cbc45a8`: +Passed on `6fcd62ba5648f90f6a78a0819f54e4a881cc2a87`: - `cargo fmt --all`. - `cargo fmt --all --check`. - `cargo check -p rustfs --lib`. +- `cargo test -p rustfs app::context::compat::tests::resolver_helpers_are_context_first_and_fallback_when_context_is_absent --lib`. - `git diff --check`. - `./scripts/check_architecture_migration_rules.sh`. - `./scripts/check_layer_dependencies.sh`. @@ -502,12 +510,12 @@ Passed on `bed875d3e0d114398da8b2f94d4473293cbc45a8`: 111 skipped, plus doctests. Notes: -- This slice migrates one coherent app usecase consumer group. +- This slice migrates one coherent admin object-store consumer group. - Global object-layer fallback remains available until the planned cleanup phase. ## Handoff Notes -- CTX-004 is complete. +- CTX-005 is complete. - Keep the next consumer or boot wrapper PR scoped to one PR type and preserve global fallback until the planned cleanup phase. diff --git a/rustfs/src/admin/handlers/account_info.rs b/rustfs/src/admin/handlers/account_info.rs index 875b1fbcd..bdc1ebc38 100644 --- a/rustfs/src/admin/handlers/account_info.rs +++ b/rustfs/src/admin/handlers/account_info.rs @@ -14,6 +14,7 @@ use crate::admin::auth::authenticate_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; +use crate::app::context::resolve_object_store_handle; use crate::auth::get_condition_values; use crate::server::{ADMIN_PREFIX, RemoteAddr}; use http::{HeaderMap, HeaderValue}; @@ -21,7 +22,6 @@ use hyper::{Method, StatusCode}; use matchit::Params; use rustfs_credentials::get_global_action_cred; use rustfs_ecstore::bucket::versioning_sys::BucketVersioningSys; -use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::store_api::BucketOperations; use rustfs_policy::policy::BucketPolicy; use rustfs_policy::policy::default::DEFAULT_POLICIES; @@ -61,7 +61,7 @@ fn resolve_bucket_access(can_list_bucket: bool, can_get_bucket_location: bool, c #[async_trait::async_trait] impl Operation for AccountInfoHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; diff --git a/rustfs/src/admin/handlers/audit_runtime_config.rs b/rustfs/src/admin/handlers/audit_runtime_config.rs index a8b6b2ca4..9fcc223e2 100644 --- a/rustfs/src/admin/handlers/audit_runtime_config.rs +++ b/rustfs/src/admin/handlers/audit_runtime_config.rs @@ -13,13 +13,14 @@ // limitations under the License. use crate::admin::handlers::target_descriptor::AdminTargetSpec; +use crate::app::context::resolve_object_store_handle; use rustfs_audit::{audit_system, start_audit_system as start_global_audit_system, system::AuditSystemState}; use rustfs_config::DEFAULT_DELIMITER; use rustfs_config::server_config::Config; use s3s::{S3Result, s3_error}; pub(crate) async fn load_server_config_from_store() -> S3Result { - let Some(store) = rustfs_ecstore::global::new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Ok(Config::new()); }; @@ -77,7 +78,7 @@ pub(crate) async fn update_audit_config_and_reload(specs: &[AdminTargetSpec], where F: FnMut(&mut Config) -> bool, { - let Some(store) = rustfs_ecstore::global::new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "server storage not initialized")); }; diff --git a/rustfs/src/admin/handlers/bucket_meta.rs b/rustfs/src/admin/handlers/bucket_meta.rs index 28bc629de..ab9e078fb 100644 --- a/rustfs/src/admin/handlers/bucket_meta.rs +++ b/rustfs/src/admin/handlers/bucket_meta.rs @@ -17,6 +17,7 @@ use crate::{ auth::validate_admin_request, router::{AdminOperation, Operation, S3Router}, }, + app::context::resolve_object_store_handle, auth::{check_key_valid, get_session_token}, server::{ADMIN_PREFIX, RemoteAddr}, }; @@ -37,7 +38,6 @@ use rustfs_ecstore::{ target::BucketTargets, }, error::StorageError, - new_object_layer_fn, store_api::BucketOperations, }; use rustfs_policy::policy::{ @@ -134,7 +134,7 @@ impl Operation for ExportBucketMetadata { ) .await?; - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "object store is not initialized")); }; @@ -520,7 +520,7 @@ impl Operation for ImportBucketMetadata { }; } - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "object store is not initialized")); }; diff --git a/rustfs/src/admin/handlers/config_admin.rs b/rustfs/src/admin/handlers/config_admin.rs index df732c512..acae5ed4c 100644 --- a/rustfs/src/admin/handlers/config_admin.rs +++ b/rustfs/src/admin/handlers/config_admin.rs @@ -19,6 +19,7 @@ use crate::admin::service::config::{ validate_server_config, }; use crate::admin::utils::{encode_compatible_admin_payload, is_compat_admin_request, read_compatible_admin_body}; +use crate::app::context::resolve_object_store_handle; use crate::auth::{check_key_valid, get_session_token}; use crate::error::ApiError; use crate::server::{ADMIN_PREFIX, RemoteAddr}; @@ -72,7 +73,6 @@ use rustfs_ecstore::config::com::STORAGE_CLASS_SUB_SYS; use rustfs_ecstore::config::com::{delete_config, read_config, read_config_without_migrate, save_config, save_server_config}; use rustfs_ecstore::config::storageclass::{INLINE_BLOCK_ENV, OPTIMIZE_ENV, RRS_ENV, STANDARD_ENV}; use rustfs_ecstore::disk::RUSTFS_META_BUCKET; -use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::store_api::ListOperations; use rustfs_policy::policy::action::{Action, AdminAction}; use s3s::header::CONTENT_TYPE; @@ -710,7 +710,7 @@ fn success_response(config_applied: bool) -> S3Result S3Result> { - new_object_layer_fn().ok_or_else(|| s3_error!(InternalError, "server storage not initialized")) + resolve_object_store_handle().ok_or_else(|| s3_error!(InternalError, "server storage not initialized")) } async fn load_server_config_from_store() -> S3Result { diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index ac1521202..431c59fdf 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -14,6 +14,7 @@ use crate::admin::auth::{authenticate_request, validate_admin_request}; use crate::admin::router::{AdminOperation, Operation, S3Router}; +use crate::app::context::resolve_object_store_handle; use crate::server::ADMIN_PREFIX; use crate::server::RemoteAddr; use bytes::Bytes; @@ -23,7 +24,6 @@ use matchit::Params; use rustfs_common::heal_channel::{HealChannelPriority, HealChannelRequest, HealOpts}; use rustfs_config::MAX_HEAL_REQUEST_SIZE; use rustfs_ecstore::bucket::utils::is_valid_object_prefix; -use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::store_api::HealOperations; use rustfs_ecstore::store_utils::is_reserved_or_invalid_bucket; use rustfs_policy::policy::action::{Action, AdminAction}; @@ -370,7 +370,7 @@ impl Operation for HealHandler { // The heal channel currently models bucket/object work. Root heal reuses the // existing format-heal path directly so `/v3/heal/` is accepted intentionally. if should_handle_root_heal_directly(&hip) { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "server not initialized")); }; @@ -560,7 +560,7 @@ impl Operation for BackgroundHealStatusHandler { warn!("handle BackgroundHealStatusHandler"); validate_heal_admin_request(&req).await?; - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "server not initialized")); }; diff --git a/rustfs/src/admin/handlers/kms_dynamic.rs b/rustfs/src/admin/handlers/kms_dynamic.rs index a65b1b140..559870d10 100644 --- a/rustfs/src/admin/handlers/kms_dynamic.rs +++ b/rustfs/src/admin/handlers/kms_dynamic.rs @@ -16,14 +16,13 @@ use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; -use crate::app::context::resolve_kms_runtime_service_manager; +use crate::app::context::{resolve_kms_runtime_service_manager, resolve_object_store_handle}; use crate::auth::{check_key_valid, get_session_token}; use crate::server::{ADMIN_PREFIX, RemoteAddr}; use hyper::{Method, StatusCode}; use matchit::Params; use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; use rustfs_ecstore::config::com::{read_config, save_config}; -use rustfs_ecstore::new_object_layer_fn; use rustfs_kms::{ ConfigureKmsRequest, ConfigureKmsResponse, KmsConfig, KmsConfigSummary, KmsServiceStatus, KmsStatusResponse, StartKmsRequest, StartKmsResponse, StopKmsResponse, @@ -104,7 +103,7 @@ fn normalize_configure_request_auth( /// Save KMS configuration to cluster storage #[instrument(skip(config))] async fn save_kms_config(config: &KmsConfig) -> Result<(), String> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err("Storage layer not initialized".to_string()); }; @@ -128,7 +127,7 @@ async fn save_kms_config(config: &KmsConfig) -> Result<(), String> { /// Load KMS configuration from cluster storage #[instrument] pub async fn load_kms_config() -> Option { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { warn!( component = LOG_COMPONENT_ADMIN, subsystem = LOG_SUBSYSTEM_KMS, diff --git a/rustfs/src/admin/handlers/oidc.rs b/rustfs/src/admin/handlers/oidc.rs index 178af0819..847e77421 100644 --- a/rustfs/src/admin/handlers/oidc.rs +++ b/rustfs/src/admin/handlers/oidc.rs @@ -15,6 +15,7 @@ use super::sts::create_oidc_sts_credentials; use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; +use crate::app::context::resolve_object_store_handle; use crate::auth::{check_key_valid, get_session_token}; use crate::server::{ADMIN_PREFIX, MINIO_ADMIN_PREFIX, RemoteAddr}; use http::StatusCode; @@ -30,7 +31,6 @@ use rustfs_config::server_config::Config as ServerConfig; use rustfs_config::server_config::get_global_server_config; use rustfs_config::{DEFAULT_DELIMITER, ENABLE_KEY, EnableState, MAX_ADMIN_REQUEST_BODY_SIZE}; use rustfs_ecstore::config::com::{read_config_without_migrate, save_server_config}; -use rustfs_ecstore::new_object_layer_fn; use rustfs_policy::policy::action::{Action, AdminAction}; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; use serde::de::DeserializeOwned; @@ -792,7 +792,7 @@ fn json_response(status: StatusCode, payload: &T) -> S3Result S3Result { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "storage layer not initialized")); }; @@ -802,7 +802,7 @@ async fn load_server_config_from_store() -> S3Result { } async fn save_server_config_to_store(config: &ServerConfig) -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "storage layer not initialized")); }; diff --git a/rustfs/src/admin/handlers/pools.rs b/rustfs/src/admin/handlers/pools.rs index 64695b69c..15689ea21 100644 --- a/rustfs/src/admin/handlers/pools.rs +++ b/rustfs/src/admin/handlers/pools.rs @@ -27,13 +27,12 @@ use crate::{ router::{AdminOperation, Operation, S3Router}, }, app::admin_usecase::{DefaultAdminUsecase, QueryPoolStatusRequest}, - app::context::resolve_endpoints_handle, + app::context::{resolve_endpoints_handle, resolve_object_store_handle}, auth::{check_key_valid, get_session_token}, error::ApiError, server::{ADMIN_PREFIX, RemoteAddr}, }; use hyper::Method; -use rustfs_ecstore::new_object_layer_fn; use std::collections::HashSet; @@ -284,7 +283,7 @@ impl Operation for StartDecommission { return Err(s3_error!(NotImplemented)); } - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(decommission_admin_not_initialized_error("start decommission")); }; @@ -399,7 +398,7 @@ impl Operation for CancelDecommission { return Err(pool_admin_pool_not_found_error("cancel decommission", &query.pool)); }; - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(decommission_admin_not_initialized_error("cancel decommission")); }; diff --git a/rustfs/src/admin/handlers/rebalance.rs b/rustfs/src/admin/handlers/rebalance.rs index f4ec14823..8f8476d1d 100644 --- a/rustfs/src/admin/handlers/rebalance.rs +++ b/rustfs/src/admin/handlers/rebalance.rs @@ -17,6 +17,7 @@ use crate::{ auth::validate_admin_request, router::{AdminOperation, Operation, S3Router}, }, + app::context::resolve_object_store_handle, auth::{check_key_valid, get_session_token}, server::{ADMIN_PREFIX, RemoteAddr}, }; @@ -26,7 +27,6 @@ use matchit::Params; use rustfs_ecstore::rebalance::RebalanceMeta; use rustfs_ecstore::{ error::StorageError, - new_object_layer_fn, notification_sys::get_global_notification_sys, rebalance::{DiskStat, RebalSaveOpt}, store_api::BucketOperations, @@ -249,7 +249,7 @@ impl Operation for RebalanceStart { ) .await?; - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "object layer is not initialized")); }; @@ -367,7 +367,7 @@ impl Operation for RebalanceStatus { ) .await?; - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "object layer is not initialized")); }; @@ -452,7 +452,7 @@ impl Operation for RebalanceStop { ) .await?; - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(s3_error!(InternalError, "object layer is not initialized")); }; diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index 60304b07b..fc40b2653 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -16,6 +16,7 @@ use crate::admin::auth::validate_admin_request; use crate::admin::handlers::site_replication::site_replication_peer_deployment_id_for_endpoint; use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::utils::read_compatible_admin_body; +use crate::app::context::resolve_object_store_handle; use crate::auth::{check_key_valid, get_session_token}; use crate::error::ApiError; use crate::server::{ADMIN_PREFIX, RemoteAddr}; @@ -33,7 +34,6 @@ use rustfs_ecstore::bucket::replication::GLOBAL_REPLICATION_STATS; use rustfs_ecstore::bucket::target::BucketTarget; use rustfs_ecstore::error::StorageError; use rustfs_ecstore::global::global_rustfs_port; -use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::store_api::BucketOperations; use rustfs_policy::policy::action::{Action, AdminAction}; use rustfs_storage_api::BucketOptions; @@ -136,7 +136,7 @@ impl Operation for GetReplicationMetricsHandler { return Err(s3_error!(InvalidRequest, "bucket is required")); } - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; @@ -194,7 +194,7 @@ impl Operation for SetRemoteTargetHandler { return Err(s3_error!(InvalidRequest, "bucket is required")); } - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; @@ -329,7 +329,7 @@ impl Operation for ListRemoteTargetHandler { return Err(s3_error!(InvalidRequest, "bucket is required")); } - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not initialized".to_string())); }; @@ -389,7 +389,7 @@ impl Operation for RemoveRemoteTargetHandler { return Err(s3_error!(InvalidRequest, "arn is required")); }; - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not initialized".to_string())); }; diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 8f3d68e56..0315ad9d6 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -19,6 +19,7 @@ use crate::admin::site_replication_identity::{ site_identity_key, }; use crate::admin::utils::{encode_compatible_admin_payload, read_compatible_admin_body}; +use crate::app::context::resolve_object_store_handle; use crate::auth::{check_key_valid, get_session_token}; use crate::config::get_config_snapshot; use crate::error::ApiError; @@ -49,7 +50,6 @@ use rustfs_ecstore::bucket::versioning::VersioningApi; use rustfs_ecstore::config::com::{delete_config, read_config, save_config}; use rustfs_ecstore::error::Error as StorageError; use rustfs_ecstore::global::{get_global_deployment_id, get_global_endpoints_opt, get_global_region, global_rustfs_port}; -use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::store_api::BucketOperations; use rustfs_iam::error::is_err_no_such_service_account; use rustfs_iam::store::{MappedPolicy, UserType}; @@ -482,7 +482,7 @@ async fn read_site_replication_json( } async fn load_site_replication_state() -> S3Result { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; @@ -511,7 +511,7 @@ async fn load_site_replication_state() -> S3Result { } async fn save_site_replication_state(state: &SiteReplicationState) -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; @@ -527,7 +527,7 @@ async fn save_site_replication_state(state: &SiteReplicationState) -> S3Result<( } async fn clear_site_replication_state() -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; @@ -1244,7 +1244,7 @@ pub async fn site_replication_make_bucket_hook(bucket: &str, lock_enabled: bool) ensure_site_replication_bucket_targets(bucket, &state, &local_peer, None).await?; ensure_site_replication_bucket_replication_config(bucket, &state, &local_peer).await?; - let created_at = new_object_layer_fn() + let created_at = resolve_object_store_handle() .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))? .get_bucket_info(bucket, &BucketOptions::default()) .await @@ -1312,7 +1312,7 @@ fn maybe_time(value: OffsetDateTime) -> Option { } async fn build_sr_info(state: &SiteReplicationState, local_peer: &PeerInfo) -> S3Result { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; @@ -2542,7 +2542,7 @@ pub async fn site_replication_peer_deployment_id_for_endpoint(endpoint: &str) -> /// kick a resync toward every remote peer so pre-existing objects back-fill. Errors are logged /// but never abort the caller — the admin can run a manual resync if needed. async fn backfill_existing_buckets_after_add(state: &SiteReplicationState, local_peer: &PeerInfo) { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return; }; let buckets = match store.list_bucket(&BucketOptions::default()).await { @@ -2847,7 +2847,7 @@ async fn ensure_site_replication_bucket_versioning(bucket: &str) -> S3Result<()> } async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; @@ -3447,7 +3447,7 @@ impl Operation for SRPeerBucketOpsHandler { .cloned() .ok_or_else(|| s3_error!(InvalidRequest, "operation is required"))?; - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; @@ -3704,7 +3704,7 @@ impl Operation for SiteReplicationResyncOpHandler { if !state.peers.contains_key(&peer.deployment_id) { return Err(s3_error!(InvalidRequest, "site replication peer not found")); } - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let buckets = store.list_bucket(&BucketOptions::default()).await.map_err(ApiError::from)?; diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index bf332d851..47bc0bcf4 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -16,6 +16,7 @@ use crate::admin::{ auth::{AdminResourceScope, validate_admin_request, validate_admin_request_with_bucket_object}, router::{AdminOperation, Operation, S3Router}, }; +use crate::app::context::resolve_object_store_handle; use crate::auth::{check_key_valid, get_session_token}; use crate::server::{RemoteAddr, TABLE_CATALOG_COMPAT_PREFIX, TABLE_CATALOG_PREFIX}; use crate::table_catalog::{DEFAULT_WAREHOUSE_ID, TableCatalogStore}; @@ -23,7 +24,7 @@ use http::{HeaderMap, HeaderValue, StatusCode}; use hyper::Method; use matchit::Params; use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; -use rustfs_ecstore::{bucket::metadata_sys, new_object_layer_fn, store::ECStore}; +use rustfs_ecstore::{bucket::metadata_sys, store::ECStore}; use rustfs_iam::{manager::get_token_signing_key, sys::SESSION_POLICY_NAME}; use rustfs_policy::{ auth::get_new_credentials_with_metadata, @@ -778,7 +779,7 @@ fn job_id_from_params(params: &Params<'_, '_>) -> S3Result { } fn table_catalog_backend() -> S3Result> { - let store = new_object_layer_fn().ok_or_else(|| s3_error!(InternalError, "object store not initialized"))?; + let store = resolve_object_store_handle().ok_or_else(|| s3_error!(InternalError, "object store not initialized"))?; Ok(crate::table_catalog::EcStoreTableCatalogObjectBackend::new(store)) } diff --git a/rustfs/src/admin/router.rs b/rustfs/src/admin/router.rs index 98e2c744b..dc5bc2ae6 100644 --- a/rustfs/src/admin/router.rs +++ b/rustfs/src/admin/router.rs @@ -14,6 +14,7 @@ use crate::admin::console::{is_console_path, make_console_server}; use crate::admin::handlers::oidc::is_oidc_path; +use crate::app::context::resolve_object_store_handle; use crate::app::object_usecase::DefaultObjectUsecase; use crate::auth::{check_key_valid, get_session_token}; use crate::error::ApiError; @@ -57,13 +58,10 @@ use rustfs_ecstore::bucket::versioning::VersioningApi; use rustfs_ecstore::bucket::versioning_sys::BucketVersioningSys; use rustfs_ecstore::config::com::read_config_without_migrate; use rustfs_ecstore::global::GLOBAL_BOOT_TIME; +use rustfs_ecstore::global::{get_global_bucket_monitor, get_global_deployment_id, get_global_region}; use rustfs_ecstore::notification_sys::get_global_notification_sys; use rustfs_ecstore::rpc::PeerRestClient; use rustfs_ecstore::store_api::BucketOperations; -use rustfs_ecstore::{ - global::{get_global_bucket_monitor, get_global_deployment_id, get_global_region}, - new_object_layer_fn, -}; use rustfs_filemeta::{ReplicationStatusType, ReplicationType}; use rustfs_madmin::utils::parse_duration; use rustfs_notify::{Event as NotificationEvent, notification_system}; @@ -630,7 +628,7 @@ async fn load_current_server_config() -> S3Result { return Ok(system.config.read().await.clone()); } - if let Some(store) = new_object_layer_fn() { + if let Some(store) = resolve_object_store_handle() { match read_config_without_migrate(store).await { Ok(config) => return Ok(config), Err(err) => { @@ -1398,7 +1396,7 @@ fn build_listen_notification_response(uri: &Uri, bucket: Option<&str>) -> S3Resu } async fn ensure_replication_bucket_exists(bucket: &str) -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init")); }; @@ -2275,7 +2273,7 @@ async fn handle_misc_extension_request(req: &mut S3Request, route: &MiscEx } MiscExtRoute::ListenNotification { bucket } => { if let Some(bucket_name) = bucket { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init")); }; store diff --git a/rustfs/src/admin/service/config.rs b/rustfs/src/admin/service/config.rs index 6afde6683..9abbaea2d 100644 --- a/rustfs/src/admin/service/config.rs +++ b/rustfs/src/admin/service/config.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use crate::app::context::resolve_object_store_handle; use rustfs_audit::reload_audit_config; use rustfs_config::audit::{AUDIT_MQTT_SUB_SYS, AUDIT_REDIS_DEFAULT_CHANNEL, AUDIT_WEBHOOK_SUB_SYS}; use rustfs_config::notify::{NOTIFY_MQTT_SUB_SYS, NOTIFY_REDIS_DEFAULT_CHANNEL, NOTIFY_WEBHOOK_SUB_SYS}; @@ -23,7 +24,6 @@ use rustfs_config::{HEAL_SUB_SYS, SCANNER_SUB_SYS}; use rustfs_ecstore::config::com::{STORAGE_CLASS_SUB_SYS, read_config_without_migrate}; use rustfs_ecstore::config::set_global_storage_class; use rustfs_ecstore::config::storageclass; -use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::notification_sys::get_global_notification_sys; use rustfs_iam::oidc::load_oidc_provider_configs_from_server_config; use rustfs_storage_api::StorageAdminApi; @@ -67,7 +67,7 @@ fn invalid_request(message: impl Into) -> S3Error { } async fn apply_storage_class_runtime_config(config: &ServerConfig) -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(internal_error("storage layer not initialized")); }; @@ -92,7 +92,7 @@ fn validate_storage_class_kvs(kvs: &KVS, set_drive_counts: &[usize]) -> S3Result } async fn validate_storage_class_config(config: &ServerConfig) -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(internal_error("storage layer not initialized")); }; @@ -297,7 +297,7 @@ pub async fn reload_dynamic_config_runtime_state(sub_system: &str) -> S3Result<( return Err(internal_error(format!("unsupported dynamic config subsystem: {sub_system}"))); } - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(internal_error("storage layer not initialized")); }; @@ -313,7 +313,7 @@ pub async fn reload_dynamic_config_runtime_state(sub_system: &str) -> S3Result<( } pub async fn reload_runtime_config_snapshot() -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(internal_error("storage layer not initialized")); }; diff --git a/rustfs/src/admin/service/site_replication.rs b/rustfs/src/admin/service/site_replication.rs index 71acbcaa7..005a0bee8 100644 --- a/rustfs/src/admin/service/site_replication.rs +++ b/rustfs/src/admin/service/site_replication.rs @@ -13,9 +13,9 @@ // limitations under the License. use crate::admin::site_replication_identity::{deployment_id_for_endpoint, normalize_peer_map_by_identity_with}; +use crate::app::context::resolve_object_store_handle; use rustfs_ecstore::config::com::{read_config, save_config}; use rustfs_ecstore::error::Error as StorageError; -use rustfs_ecstore::new_object_layer_fn; use rustfs_madmin::PeerInfo; use s3s::{S3Error, S3ErrorCode, S3Result}; use serde_json::{Map, Value}; @@ -91,7 +91,7 @@ fn normalize_site_replication_state_json(data: &[u8]) -> Result>, /// RustFS does not currently keep a separate in-memory cache for this state, /// so "reload" means validating that the persisted JSON is readable. pub async fn reload_site_replication_runtime_state() -> S3Result<()> { - let Some(store) = new_object_layer_fn() else { + let Some(store) = resolve_object_store_handle() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; diff --git a/rustfs/src/app/context/compat.rs b/rustfs/src/app/context/compat.rs index 2c447242e..131dd47a5 100644 --- a/rustfs/src/app/context/compat.rs +++ b/rustfs/src/app/context/compat.rs @@ -22,6 +22,7 @@ use crate::config::RustFSBufferConfig; use rustfs_config::server_config::Config; use rustfs_ecstore::bucket::metadata_sys::BucketMetadataSys; use rustfs_ecstore::endpoints::EndpointServerPools; +use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::store::ECStore; use rustfs_ecstore::tier::tier::TierConfigMgr; use rustfs_kms::KmsServiceManager; @@ -38,9 +39,9 @@ pub fn resolve_bucket_metadata_handle() -> Option> resolve_bucket_metadata_handle_with(get_global_app_context(), || default_bucket_metadata_interface().handle()) } -/// Resolve object store handle from AppContext. +/// Resolve object store handle using AppContext-first precedence. pub fn resolve_object_store_handle() -> Option> { - resolve_object_store_handle_with(get_global_app_context()) + resolve_object_store_handle_with(get_global_app_context(), new_object_layer_fn) } /// Resolve endpoints using AppContext-first precedence. @@ -76,8 +77,11 @@ fn resolve_bucket_metadata_handle_with( .or_else(fallback) } -fn resolve_object_store_handle_with(context: Option>) -> Option> { - context.map(|context| context.object_store()) +fn resolve_object_store_handle_with( + context: Option>, + fallback: impl FnOnce() -> Option>, +) -> Option> { + context.map(|context| context.object_store()).or_else(fallback) } fn resolve_endpoints_handle_with( @@ -295,7 +299,7 @@ mod tests { &bucket_metadata )); assert!(Arc::ptr_eq( - &resolve_object_store_handle_with(Some(context.clone())).expect("context object store"), + &resolve_object_store_handle_with(Some(context.clone()), || None).expect("context object store"), &object_store )); assert_eq!( @@ -326,7 +330,10 @@ mod tests { &resolve_bucket_metadata_handle_with(None, || Some(bucket_metadata.clone())).expect("fallback bucket metadata"), &bucket_metadata )); - assert!(resolve_object_store_handle_with(None).is_none()); + assert!(Arc::ptr_eq( + &resolve_object_store_handle_with(None, || Some(object_store.clone())).expect("fallback object store"), + &object_store + )); assert_eq!( resolve_endpoints_handle_with(None, || Some(endpoints.clone())) .expect("fallback endpoints")