refactor: centralize ecstore bucket runtime sources (#3810)

This commit is contained in:
Zhengchao An
2026-06-24 09:41:11 +08:00
committed by GitHub
parent 17ae335313
commit 3f624bf06c
7 changed files with 59 additions and 25 deletions
@@ -891,7 +891,7 @@ impl TransitionState {
tokio::spawn(async move {
Self::inc_counter(&state.compensation_running_tasks);
state.record_scanner_transition_state();
let Some(api) = crate::global::resolve_object_store_handle() else {
let Some(api) = runtime_sources::object_store_handle() else {
scheduled.lock().unwrap().remove(&bucket);
Self::add_counter(&state.compensation_running_tasks, -1);
state.record_scanner_transition_state();
@@ -1919,7 +1919,7 @@ pub async fn enqueue_immediate_expiry(oi: &ObjectInfo, src: LcEventSrc) {
let Some(lifecycle) = runtime_sources::bucket_lifecycle_config(&oi.bucket).await else {
return;
};
let Some(api) = crate::global::resolve_object_store_handle() else {
let Some(api) = runtime_sources::object_store_handle() else {
return;
};
+2 -2
View File
@@ -20,7 +20,7 @@ use crate::bucket::utils::deserialize;
use crate::config::com::{read_config, save_config};
use crate::disk::BUCKET_META_PREFIX;
use crate::error::{Error, Result};
use crate::global::resolve_object_store_handle;
use crate::runtime_sources;
use crate::store::ECStore;
use byteorder::{BigEndian, ByteOrder, LittleEndian};
use rustfs_policy::policy::BucketPolicy;
@@ -765,7 +765,7 @@ impl BucketMetadata {
}
pub async fn save(&mut self) -> Result<()> {
let Some(store) = resolve_object_store_handle() else {
let Some(store) = runtime_sources::object_store_handle() else {
return Err(Error::other("errServerNotInitialized"));
};
+10 -11
View File
@@ -19,7 +19,7 @@ use crate::bucket::bucket_target_sys::BucketTargetSys;
use crate::bucket::metadata::{BUCKET_LIFECYCLE_CONFIG, load_bucket_metadata_parse};
use crate::bucket::utils::{deserialize, is_meta_bucketname};
use crate::error::{Error, Result, is_err_bucket_not_found};
use crate::global::{GLOBAL_Endpoints, is_dist_erasure, is_erasure, resolve_object_store_handle};
use crate::runtime_sources;
use crate::store::ECStore;
use futures::future::join_all;
use lazy_static::lazy_static;
@@ -263,13 +263,9 @@ impl BucketMetadataSys {
let _ = self.init_internal(buckets).await;
}
async fn init_internal(&self, buckets: Vec<String>) -> Result<()> {
let count = {
if let Some(endpoints) = GLOBAL_Endpoints.get() {
endpoints.es_count() * 10
} else {
return Err(Error::other("GLOBAL_Endpoints not init"));
}
};
let count = runtime_sources::endpoint_erasure_set_count()
.map(|count| count * 10)
.ok_or_else(|| Error::other("GLOBAL_Endpoints not init"))?;
let mut failed_buckets: HashSet<String> = HashSet::new();
let mut buckets = buckets.as_slice();
@@ -288,7 +284,7 @@ impl BucketMetadataSys {
let mut initialized = self.initialized.write().await;
*initialized = true;
if is_dist_erasure().await {
if runtime_sources::setup_is_dist_erasure().await {
// TODO: refresh_buckets_metadata_loop
}
@@ -406,7 +402,7 @@ impl BucketMetadataSys {
}
async fn update_and_parse(&mut self, bucket: &str, config_file: &str, data: Vec<u8>, parse: bool) -> Result<OffsetDateTime> {
let Some(store) = resolve_object_store_handle() else {
let Some(store) = runtime_sources::object_store_handle() else {
return Err(Error::other("errServerNotInitialized"));
};
@@ -417,7 +413,10 @@ impl BucketMetadataSys {
let mut bm = match load_bucket_metadata_parse(store, bucket, parse).await {
Ok(res) => res,
Err(err) => {
if !is_erasure().await && !is_dist_erasure().await && is_err_bucket_not_found(&err) {
if !runtime_sources::setup_is_erasure().await
&& !runtime_sources::setup_is_dist_erasure().await
&& is_err_bucket_not_found(&err)
{
BucketMetadata::new(bucket)
} else {
error!("load bucket metadata failed: {}", err);
@@ -28,7 +28,6 @@ use crate::config::com::save_config;
use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET};
use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_found};
use crate::event_notification::{EventArgs, send_event};
use crate::global::resolve_object_store_handle;
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
use crate::runtime_sources;
use crate::set_disk::get_lock_acquire_timeout;
@@ -1758,7 +1757,7 @@ impl ObjectInfoExt for ObjectInfo {
}
pub async fn must_replicate(bucket: &str, object: &str, mopts: MustReplicateOptions) -> ReplicateDecision {
if resolve_object_store_handle().is_none() {
if runtime_sources::object_store_handle().is_none() {
return ReplicateDecision::default();
}
+14 -2
View File
@@ -28,8 +28,8 @@ use crate::{
GLOBAL_BOOT_TIME, GLOBAL_EventNotifier, GLOBAL_IsErasureSD, GLOBAL_LOCAL_DISK_ID_MAP, GLOBAL_LOCAL_DISK_MAP,
GLOBAL_LOCAL_DISK_SET_DRIVES, GLOBAL_LifecycleSys, GLOBAL_LocalNodeName, GLOBAL_RootDiskThreshold, GLOBAL_TierConfigMgr,
TypeLocalDiskSetDrives, get_global_bucket_monitor, get_global_deployment_id, get_global_endpoints,
get_global_endpoints_opt, init_global_bucket_monitor, is_first_cluster_node_local, resolve_object_store_handle,
set_global_deployment_id,
get_global_endpoints_opt, init_global_bucket_monitor, is_dist_erasure, is_erasure, is_first_cluster_node_local,
resolve_object_store_handle, set_global_deployment_id,
},
notification_sys::{NotificationSys, get_global_notification_sys},
store::ECStore,
@@ -59,6 +59,10 @@ pub(crate) fn endpoint_pools() -> Option<EndpointServerPools> {
get_global_endpoints_opt()
}
pub(crate) fn endpoint_erasure_set_count() -> Option<usize> {
endpoint_pools().map(|endpoints| endpoints.es_count())
}
pub(crate) fn endpoint_pool_is_local(pool_index: usize) -> bool {
get_global_endpoints()
.as_ref()
@@ -70,6 +74,14 @@ pub(crate) async fn first_cluster_node_is_local() -> bool {
is_first_cluster_node_local().await
}
pub(crate) async fn setup_is_erasure() -> bool {
is_erasure().await
}
pub(crate) async fn setup_is_dist_erasure() -> bool {
is_dist_erasure().await
}
pub(crate) async fn local_node_name() -> String {
GLOBAL_LOCAL_NODE_NAME.read().await.clone()
}
+5 -1
View File
@@ -77,7 +77,11 @@ fn check() -> Result<(), OpaConfigError> {
Ok(())
}
async fn validate(config: &Args) -> Result<(), OpaConfigError> {
let client = reqwest::Client::new();
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(5))
.connect_timeout(Duration::from_secs(1))
.build()
.map_err(OpaConfigError::Connection)?;
match client.post(&config.url).send().await {
Ok(resp) => {