refactor: route readiness through app context (#3770)

This commit is contained in:
Zhengchao An
2026-06-23 08:28:19 +08:00
committed by GitHub
parent 7499dd085d
commit 087f794901
3 changed files with 61 additions and 18 deletions
+18 -3
View File
@@ -41,6 +41,13 @@ pub fn resolve_kms_runtime_service_manager() -> Option<Arc<KmsServiceManager>> {
resolve_kms_runtime_service_manager_with(get_global_app_context(), || default_kms_runtime_interface().service_manager())
}
/// Resolve IAM readiness using AppContext-first precedence.
pub fn resolve_iam_ready() -> bool {
resolve_iam_ready_with(get_global_app_context(), || {
rustfs_iam::get_global_iam_sys().is_some_and(|sys| sys.is_ready())
})
}
/// Resolve bucket metadata handle using AppContext-first precedence.
pub fn resolve_bucket_metadata_handle() -> Option<Arc<RwLock<BucketMetadataSys>>> {
resolve_bucket_metadata_handle_with(get_global_app_context(), || default_bucket_metadata_interface().handle())
@@ -93,6 +100,10 @@ fn resolve_kms_runtime_service_manager_with(
.or_else(fallback)
}
fn resolve_iam_ready_with(context: Option<Arc<AppContext>>, fallback: impl FnOnce() -> bool) -> bool {
context.map_or_else(fallback, |context| context.iam().is_ready())
}
fn resolve_bucket_metadata_handle_with(
context: Option<Arc<AppContext>>,
fallback: impl FnOnce() -> Option<Arc<RwLock<BucketMetadataSys>>>,
@@ -154,7 +165,9 @@ mod tests {
use tempfile::TempDir;
use tokio_util::sync::CancellationToken;
struct TestIamInterface;
struct TestIamInterface {
ready: bool,
}
impl IamInterface for TestIamInterface {
fn handle(&self) -> Arc<IamSys<ObjectStore>> {
@@ -162,7 +175,7 @@ mod tests {
}
fn is_ready(&self) -> bool {
true
self.ready
}
}
@@ -299,7 +312,7 @@ mod tests {
let context = Arc::new(AppContext::with_test_interfaces(
object_store.clone(),
AppContextTestInterfaces {
iam: Arc::new(TestIamInterface),
iam: Arc::new(TestIamInterface { ready: true }),
kms: Arc::new(TestKmsInterface {
kms: context_kms.clone(),
}),
@@ -329,6 +342,7 @@ mod tests {
.expect("context KMS runtime"),
&context_kms
));
assert!(resolve_iam_ready_with(Some(context.clone()), || false));
assert!(Arc::ptr_eq(
&resolve_bucket_metadata_handle_with(Some(context.clone()), || None).expect("context bucket metadata"),
&bucket_metadata
@@ -361,6 +375,7 @@ mod tests {
&resolve_kms_runtime_service_manager_with(None, || Some(fallback_kms.clone())).expect("fallback KMS runtime"),
&fallback_kms
));
assert!(!resolve_iam_ready_with(None, || false));
assert!(Arc::ptr_eq(
&resolve_bucket_metadata_handle_with(None, || Some(bucket_metadata.clone())).expect("fallback bucket metadata"),
&bucket_metadata
+5 -8
View File
@@ -12,12 +12,10 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::app::context::{resolve_endpoints_handle, resolve_iam_ready};
use crate::server::{ServiceState, ServiceStateManager};
use crate::server::{has_path_prefix, is_table_catalog_path};
use crate::storage::{
Endpoint, EndpointServerPools, get_global_endpoints_opt, get_global_lock_clients, is_dist_erasure,
resolve_object_store_handle,
};
use crate::storage::{Endpoint, EndpointServerPools, get_global_lock_clients, is_dist_erasure, resolve_object_store_handle};
#[cfg(test)]
use crate::storage::{Endpoints, PoolEndpoints};
use bytes::Bytes;
@@ -27,7 +25,6 @@ use http_body_util::{BodyExt, Full};
use hyper::body::Incoming;
use metrics::{counter, gauge};
use rustfs_common::GlobalReadiness;
use rustfs_iam::get_global_iam_sys;
use rustfs_madmin::{Disk, StorageInfo};
use rustfs_storage_api::StorageAdminApi;
use std::future::Future;
@@ -441,7 +438,7 @@ pub async fn collect_dependency_readiness() -> DependencyReadiness {
}
pub async fn collect_dependency_readiness_report() -> DependencyReadinessReport {
let iam_ready_raw = get_global_iam_sys().is_some_and(|sys| sys.is_ready());
let iam_ready_raw = resolve_iam_ready();
let storage_ready = if let Some(cached) = load_cached_storage_readiness().await {
cached
} else {
@@ -475,7 +472,7 @@ async fn collect_lock_quorum_status() -> LockQuorumStatus {
}
async fn collect_dependency_readiness_uncached() -> DependencyReadiness {
let iam_ready_raw = get_global_iam_sys().is_some_and(|sys| sys.is_ready());
let iam_ready_raw = resolve_iam_ready();
let storage_ready = collect_storage_readiness_uncached().await;
let lock_quorum_status = collect_lock_quorum_status_uncached().await;
@@ -584,7 +581,7 @@ async fn collect_lock_quorum_status_uncached() -> LockQuorumStatus {
};
}
let Some(pool_endpoints) = get_global_endpoints_opt() else {
let Some(pool_endpoints) = resolve_endpoints_handle() else {
return LockQuorumStatus::default();
};
let Some(lock_clients) = get_global_lock_clients() else {