diff --git a/docs/architecture/migration-progress.md b/docs/architecture/migration-progress.md index 4037b8d42..2a456e3e9 100644 --- a/docs/architecture/migration-progress.md +++ b/docs/architecture/migration-progress.md @@ -5,13 +5,14 @@ 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-admin-topology-context-resolvers` -- Baseline: completed `C-011/C-012/C-013/API-055/API-059/API-079/API-080/API-081/API-082/API-083/API-084/API-085/API-086/API-087/API-088/API-089/API-090/API-091/API-092/API-093/API-094/API-095/API-096/API-097/API-098/API-099/API-100/API-101/API-102/API-103/API-104/API-105/API-106/API-107/API-108/API-109/API-110/API-111/API-112/API-113/API-114/API-115/API-116/API-117/API-118/API-119/API-120/API-121/API-122/API-123/API-124/API-125/API-126/API-127/API-128/API-129/API-130/API-131/API-132/API-133/API-134/API-135/API-136/API-137/API-138/API-139/API-140/API-141/API-142/API-143/API-144/API-145/API-146/API-147/API-148/API-149/API-150/API-151/API-152/API-153/API-154/API-155/API-156/API-157/API-158/API-159/API-160`. -- Based on: API-160 branch while PR #3773 is pending. +- Branch: `overtrue/arch-admin-peer-systems-context-resolvers` +- Baseline: completed `C-011/C-012/C-013/API-055/API-059/API-079/API-080/API-081/API-082/API-083/API-084/API-085/API-086/API-087/API-088/API-089/API-090/API-091/API-092/API-093/API-094/API-095/API-096/API-097/API-098/API-099/API-100/API-101/API-102/API-103/API-104/API-105/API-106/API-107/API-108/API-109/API-110/API-111/API-112/API-113/API-114/API-115/API-116/API-117/API-118/API-119/API-120/API-121/API-122/API-123/API-124/API-125/API-126/API-127/API-128/API-129/API-130/API-131/API-132/API-133/API-134/API-135/API-136/API-137/API-138/API-139/API-140/API-141/API-142/API-143/API-144/API-145/API-146/API-147/API-148/API-149/API-150/API-151/API-152/API-153/API-154/API-155/API-156/API-157/API-158/API-159/API-160/API-161`. +- Based on: API-161 branch while PR #3774 is pending. - PR type for this branch: `consumer-migration` - Runtime behavior changes: none. -- Rust code changes: route admin topology endpoint, deployment-id, and runtime - port consumers through AppContext resolvers with legacy global fallback. +- Rust code changes: route admin and app notification-system, bucket-monitor, + and replication-pool consumers through AppContext resolvers with legacy + global fallback. - CI/script changes: lock completed owner and test/fuzz boundaries against bare/glob imports, scattered raw ECStore facade subpaths, and startup runtime/root-server/table/S3/app shared/app bucket/app ECStore/admin facade @@ -20,7 +21,7 @@ Status values: `[ ]` not started, `[~]` in progress, `[x]` complete, `[!]` block and storage owner thin bridge regressions, plus app context and notify event-bridge thin module regressions; accept the reviewed AppContext resolver reverse dependencies in the layer baseline. -- Docs changes: record the API-136/API-137/API-138/API-139/API-140/API-141/API-142/API-143/API-144/API-145/API-146/API-147/API-148/API-149/API-150/API-151/API-152/API-153/API-154/API-155/API-156/API-157/API-158/API-159/API-160/API-161 owner facade cleanup. +- Docs changes: record the API-136/API-137/API-138/API-139/API-140/API-141/API-142/API-143/API-144/API-145/API-146/API-147/API-148/API-149/API-150/API-151/API-152/API-153/API-154/API-155/API-156/API-157/API-158/API-159/API-160/API-161/API-162 owner facade cleanup. ## Phase 0 Tasks @@ -4327,6 +4328,23 @@ Status values: `[ ]` not started, `[~]` in progress, `[x]` complete, `[!]` block migration guard, formatting, diff hygiene, Rust risk scan, branch freshness check, pre-commit, and three-expert review. +- [x] `API-162` Route admin peer-system reads through AppContext resolvers. + - Do: add notification-system, bucket-monitor, and replication-pool + AppContext interfaces, default legacy adapters, and route admin and app + peer-system consumers through resolvers. + - Acceptance: admin tier/rebalance/config/router/site-replication and app + bucket metadata reload paths no longer directly read notification system, + bucket monitor, or replication pool globals outside AppContext default + adapters. + - Must preserve: tier config propagation, dynamic config reload propagation, + config snapshot refresh, live event peer listing, replication bandwidth + metric collection, replication resync start/status/cancel behavior, + rebalance stop fallback, and legacy global fallback when AppContext is + absent. + - Verification: RustFS compile coverage, targeted context resolver tests, + migration guard, formatting, diff hygiene, residual global-read scan, Rust + risk scan, branch freshness check, and three-expert review. + ## Next PRs 1. `consumer-migration`: continue reducing direct global reads behind AppContext resolver boundaries. @@ -4366,11 +4384,30 @@ Status values: `[ ]` not started, `[~]` in progress, `[x]` complete, `[!]` block | Migration preservation | pass | Site replication endpoint inference, same-target checks, same-deployment stale target detection, and runtime-port fallback keep legacy behavior. | | Testing/verification | pass | RustFS focused compile, targeted context tests, formatting, migration guard, diff hygiene, Rust risk scan, and pre-commit passed for API-161. | | Testing/verification | pass | CI follow-up: layer dependency baseline accepts the reviewed AppContext resolver reverse dependencies, and the layer guard passes. | +| Quality/architecture | pass | API-162 keeps admin peer-system notification, bucket-monitor, and replication-pool reads behind AppContext resolver boundaries. | +| Migration preservation | pass | Tier/rebalance/config propagation, live event peers, replication metrics, and resync operations keep legacy fallback behavior. | +| Testing/verification | pass | RustFS focused compile, targeted context tests, formatting, migration guard, diff hygiene, residual scan, Rust risk scan, and pre-commit passed for API-162. | ## Verification Notes Passed before push: +- Issue #660 API-162 current slice: + - `cargo check --tests -p rustfs`: passed. + - `cargo test -p rustfs resolver_helpers_are_context_first_and_fallback_when_context_is_absent --lib`: + passed. + - `cargo fmt --all`: passed. + - `cargo fmt --all --check`: passed. + - `git diff --check`: passed. + - `bash -n scripts/check_architecture_migration_rules.sh`: passed. + - `./scripts/check_architecture_migration_rules.sh`: passed. + - `make pre-commit`: passed. + - AppContext peer-system resolver scan: passed; direct admin and app + notification-system, bucket-monitor, and replication-pool global reads are + isolated to AppContext default adapters. + - Rust risk scan: no new production unwrap/expect, panic/todo/unsafe, or cast + risks added. + - Issue #660 API-161 current slice: - `cargo check --tests -p rustfs`: passed. - `cargo test -p rustfs resolver_helpers_are_context_first_and_fallback_when_context_is_absent --lib`: diff --git a/rustfs/src/admin/handlers/rebalance.rs b/rustfs/src/admin/handlers/rebalance.rs index 4a80c8d39..b7597008b 100644 --- a/rustfs/src/admin/handlers/rebalance.rs +++ b/rustfs/src/admin/handlers/rebalance.rs @@ -14,14 +14,14 @@ use super::super::{ DiskStat, ECStore, NotificationSys, RebalSaveOpt, RebalanceCleanupWarnings, RebalanceMeta, RebalanceStopPropagationRecord, - StorageError, decode_rebalance_stop_propagation_record, get_global_notification_sys, + StorageError, decode_rebalance_stop_propagation_record, }; use crate::{ admin::{ auth::validate_admin_request, router::{AdminOperation, Operation, S3Router}, }, - app::context::resolve_object_store_handle, + app::context::{resolve_notification_system, resolve_object_store_handle}, auth::{check_key_valid, get_session_token}, server::{ADMIN_PREFIX, RemoteAddr}, }; @@ -496,7 +496,7 @@ impl Operation for RebalanceStart { rebalance_id = %id, "admin rebalance state" ); - if let Some(notification_sys) = get_global_notification_sys() { + if let Some(notification_sys) = resolve_notification_system() { info!( event = EVENT_ADMIN_REBALANCE_STATE, component = LOG_COMPONENT_ADMIN, @@ -747,7 +747,7 @@ impl Operation for RebalanceStop { return Err(s3_error!(NoSuchResource, "pool rebalance is not started")); } - let notification_sys = get_global_notification_sys(); + let notification_sys = resolve_notification_system(); let stop_attempt_at = OffsetDateTime::now_utc(); let mut stop_failures = Vec::new(); if let Some(notification_sys) = notification_sys { diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index e3a6973ca..68f160dc0 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -20,8 +20,7 @@ use super::super::metadata::{ BUCKET_SSECONFIG, BUCKET_TAGGING_CONFIG, BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG, }; use super::super::metadata_sys; -use super::super::replication::GLOBAL_REPLICATION_STATS; -use super::super::replication::{ResyncOpts, get_global_replication_pool}; +use super::super::replication::{GLOBAL_REPLICATION_STATS, ResyncOpts}; use super::super::target::{ARN, BucketTarget, BucketTargetType, BucketTargets, Credentials}; use super::super::{AdminReplicationConfigExt as _, AdminVersioningConfigExt as _}; use super::super::{delete_admin_config, read_admin_config, save_admin_config}; @@ -33,8 +32,8 @@ use crate::admin::site_replication_identity::{ }; use crate::admin::utils::{encode_compatible_admin_payload, read_compatible_admin_body}; use crate::app::context::{ - resolve_deployment_id, resolve_endpoints_handle, resolve_object_store_handle, resolve_region, resolve_runtime_port, - resolve_server_config, + resolve_deployment_id, resolve_endpoints_handle, resolve_object_store_handle, resolve_region, + resolve_replication_pool_handle, resolve_runtime_port, resolve_server_config, }; use crate::auth::{check_key_valid, get_session_token}; use crate::config::get_config_snapshot; @@ -3219,7 +3218,7 @@ async fn start_site_bucket_resync(bucket: &str, peer: &PeerInfo, resync_id: &str } BucketTargetSys::get().update_all_targets(bucket, Some(&targets)).await; - let Some(pool) = get_global_replication_pool() else { + let Some(pool) = resolve_replication_pool_handle() else { bucket_status.status = "failed".to_string(); bucket_status.err_detail = "replication pool is not initialized".to_string(); return bucket_status; @@ -3287,7 +3286,7 @@ async fn cancel_site_bucket_resync(bucket: &str, peer: &PeerInfo, resync_id: &st } BucketTargetSys::get().update_all_targets(bucket, Some(&targets)).await; - let Some(pool) = get_global_replication_pool() else { + let Some(pool) = resolve_replication_pool_handle() else { bucket_status.status = "failed".to_string(); bucket_status.err_detail = "replication pool is not initialized".to_string(); return bucket_status; diff --git a/rustfs/src/admin/handlers/tier.rs b/rustfs/src/admin/handlers/tier.rs index a0ce70206..3298f8197 100644 --- a/rustfs/src/admin/handlers/tier.rs +++ b/rustfs/src/admin/handlers/tier.rs @@ -17,14 +17,14 @@ use super::super::lifecycle::bucket_lifecycle_ops::GLOBAL_TransitionState; use super::super::{ AdminError, DailyAllTierStats, ERR_TIER_ALREADY_EXISTS, ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY, ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_MISSING_CREDENTIALS, ERR_TIER_NAME_NOT_UPPERCASE, - ERR_TIER_NOT_FOUND, TierConfig, TierCreds, TierType, get_global_notification_sys, storageclass, + ERR_TIER_NOT_FOUND, TierConfig, TierCreds, TierType, storageclass, }; use crate::{ admin::{ auth::validate_admin_request, router::{AdminOperation, Operation, S3Router}, }, - app::context::{resolve_object_store_handle, resolve_tier_config_handle}, + app::context::{resolve_notification_system, resolve_object_store_handle, resolve_tier_config_handle}, auth::{check_key_valid, get_session_token}, server::{ADMIN_PREFIX, RemoteAddr}, storage::request_context::spawn_traced, @@ -80,7 +80,7 @@ pub struct AddTierQuery { pub struct AddTier {} fn spawn_transition_tier_config_propagation(action: &'static str) { - if let Some(notification_sys) = get_global_notification_sys() { + if let Some(notification_sys) = resolve_notification_system() { spawn_traced(async move { for peer_result in notification_sys.load_transition_tier_config().await { if let Some(err) = peer_result.err { diff --git a/rustfs/src/admin/mod.rs b/rustfs/src/admin/mod.rs index 711e3c9a9..fc692d96a 100644 --- a/rustfs/src/admin/mod.rs +++ b/rustfs/src/admin/mod.rs @@ -136,7 +136,7 @@ mod ecstore_error { } mod ecstore_global { - pub(crate) use crate::storage::ecstore_global::{GLOBAL_BOOT_TIME, get_global_bucket_monitor}; + pub(crate) use crate::storage::ecstore_global::GLOBAL_BOOT_TIME; } #[allow(unused_imports)] @@ -149,7 +149,7 @@ mod ecstore_metrics { } mod ecstore_notification { - pub(crate) use crate::storage::ecstore_notification::{NotificationSys, get_global_notification_sys}; + pub(crate) use crate::storage::ecstore_notification::NotificationSys; } #[allow(unused_imports)] @@ -250,7 +250,6 @@ impl AdminVersioningConfigExt for s3s::dto::VersioningConfiguration { pub(crate) mod bandwidth { pub(crate) mod monitor { pub(crate) type BandwidthDetails = super::super::ecstore_bucket::bandwidth::monitor::BandwidthDetails; - pub(crate) type Monitor = super::super::ecstore_bucket::bandwidth::monitor::Monitor; } } @@ -399,7 +398,6 @@ pub(crate) mod replication { pub(crate) type BucketReplicationResyncStatus = super::ecstore_bucket::replication::BucketReplicationResyncStatus; pub(crate) type BucketStats = super::ecstore_bucket::replication::BucketStats; - pub(crate) type DynReplicationPool = super::ecstore_bucket::replication::DynReplicationPool; pub(crate) type ObjectOpts = super::ecstore_bucket::replication::ObjectOpts; pub(crate) type ReplicationStats = super::ecstore_bucket::replication::ReplicationStats; pub(crate) type ResyncOpts = super::ecstore_bucket::replication::ResyncOpts; @@ -417,10 +415,6 @@ pub(crate) mod replication { super::ecstore_bucket::replication::GLOBAL_REPLICATION_STATS.get() } } - - pub(crate) fn get_global_replication_pool() -> Option> { - super::ecstore_bucket::replication::get_global_replication_pool() - } } pub(crate) mod target { @@ -510,10 +504,6 @@ pub(crate) async fn load_data_usage_from_backend( ecstore_data_usage::load_data_usage_from_backend(store).await } -pub(crate) fn get_global_bucket_monitor() -> Option> { - ecstore_global::get_global_bucket_monitor() -} - pub(crate) async fn collect_local_metrics( types: MetricType, opts: &CollectMetricsOpts, @@ -521,10 +511,6 @@ pub(crate) async fn collect_local_metrics( ecstore_metrics::collect_local_metrics(types, opts).await } -pub(crate) fn get_global_notification_sys() -> Option<&'static NotificationSys> { - ecstore_notification::get_global_notification_sys() -} - pub(crate) struct BootTimeCompat; pub(crate) static GLOBAL_BOOT_TIME: BootTimeCompat = BootTimeCompat; diff --git a/rustfs/src/admin/router.rs b/rustfs/src/admin/router.rs index 02442841a..6144efdc9 100644 --- a/rustfs/src/admin/router.rs +++ b/rustfs/src/admin/router.rs @@ -16,20 +16,19 @@ use super::GLOBAL_BOOT_TIME; use super::PeerRestClient; use super::bandwidth::monitor::BandwidthDetails; use super::bucket_target_sys::{BucketTargetSys, PutObjectOptions, RemoveObjectOptions, S3ClientError, TargetClient}; -use super::get_global_bucket_monitor; -use super::get_global_notification_sys; use super::metadata::BUCKET_TARGETS_FILE; use super::metadata_sys; use super::read_admin_config_without_migrate; -use super::replication::{ - BucketReplicationResyncStatus, BucketStats, GLOBAL_REPLICATION_STATS, ObjectOpts, ResyncOpts, get_global_replication_pool, -}; +use super::replication::{BucketReplicationResyncStatus, BucketStats, GLOBAL_REPLICATION_STATS, ObjectOpts, ResyncOpts}; use super::target::{BucketTarget, BucketTargetType, BucketTargets}; use super::versioning_sys::BucketVersioningSys; use super::{AdminReplicationConfigExt as _, AdminVersioningConfigExt as _}; use crate::admin::console::{is_console_path, make_console_server}; use crate::admin::handlers::oidc::is_oidc_path; -use crate::app::context::{resolve_deployment_id, resolve_object_store_handle, resolve_region, resolve_server_config}; +use crate::app::context::{ + resolve_bucket_monitor_handle, resolve_deployment_id, resolve_notification_system, resolve_object_store_handle, + resolve_region, resolve_replication_pool_handle, resolve_server_config, +}; use crate::app::object_usecase::DefaultObjectUsecase; use crate::auth::{check_key_valid, get_session_token}; use crate::error::ApiError; @@ -1189,7 +1188,7 @@ fn serialize_listen_notification_event(event: &NotificationEvent) -> S3Result Vec { - get_global_notification_sys() + resolve_notification_system() .map(|system| { system .peer_clients @@ -1446,7 +1445,7 @@ fn replication_metrics_uptime_seconds() -> i64 { } fn collect_replication_metrics_bandwidth(bucket: &str) -> HashMap { - get_global_bucket_monitor() + resolve_bucket_monitor_handle() .map(|monitor| { monitor .get_report(|name| name == bucket) @@ -2150,7 +2149,7 @@ async fn start_replication_resync(bucket: &str, reset: &ReplicationResetStartReq .map_err(ApiError::from)?; BucketTargetSys::get().update_all_targets(bucket, Some(&targets)).await; - let Some(pool) = get_global_replication_pool() else { + let Some(pool) = resolve_replication_pool_handle() else { return Err(s3_error!(InternalError, "replication pool is not initialized")); }; @@ -2170,7 +2169,7 @@ async fn start_replication_resync(bucket: &str, reset: &ReplicationResetStartReq } async fn load_replication_resync_status(bucket: &str) -> S3Result { - let Some(pool) = get_global_replication_pool() else { + let Some(pool) = resolve_replication_pool_handle() else { return Err(s3_error!(InternalError, "replication pool is not initialized")); }; diff --git a/rustfs/src/admin/service/config.rs b/rustfs/src/admin/service/config.rs index 5c847972d..62f20a088 100644 --- a/rustfs/src/admin/service/config.rs +++ b/rustfs/src/admin/service/config.rs @@ -12,11 +12,10 @@ // See the License for the specific language governing permissions and // limitations under the License. -use super::super::get_global_notification_sys; use super::super::set_global_storage_class; use super::super::storageclass; use super::super::{STORAGE_CLASS_SUB_SYS, read_admin_config_without_migrate}; -use crate::app::context::resolve_object_store_handle; +use crate::app::context::{resolve_notification_system, 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}; @@ -345,7 +344,7 @@ pub async fn signal_dynamic_config_reload(sub_system: &str) { return; } - let Some(notification_sys) = get_global_notification_sys() else { + let Some(notification_sys) = resolve_notification_system() else { return; }; @@ -357,7 +356,7 @@ pub async fn signal_dynamic_config_reload(sub_system: &str) { } pub async fn signal_config_snapshot_reload() { - let Some(notification_sys) = get_global_notification_sys() else { + let Some(notification_sys) = resolve_notification_system() else { return; }; diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index 2684f511f..1091d02d7 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -16,7 +16,6 @@ use super::ECStore; use super::StorageError; -use super::get_global_notification_sys; use super::object_api_utils::to_s3s_etag; use super::{AppObjectLockConfigExt as _, AppVersioningConfigExt as _}; use super::{ @@ -40,7 +39,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::context::{ - AppContext, get_global_app_context, resolve_notify_interface_for_context, resolve_object_store_handle_for_context, + AppContext, get_global_app_context, resolve_notification_system, resolve_notify_interface_for_context, + resolve_object_store_handle_for_context, }; use crate::auth::get_condition_values_with_client_info; use crate::error::ApiError; @@ -207,7 +207,7 @@ fn notify_bucket_metadata_reload( request_context: Option, ) { spawn_background_with_context(request_context, async move { - if let Some(notification_sys) = get_global_notification_sys() + if let Some(notification_sys) = resolve_notification_system() && let Err(err) = notification_sys.load_bucket_metadata(&bucket).await { warn!(bucket = %bucket, error = %err, "failed to notify peers after {operation}"); diff --git a/rustfs/src/app/context.rs b/rustfs/src/app/context.rs index 2f1b1eae8..b29cfd804 100644 --- a/rustfs/src/app/context.rs +++ b/rustfs/src/app/context.rs @@ -30,6 +30,7 @@ use super::EndpointServerPools; use super::TierConfigMgr; use super::metadata_sys::BucketMetadataSys; use super::new_object_layer_fn; +use super::{BucketBandwidthMonitor, DynReplicationPool, NotificationSys}; use crate::config::RustFSBufferConfig; use rustfs_config::server_config::Config; use rustfs_credentials::Credentials; @@ -79,11 +80,26 @@ pub fn resolve_notify_interface_for_context(context: Option<&AppContext>) -> Arc .unwrap_or_else(default_notify_interface) } +/// Resolve notification system handle using AppContext-first precedence. +pub fn resolve_notification_system() -> Option<&'static NotificationSys> { + resolve_notification_system_with(get_global_app_context(), || default_notification_system_interface().handle()) +} + /// Resolve endpoints using AppContext-first precedence. pub fn resolve_endpoints_handle() -> Option { resolve_endpoints_handle_with(get_global_app_context(), || default_endpoints_interface().handle()) } +/// Resolve bucket bandwidth monitor using AppContext-first precedence. +pub fn resolve_bucket_monitor_handle() -> Option> { + resolve_bucket_monitor_handle_with(get_global_app_context(), || default_bucket_monitor_interface().handle()) +} + +/// Resolve replication pool handle using AppContext-first precedence. +pub fn resolve_replication_pool_handle() -> Option> { + resolve_replication_pool_handle_with(get_global_app_context(), || default_replication_pool_interface().handle()) +} + /// Resolve deployment identity using AppContext-first precedence. pub fn resolve_deployment_id() -> Option { resolve_deployment_id_with(get_global_app_context(), || default_deployment_id_interface().get()) @@ -158,6 +174,33 @@ fn resolve_bucket_metadata_handle_with( .or_else(fallback) } +fn resolve_notification_system_with( + context: Option>, + fallback: impl FnOnce() -> Option<&'static NotificationSys>, +) -> Option<&'static NotificationSys> { + context + .and_then(|context| context.notification_system().handle()) + .or_else(fallback) +} + +fn resolve_bucket_monitor_handle_with( + context: Option>, + fallback: impl FnOnce() -> Option>, +) -> Option> { + context + .and_then(|context| context.bucket_monitor().handle()) + .or_else(fallback) +} + +fn resolve_replication_pool_handle_with( + context: Option>, + fallback: impl FnOnce() -> Option>, +) -> Option> { + context + .and_then(|context| context.replication_pool().handle()) + .or_else(fallback) +} + #[cfg(test)] fn resolve_object_store_handle_with( context: Option>, @@ -242,7 +285,10 @@ mod tests { use super::super::{Endpoints, PoolEndpoints}; use super::*; use crate::app::context::global::AppContextTestInterfaces; - use crate::app::context::handles::default_notify_interface; + use crate::app::context::handles::{ + default_bucket_monitor_interface, default_notification_system_interface, default_notify_interface, + default_replication_pool_interface, + }; use crate::app::context::interfaces::{ ActionCredentialInterface, BucketMetadataInterface, BufferConfigInterface, DeploymentIdInterface, EndpointsInterface, IamInterface, KmsInterface, KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, RegionInterface, @@ -490,9 +536,12 @@ mod tests { kms: Some(context_kms.clone()), }), notify: default_notify_interface(), + notification_system: default_notification_system_interface(), bucket_metadata: Arc::new(TestBucketMetadataInterface { metadata: Some(bucket_metadata.clone()), }), + bucket_monitor: default_bucket_monitor_interface(), + replication_pool: default_replication_pool_interface(), endpoints: Arc::new(TestEndpointsInterface { endpoints: Some(endpoints.clone()), }), diff --git a/rustfs/src/app/context/global.rs b/rustfs/src/app/context/global.rs index e42f54bd2..81cbebbcb 100644 --- a/rustfs/src/app/context/global.rs +++ b/rustfs/src/app/context/global.rs @@ -15,14 +15,17 @@ use super::super::{ECStore, set_object_store_resolver}; use super::handles::{ IamHandle, KmsHandle, default_action_credential_interface, default_bucket_metadata_interface, - default_buffer_config_interface, default_deployment_id_interface, default_endpoints_interface, default_kms_runtime_interface, - default_local_node_name_interface, default_lock_client_interface, default_notify_interface, default_region_interface, - default_runtime_port_interface, default_server_config_interface, default_tier_config_interface, + default_bucket_monitor_interface, default_buffer_config_interface, default_deployment_id_interface, + default_endpoints_interface, default_kms_runtime_interface, default_local_node_name_interface, default_lock_client_interface, + default_notification_system_interface, default_notify_interface, default_region_interface, + default_replication_pool_interface, default_runtime_port_interface, default_server_config_interface, + default_tier_config_interface, }; use super::interfaces::{ - ActionCredentialInterface, BucketMetadataInterface, BufferConfigInterface, DeploymentIdInterface, EndpointsInterface, - IamInterface, KmsInterface, KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, NotifyInterface, - RegionInterface, RuntimePortInterface, ServerConfigInterface, TierConfigInterface, + ActionCredentialInterface, BucketMetadataInterface, BucketMonitorInterface, BufferConfigInterface, DeploymentIdInterface, + EndpointsInterface, IamInterface, KmsInterface, KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, + NotificationSystemInterface, NotifyInterface, RegionInterface, ReplicationPoolInterface, RuntimePortInterface, + ServerConfigInterface, TierConfigInterface, }; use rustfs_iam::{store::object::ObjectStore, sys::IamSys}; use rustfs_kms::KmsServiceManager; @@ -37,7 +40,10 @@ pub struct AppContext { kms: Arc, kms_runtime: Arc, notify: Arc, + notification_system: Arc, bucket_metadata: Arc, + bucket_monitor: Arc, + replication_pool: Arc, endpoints: Arc, deployment_id: Arc, runtime_port: Arc, @@ -58,7 +64,10 @@ impl AppContext { kms, kms_runtime: default_kms_runtime_interface(), notify: default_notify_interface(), + notification_system: default_notification_system_interface(), bucket_metadata: default_bucket_metadata_interface(), + bucket_monitor: default_bucket_monitor_interface(), + replication_pool: default_replication_pool_interface(), endpoints: default_endpoints_interface(), deployment_id: default_deployment_id_interface(), runtime_port: default_runtime_port_interface(), @@ -101,10 +110,22 @@ impl AppContext { self.notify.clone() } + pub fn notification_system(&self) -> Arc { + self.notification_system.clone() + } + pub fn bucket_metadata(&self) -> Arc { self.bucket_metadata.clone() } + pub fn bucket_monitor(&self) -> Arc { + self.bucket_monitor.clone() + } + + pub fn replication_pool(&self) -> Arc { + self.replication_pool.clone() + } + pub fn endpoints(&self) -> Arc { self.endpoints.clone() } @@ -152,7 +173,10 @@ pub(super) struct AppContextTestInterfaces { pub(super) kms: Arc, pub(super) kms_runtime: Arc, pub(super) notify: Arc, + pub(super) notification_system: Arc, pub(super) bucket_metadata: Arc, + pub(super) bucket_monitor: Arc, + pub(super) replication_pool: Arc, pub(super) endpoints: Arc, pub(super) deployment_id: Arc, pub(super) runtime_port: Arc, @@ -174,7 +198,10 @@ impl AppContext { kms: interfaces.kms, kms_runtime: interfaces.kms_runtime, notify: interfaces.notify, + notification_system: interfaces.notification_system, bucket_metadata: interfaces.bucket_metadata, + bucket_monitor: interfaces.bucket_monitor, + replication_pool: interfaces.replication_pool, endpoints: interfaces.endpoints, deployment_id: interfaces.deployment_id, runtime_port: interfaces.runtime_port, diff --git a/rustfs/src/app/context/handles.rs b/rustfs/src/app/context/handles.rs index df7025430..2c6935097 100644 --- a/rustfs/src/app/context/handles.rs +++ b/rustfs/src/app/context/handles.rs @@ -16,13 +16,14 @@ use super::super::EndpointServerPools; use super::super::TierConfigMgr; use super::super::metadata_sys::{BucketMetadataSys, get_global_bucket_metadata_sys}; use super::super::{ - get_global_deployment_id, get_global_endpoints_opt, get_global_lock_client, get_global_region, get_global_tier_config_mgr, - global_rustfs_port, + get_global_bucket_monitor, get_global_deployment_id, get_global_endpoints_opt, get_global_lock_client, + get_global_notification_sys, get_global_region, get_global_replication_pool, get_global_tier_config_mgr, global_rustfs_port, }; use super::interfaces::{ - ActionCredentialInterface, BucketMetadataInterface, BufferConfigInterface, DeploymentIdInterface, EndpointsInterface, - IamInterface, KmsInterface, KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, NotifyInterface, - RegionInterface, RuntimePortInterface, ServerConfigInterface, TierConfigInterface, + ActionCredentialInterface, BucketMetadataInterface, BucketMonitorInterface, BufferConfigInterface, DeploymentIdInterface, + EndpointsInterface, IamInterface, KmsInterface, KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, + NotificationSystemInterface, NotifyInterface, RegionInterface, ReplicationPoolInterface, RuntimePortInterface, + ServerConfigInterface, TierConfigInterface, }; use crate::config::{RustFSBufferConfig, get_global_buffer_config}; use async_trait::async_trait; @@ -112,6 +113,16 @@ impl NotifyInterface for NotifyHandle { } } +/// Default notification system handle adapter. +#[derive(Default)] +pub struct NotificationSystemHandle; + +impl NotificationSystemInterface for NotificationSystemHandle { + fn handle(&self) -> Option<&'static super::super::NotificationSys> { + get_global_notification_sys() + } +} + /// Default bucket metadata interface adapter. #[derive(Default)] pub struct BucketMetadataHandle; @@ -122,6 +133,26 @@ impl BucketMetadataInterface for BucketMetadataHandle { } } +/// Default bucket monitor interface adapter. +#[derive(Default)] +pub struct BucketMonitorHandle; + +impl BucketMonitorInterface for BucketMonitorHandle { + fn handle(&self) -> Option> { + get_global_bucket_monitor() + } +} + +/// Default replication pool interface adapter. +#[derive(Default)] +pub struct ReplicationPoolHandle; + +impl ReplicationPoolInterface for ReplicationPoolHandle { + fn handle(&self) -> Option> { + get_global_replication_pool() + } +} + /// Default endpoints interface adapter. #[derive(Default)] pub struct EndpointsHandle; @@ -227,6 +258,10 @@ pub fn default_notify_interface() -> Arc { Arc::new(NotifyHandle) } +pub fn default_notification_system_interface() -> Arc { + Arc::new(NotificationSystemHandle) +} + pub fn default_kms_runtime_interface() -> Arc { Arc::new(KmsRuntimeHandle) } @@ -235,6 +270,14 @@ pub fn default_bucket_metadata_interface() -> Arc { Arc::new(BucketMetadataHandle) } +pub fn default_bucket_monitor_interface() -> Arc { + Arc::new(BucketMonitorHandle) +} + +pub fn default_replication_pool_interface() -> Arc { + Arc::new(ReplicationPoolHandle) +} + pub fn default_endpoints_interface() -> Arc { Arc::new(EndpointsHandle) } diff --git a/rustfs/src/app/context/interfaces.rs b/rustfs/src/app/context/interfaces.rs index 12bcf2d9a..d83ab862c 100644 --- a/rustfs/src/app/context/interfaces.rs +++ b/rustfs/src/app/context/interfaces.rs @@ -15,6 +15,7 @@ use super::super::EndpointServerPools; use super::super::TierConfigMgr; use super::super::metadata_sys::BucketMetadataSys; +use super::super::{BucketBandwidthMonitor, DynReplicationPool, NotificationSys}; use crate::config::RustFSBufferConfig; use async_trait::async_trait; use rustfs_config::server_config::Config; @@ -60,11 +61,26 @@ pub trait NotifyInterface: Send + Sync { async fn clear_bucket_notification_rules(&self, bucket_name: &str) -> Result<(), NotificationError>; } +/// Notification system handle interface for admin peer orchestration. +pub trait NotificationSystemInterface: Send + Sync { + fn handle(&self) -> Option<&'static NotificationSys>; +} + /// Bucket metadata interface for application-layer use-cases. pub trait BucketMetadataInterface: Send + Sync { fn handle(&self) -> Option>>; } +/// Bucket bandwidth monitor interface for admin metric integration. +pub trait BucketMonitorInterface: Send + Sync { + fn handle(&self) -> Option>; +} + +/// Replication pool interface for admin resync integration. +pub trait ReplicationPoolInterface: Send + Sync { + fn handle(&self) -> Option>; +} + /// Endpoints interface for application-layer use-cases. pub trait EndpointsInterface: Send + Sync { fn handle(&self) -> Option; diff --git a/rustfs/src/app/mod.rs b/rustfs/src/app/mod.rs index 5b3740ad6..4a0c1ba0b 100644 --- a/rustfs/src/app/mod.rs +++ b/rustfs/src/app/mod.rs @@ -94,10 +94,12 @@ pub(crate) const MIN_DISK_COMPRESSIBLE_SIZE: usize = ecstore_compression::MIN_DI pub(crate) type DiskError = crate::storage::DiskError; pub(crate) type DynReader = crate::storage::DynReader; +pub(crate) type DynReplicationPool = crate::storage::DynReplicationPool; pub(crate) type ECStore = crate::storage::ECStore; pub(crate) type EndpointServerPools = crate::storage::EndpointServerPools; pub(crate) type HashReader = crate::storage::HashReader; pub(crate) type NotificationSys = crate::storage::NotificationSys; +pub(crate) type BucketBandwidthMonitor = crate::storage::BucketBandwidthMonitor; pub(crate) type ObjectStoreResolver = crate::storage::ObjectStoreResolver; pub(crate) type ObjectInfo = ::ObjectInfo; pub(crate) type ObjectOptions = ::ObjectOptions; @@ -662,6 +664,14 @@ pub(crate) fn get_global_notification_sys() -> Option<&'static NotificationSys> crate::storage::get_global_notification_sys() } +pub(crate) fn get_global_bucket_monitor() -> Option> { + crate::storage::get_global_bucket_monitor() +} + +pub(crate) fn get_global_replication_pool() -> Option> { + crate::storage::get_global_replication_pool() +} + #[cfg(test)] pub(crate) fn boxed_reader(reader: R) -> DynReader where diff --git a/rustfs/src/storage/mod.rs b/rustfs/src/storage/mod.rs index 35f265bdc..891f6d1da 100644 --- a/rustfs/src/storage/mod.rs +++ b/rustfs/src/storage/mod.rs @@ -210,6 +210,7 @@ pub(crate) type BucketMetadata = ecstore_bucket::metadata::BucketMetadata; #[cfg(test)] pub(crate) type BucketMetadataSys = ecstore_bucket::metadata_sys::BucketMetadataSys; pub(crate) type BucketVersioningSys = ecstore_bucket::versioning_sys::BucketVersioningSys; +pub(crate) type BucketBandwidthMonitor = ecstore_bucket::bandwidth::monitor::Monitor; pub(crate) type CheckPartsResp = ecstore_disk::CheckPartsResp; pub(crate) type CollectMetricsOpts = ecstore_metrics::CollectMetricsOpts; pub(crate) type DeleteOptions = ecstore_disk::DeleteOptions; @@ -809,6 +810,10 @@ pub(crate) fn get_global_lock_clients() ecstore_global::get_global_lock_clients() } +pub(crate) fn get_global_bucket_monitor() -> Option> { + ecstore_global::get_global_bucket_monitor() +} + pub(crate) fn get_global_endpoints_opt() -> Option { ecstore_global::get_global_endpoints_opt() }