From 73efd4f4938fb74f8e063b2720fa09b8971967c5 Mon Sep 17 00:00:00 2001 From: weisd Date: Thu, 28 Nov 2024 14:47:01 +0800 Subject: [PATCH] use Arc --- api/admin/src/handlers/list_pools.rs | 7 +- ecstore/src/admin_server_info.rs | 20 +- ecstore/src/bucket/metadata.rs | 21 +- ecstore/src/bucket/metadata_sys.rs | 145 ++++++------ ecstore/src/bucket/policy_sys.rs | 2 +- ecstore/src/bucket/versioning_sys.rs | 2 +- ecstore/src/config/common.rs | 31 +-- ecstore/src/config/mod.rs | 6 +- ecstore/src/disk/endpoint.rs | 1 - ecstore/src/disk/local.rs | 16 +- ecstore/src/global.rs | 21 +- ecstore/src/heal/background_heal_ops.rs | 30 +-- ecstore/src/heal/data_scanner.rs | 29 ++- ecstore/src/heal/data_usage.rs | 18 +- ecstore/src/heal/data_usage_cache.rs | 11 +- ecstore/src/heal/heal_commands.rs | 8 +- ecstore/src/heal/heal_ops.rs | 21 +- ecstore/src/pools.rs | 31 +-- ecstore/src/store.rs | 96 ++++---- madmin/src/heal_command.rs | 15 +- rustfs/src/admin/handlers.rs | 35 +-- rustfs/src/grpc.rs | 204 ++++++----------- rustfs/src/main.rs | 4 +- rustfs/src/storage/ecfs.rs | 287 ++++++++---------------- 24 files changed, 411 insertions(+), 650 deletions(-) diff --git a/api/admin/src/handlers/list_pools.rs b/api/admin/src/handlers/list_pools.rs index 277183d55..b80dd4654 100644 --- a/api/admin/src/handlers/list_pools.rs +++ b/api/admin/src/handlers/list_pools.rs @@ -2,6 +2,7 @@ use crate::error::ErrorCode; use crate::Result as LocalResult; use axum::Json; +use ecstore::new_object_layer_fn; use serde::Serialize; use time::OffsetDateTime; @@ -47,14 +48,12 @@ pub async fn handler() -> LocalResult>> { // // todo 实用oncelock作为全局变量 - let layer = ecstore::new_object_layer_fn(); - let lock = layer.read().await; - let pools = lock.as_ref().ok_or(ErrorCode::ErrNotImplemented)?; + let Some(store) = new_object_layer_fn() else { return Err(ErrorCode::ErrNotImplemented) }; // todo, 调用pool.status()接口获取每个池的数据 // let mut result = Vec::new(); - for (idx, _pool) in pools.pools.iter().enumerate() { + for (idx, _pool) in store.pools.iter().enumerate() { // 这里mock一下数据 result.push(PoolStatus { id: idx as _, diff --git a/ecstore/src/admin_server_info.rs b/ecstore/src/admin_server_info.rs index 5d4d9a8f9..8c140b6cc 100644 --- a/ecstore/src/admin_server_info.rs +++ b/ecstore/src/admin_server_info.rs @@ -155,19 +155,13 @@ pub async fn get_local_server_property() -> ServerProperties { // sensitive.insert(ENV_ROOT_USER.to_string()); // sensitive.insert(ENV_ROOT_PASSWORD.to_string()); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - match lock.as_ref() { - Some(store) => { - let storage_info = store.local_storage_info().await; - props.state = ITEM_ONLINE.to_string(); - props.disks = storage_info.disks; - } - None => { - props.state = ITEM_INITIALIZING.to_string(); - // todo: get_offline_disks - // props.disks = - } + if let Some(store) = new_object_layer_fn() { + let storage_info = store.local_storage_info().await; + props.state = ITEM_ONLINE.to_string(); + props.disks = storage_info.disks; + } else { + props.state = ITEM_INITIALIZING.to_string(); }; + props } diff --git a/ecstore/src/bucket/metadata.rs b/ecstore/src/bucket/metadata.rs index a296a973f..5bcc353f6 100644 --- a/ecstore/src/bucket/metadata.rs +++ b/ecstore/src/bucket/metadata.rs @@ -12,12 +12,13 @@ use s3s::dto::{ use serde::Serializer; use serde::{Deserialize, Serialize}; use std::collections::HashMap; +use std::sync::Arc; use time::OffsetDateTime; use tracing::{error, warn}; -use crate::config; use crate::config::common::{read_config, save_config}; use crate::error::{Error, Result}; +use crate::{config, new_object_layer_fn}; use crate::disk::BUCKET_META_PREFIX; use crate::store::ECStore; @@ -290,8 +291,10 @@ impl BucketMetadata { self.created = created } - pub async fn save(&mut self, api: &ECStore) -> Result<()> { - self.parse_all_configs(api)?; + pub async fn save(&mut self) -> Result<()> { + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; + + self.parse_all_configs(store.clone())?; let mut buf: Vec = vec![0; 4]; @@ -303,12 +306,12 @@ impl BucketMetadata { buf.extend_from_slice(&data); - save_config(api, self.save_file_path().as_str(), &buf).await?; + save_config(store, self.save_file_path().as_str(), &buf).await?; Ok(()) } - fn parse_all_configs(&mut self, _api: &ECStore) -> Result<()> { + fn parse_all_configs(&mut self, _api: Arc) -> Result<()> { if !self.policy_config_json.is_empty() { self.policy_config = Some(serde_json::from_slice(&self.policy_config_json)?); } @@ -347,12 +350,12 @@ impl BucketMetadata { } } -pub async fn load_bucket_metadata(api: &ECStore, bucket: &str) -> Result { +pub async fn load_bucket_metadata(api: Arc, bucket: &str) -> Result { load_bucket_metadata_parse(api, bucket, true).await } -pub async fn load_bucket_metadata_parse(api: &ECStore, bucket: &str, parse: bool) -> Result { - let mut bm = match read_bucket_metadata(api, bucket).await { +pub async fn load_bucket_metadata_parse(api: Arc, bucket: &str, parse: bool) -> Result { + let mut bm = match read_bucket_metadata(api.clone(), bucket).await { Ok(res) => res, Err(err) => { warn!("load_bucket_metadata_parse err {:?}", &err); @@ -375,7 +378,7 @@ pub async fn load_bucket_metadata_parse(api: &ECStore, bucket: &str, parse: bool Ok(bm) } -async fn read_bucket_metadata(api: &ECStore, bucket: &str) -> Result { +async fn read_bucket_metadata(api: Arc, bucket: &str) -> Result { if bucket.is_empty() { error!("bucket name empty"); return Err(Error::msg("invalid argument")); diff --git a/ecstore/src/bucket/metadata_sys.rs b/ecstore/src/bucket/metadata_sys.rs index fb45c269b..162862124 100644 --- a/ecstore/src/bucket/metadata_sys.rs +++ b/ecstore/src/bucket/metadata_sys.rs @@ -1,4 +1,5 @@ use std::collections::HashSet; +use std::sync::OnceLock; use std::{collections::HashMap, sync::Arc}; use crate::bucket::error::BucketMetadataError; @@ -28,123 +29,129 @@ use super::target::BucketTargets; use lazy_static::lazy_static; lazy_static! { - pub static ref GLOBAL_BucketMetadataSys: Arc> = Arc::new(RwLock::new(BucketMetadataSys::new())); + pub static ref GLOBAL_BucketMetadataSys: OnceLock>> = OnceLock::new(); } -pub async fn init_bucket_metadata_sys(api: ECStore, buckets: Vec) { - let mut sys = GLOBAL_BucketMetadataSys.write().await; - sys.init(api, buckets).await +pub async fn init_bucket_metadata_sys(api: Arc, buckets: Vec) { + let mut sys = BucketMetadataSys::new(api); + sys.init(buckets).await; + + let sys = Arc::new(RwLock::new(sys)); + + GLOBAL_BucketMetadataSys.set(sys).unwrap(); } -pub(super) async fn get_bucket_metadata_sys() -> Arc> { - GLOBAL_BucketMetadataSys.clone() +// panic if not init +pub(super) fn get_bucket_metadata_sys() -> Arc> { + GLOBAL_BucketMetadataSys.get().unwrap().clone() } -pub(crate) async fn set_bucket_metadata(bucket: String, bm: BucketMetadata) { - let sys = GLOBAL_BucketMetadataSys.write().await; - sys.set(bucket, Arc::new(bm)).await +pub async fn set_bucket_metadata(bucket: String, bm: BucketMetadata) { + let sys = get_bucket_metadata_sys(); + let lock = sys.write().await; + lock.set(bucket, Arc::new(bm)).await; } pub(crate) async fn get(bucket: &str) -> Result> { - let sys = GLOBAL_BucketMetadataSys.write().await; - sys.get(bucket).await + let sys = get_bucket_metadata_sys(); + let lock = sys.read().await; + lock.get(bucket).await } pub async fn update(bucket: &str, config_file: &str, data: Vec) -> Result { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let mut bucket_meta_sys = bucket_meta_sys_lock.write().await; bucket_meta_sys.update(bucket, config_file, data).await } pub async fn delete(bucket: &str, config_file: &str) -> Result { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let mut bucket_meta_sys = bucket_meta_sys_lock.write().await; bucket_meta_sys.delete(bucket, config_file).await } pub async fn get_tagging_config(bucket: &str) -> Result<(Tagging, OffsetDateTime)> { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.get_tagging_config(bucket).await } pub async fn get_lifecycle_config(bucket: &str) -> Result<(BucketLifecycleConfiguration, OffsetDateTime)> { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.get_lifecycle_config(bucket).await } pub async fn get_sse_config(bucket: &str) -> Result<(ServerSideEncryptionConfiguration, OffsetDateTime)> { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.get_sse_config(bucket).await } pub async fn get_object_lock_config(bucket: &str) -> Result<(ObjectLockConfiguration, OffsetDateTime)> { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.get_object_lock_config(bucket).await } pub async fn get_replication_config(bucket: &str) -> Result<(ReplicationConfiguration, OffsetDateTime)> { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.get_replication_config(bucket).await } pub async fn get_notification_config(bucket: &str) -> Result> { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.get_notification_config(bucket).await } pub async fn get_versioning_config(bucket: &str) -> Result<(VersioningConfiguration, OffsetDateTime)> { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.get_versioning_config(bucket).await } pub async fn get_config_from_disk(bucket: &str) -> Result { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.get_config_from_disk(bucket).await } pub async fn created_at(bucket: &str) -> Result { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.read().await; bucket_meta_sys.created_at(bucket).await } -#[derive(Debug, Default)] +#[derive(Debug)] pub struct BucketMetadataSys { metadata_map: RwLock>>, - api: Option, + api: Arc, initialized: RwLock, } impl BucketMetadataSys { - pub fn new() -> Self { - Self::default() + pub fn new(api: Arc) -> Self { + Self { + metadata_map: RwLock::new(HashMap::new()), + api, + initialized: RwLock::new(false), + } } - pub async fn init(&mut self, api: ECStore, buckets: Vec) { - // if api.is_none() { - // return Err(Error::msg("errServerNotInitialized")); - // } - self.api = Some(api); - + pub async fn init(&mut self, buckets: Vec) { let _ = self.init_internal(buckets).await; } async fn init_internal(&self, buckets: Vec) -> Result<()> { @@ -184,7 +191,7 @@ impl BucketMetadataSys { let mut futures = Vec::new(); for bucket in buckets.iter() { - futures.push(load_bucket_metadata(self.api.as_ref().unwrap(), bucket.as_str())); + futures.push(load_bucket_metadata(self.api.clone(), bucket.as_str())); } let results = join_all(futures).await; @@ -270,12 +277,7 @@ impl BucketMetadataSys { } async fn update_and_parse(&mut self, bucket: &str, config_file: &str, data: Vec, parse: bool) -> Result { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::msg("errServerNotInitialized")), - }; + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; if is_meta_bucketname(bucket) { return Err(Error::msg("errInvalidArgument")); @@ -300,20 +302,13 @@ impl BucketMetadataSys { } async fn save(&self, bm: BucketMetadata) -> Result<()> { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::msg("errServerNotInitialized")), - }; - if is_meta_bucketname(&bm.name) { return Err(Error::msg("errInvalidArgument")); } let mut bm = bm; - bm.save(store).await?; + bm.save().await?; self.set(bm.name.clone(), Arc::new(bm)).await; @@ -321,51 +316,39 @@ impl BucketMetadataSys { } pub async fn get_config_from_disk(&self, bucket: &str) -> Result { - if self.api.as_ref().is_none() { - return Err(Error::msg("errBucketMetadataNotInitialized")); - } - if is_meta_bucketname(bucket) { return Err(Error::msg("errInvalidArgument")); } - if let Some(api) = self.api.as_ref() { - load_bucket_metadata(api, bucket).await - } else { - Err(Error::msg("errBucketMetadataNotInitialized")) - } + load_bucket_metadata(self.api.clone(), bucket).await } pub async fn get_config(&self, bucket: &str) -> Result<(Arc, bool)> { - if let Some(api) = self.api.as_ref() { - let has_bm = { - let map = self.metadata_map.read().await; - map.get(&bucket.to_string()).cloned() + let has_bm = { + let map = self.metadata_map.read().await; + map.get(&bucket.to_string()).cloned() + }; + + if let Some(bm) = has_bm { + Ok((bm, false)) + } else { + let bm = match load_bucket_metadata(self.api.clone(), bucket).await { + Ok(res) => res, + Err(err) => { + if *self.initialized.read().await { + return Err(Error::msg("errBucketMetadataNotInitialized")); + } else { + return Err(err); + } + } }; - if let Some(bm) = has_bm { - Ok((bm, false)) - } else { - let bm = match load_bucket_metadata(api, bucket).await { - Ok(res) => res, - Err(err) => { - if *self.initialized.read().await { - return Err(Error::msg("errBucketMetadataNotInitialized")); - } else { - return Err(err); - } - } - }; + let mut map = self.metadata_map.write().await; - let mut map = self.metadata_map.write().await; + let bm = Arc::new(bm); + map.insert(bucket.to_string(), bm.clone()); - let bm = Arc::new(bm); - map.insert(bucket.to_string(), bm.clone()); - - Ok((bm, true)) - } - } else { - Err(Error::msg("errBucketMetadataNotInitialized")) + Ok((bm, true)) } } diff --git a/ecstore/src/bucket/policy_sys.rs b/ecstore/src/bucket/policy_sys.rs index 9a2501190..b72cb333f 100644 --- a/ecstore/src/bucket/policy_sys.rs +++ b/ecstore/src/bucket/policy_sys.rs @@ -22,7 +22,7 @@ impl PolicySys { args.is_owner } pub async fn get(bucket: &str) -> Result { - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.write().await; let (cfg, _) = bucket_meta_sys.get_bucket_policy(bucket).await?; diff --git a/ecstore/src/bucket/versioning_sys.rs b/ecstore/src/bucket/versioning_sys.rs index b14fc910d..b464a1606 100644 --- a/ecstore/src/bucket/versioning_sys.rs +++ b/ecstore/src/bucket/versioning_sys.rs @@ -52,7 +52,7 @@ impl BucketVersioningSys { return Ok(VersioningConfiguration::default()); } - let bucket_meta_sys_lock = get_bucket_metadata_sys().await; + let bucket_meta_sys_lock = get_bucket_metadata_sys(); let bucket_meta_sys = bucket_meta_sys_lock.write().await; let (cfg, _) = bucket_meta_sys.get_versioning_config(bucket).await?; diff --git a/ecstore/src/config/common.rs b/ecstore/src/config/common.rs index eb9d65a4b..7ce6ced42 100644 --- a/ecstore/src/config/common.rs +++ b/ecstore/src/config/common.rs @@ -1,4 +1,5 @@ use std::collections::HashSet; +use std::sync::Arc; use super::error::ConfigError; use super::{storageclass, Config, GLOBAL_StorageClass, KVS}; @@ -29,13 +30,13 @@ lazy_static! { h }; } -pub async fn read_config(api: &ECStore, file: &str) -> Result> { +pub async fn read_config(api: Arc, file: &str) -> Result> { let (data, _obj) = read_config_with_metadata(api, file, &ObjectOptions::default()).await?; Ok(data) } -async fn read_config_with_metadata(api: &ECStore, file: &str, opts: &ObjectOptions) -> Result<(Vec, ObjectInfo)> { +async fn read_config_with_metadata(api: Arc, file: &str, opts: &ObjectOptions) -> Result<(Vec, ObjectInfo)> { let range = HTTPRangeSpec::nil(); let h = HeaderMap::new(); let mut rd = api @@ -58,7 +59,7 @@ async fn read_config_with_metadata(api: &ECStore, file: &str, opts: &ObjectOptio Ok((data, rd.object_info)) } -pub async fn save_config(api: &ECStore, file: &str, data: &[u8]) -> Result<()> { +pub async fn save_config(api: Arc, file: &str, data: &[u8]) -> Result<()> { save_config_with_opts( api, file, @@ -71,7 +72,7 @@ pub async fn save_config(api: &ECStore, file: &str, data: &[u8]) -> Result<()> { .await } -async fn save_config_with_opts(api: &ECStore, file: &str, data: &[u8], opts: &ObjectOptions) -> Result<()> { +async fn save_config_with_opts(api: Arc, file: &str, data: &[u8], opts: &ObjectOptions) -> Result<()> { let _ = api .put_object( RUSTFS_META_BUCKET, @@ -87,17 +88,17 @@ fn new_server_config() -> Config { Config::new() } -async fn new_and_save_server_config(api: &ECStore) -> Result { +async fn new_and_save_server_config(api: Arc) -> Result { let mut cfg = new_server_config(); - lookup_configs(&mut cfg, api).await; + lookup_configs(&mut cfg, api.clone()).await; save_server_config(api, &cfg).await?; Ok(cfg) } -pub async fn read_config_without_migrate(api: &ECStore) -> Result { +pub async fn read_config_without_migrate(api: Arc) -> Result { let config_file = format!("{}{}{}", CONFIG_PREFIX, SLASH_SEPARATOR, CONFIG_FILE); - let data = match read_config(api, config_file.as_str()).await { + let data = match read_config(api.clone(), config_file.as_str()).await { Ok(res) => res, Err(err) => { if is_not_found(&err) { @@ -115,11 +116,11 @@ pub async fn read_config_without_migrate(api: &ECStore) -> Result { read_server_config(api, data.as_slice()).await } -async fn read_server_config(api: &ECStore, data: &[u8]) -> Result { +async fn read_server_config(api: Arc, data: &[u8]) -> Result { let cfg = { if data.is_empty() { let config_file = format!("{}{}{}", CONFIG_PREFIX, SLASH_SEPARATOR, CONFIG_FILE); - let cfg_data = match read_config(api, config_file.as_str()).await { + let cfg_data = match read_config(api.clone(), config_file.as_str()).await { Ok(res) => res, Err(err) => { if is_not_found(&err) { @@ -144,7 +145,7 @@ async fn read_server_config(api: &ECStore, data: &[u8]) -> Result { Ok(cfg.merge()) } -async fn save_server_config(api: &ECStore, cfg: &Config) -> Result<()> { +async fn save_server_config(api: Arc, cfg: &Config) -> Result<()> { let data = cfg.marshal()?; let config_file = format!("{}{}{}", CONFIG_PREFIX, SLASH_SEPARATOR, CONFIG_FILE); @@ -152,22 +153,22 @@ async fn save_server_config(api: &ECStore, cfg: &Config) -> Result<()> { save_config(api, &config_file, data.as_slice()).await } -pub async fn lookup_configs(cfg: &mut Config, api: &ECStore) { +pub async fn lookup_configs(cfg: &mut Config, api: Arc) { // TODO: from etcd if let Err(err) = apply_dynamic_config(cfg, api).await { error!("apply_dynamic_config err {:?}", &err); } } -async fn apply_dynamic_config(cfg: &mut Config, api: &ECStore) -> Result<()> { +async fn apply_dynamic_config(cfg: &mut Config, api: Arc) -> Result<()> { for key in SubSystemsDynamic.iter() { - apply_dynamic_config_for_sub_sys(cfg, api, key).await?; + apply_dynamic_config_for_sub_sys(cfg, api.clone(), key).await?; } Ok(()) } -async fn apply_dynamic_config_for_sub_sys(cfg: &mut Config, api: &ECStore, subsys: &str) -> Result<()> { +async fn apply_dynamic_config_for_sub_sys(cfg: &mut Config, api: Arc, subsys: &str) -> Result<()> { let set_drive_counts = api.set_drive_counts(); if subsys == STORAGE_CLASS_SUB_SYS { let kvs = match cfg.get_value(STORAGE_CLASS_SUB_SYS, DEFAULT_KV_KEY) { diff --git a/ecstore/src/config/mod.rs b/ecstore/src/config/mod.rs index f241d0f12..c4bae997a 100644 --- a/ecstore/src/config/mod.rs +++ b/ecstore/src/config/mod.rs @@ -10,7 +10,7 @@ use common::{lookup_configs, read_config_without_migrate, STORAGE_CLASS_SUB_SYS} use lazy_static::lazy_static; use serde::{Deserialize, Serialize}; use std::collections::HashMap; -use std::sync::OnceLock; +use std::sync::{Arc, OnceLock}; lazy_static! { pub static ref GLOBAL_StorageClass: OnceLock = OnceLock::new(); @@ -38,8 +38,8 @@ impl ConfigSys { pub fn new() -> Self { Self {} } - pub async fn init(&self, api: &ECStore) -> Result<()> { - let mut cfg = read_config_without_migrate(api).await?; + pub async fn init(&self, api: Arc) -> Result<()> { + let mut cfg = read_config_without_migrate(api.clone().clone()).await?; lookup_configs(&mut cfg, api).await; diff --git a/ecstore/src/disk/endpoint.rs b/ecstore/src/disk/endpoint.rs index 16edab64e..254a1213b 100644 --- a/ecstore/src/disk/endpoint.rs +++ b/ecstore/src/disk/endpoint.rs @@ -2,7 +2,6 @@ use crate::error::{Error, Result}; use crate::utils::net; use path_absolutize::Absolutize; use path_clean::PathClean; -use std::fs; use std::{fmt::Display, path::Path}; use url::{ParseError, Url}; diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index ac14fd5ff..46dee48a6 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -9,7 +9,7 @@ use super::{ UpdateMetadataOpts, VolumeInfo, WalkDirOptions, BUCKET_META_PREFIX, RUSTFS_META_BUCKET, STORAGE_FORMAT_FILE_BACKUP, }; use crate::bitrot::bitrot_verify; -use crate::bucket::metadata_sys::GLOBAL_BucketMetadataSys; +use crate::bucket::metadata_sys::{self}; use crate::cache_value::cache::{Cache, Opts, UpdateFn}; use crate::disk::error::{ convert_access_error, is_sys_err_handle_invalid, is_sys_err_invalid_arg, is_sys_err_is_dir, is_sys_err_not_dir, @@ -1967,23 +1967,13 @@ impl DiskAPI for LocalDisk { defer!(|| { self.scanning.fetch_sub(1, Ordering::SeqCst) }); // Check if the current bucket has replication configuration - if let Ok((rcfg, _)) = GLOBAL_BucketMetadataSys - .read() - .await - .get_replication_config(&cache.info.name) - .await - { + if let Ok((rcfg, _)) = metadata_sys::get_replication_config(&cache.info.name).await { if has_active_rules(&rcfg, "", true) { // TODO: globalBucketTargetSys } } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::msg("errServerNotInitialized")), - }; + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; let loc = self.get_disk_location(); let disks = store.get_disks(loc.pool_idx.unwrap(), loc.disk_idx.unwrap()).await?; let disk = Arc::new(LocalDisk::new(&self.endpoint(), false).await?); diff --git a/ecstore/src/global.rs b/ecstore/src/global.rs index 1aaad8a6e..fdb7f8444 100644 --- a/ecstore/src/global.rs +++ b/ecstore/src/global.rs @@ -9,7 +9,6 @@ use uuid::Uuid; use crate::{ disk::DiskStore, endpoints::{EndpointServerPools, PoolEndpoints, SetupType}, - error::{Error, Result}, heal::{background_heal_ops::HealRoutine, heal_ops::AllHealState}, store::ECStore, }; @@ -20,7 +19,7 @@ pub const DISK_FILL_FRACTION: f64 = 0.99; pub const DISK_RESERVE_FRACTION: f64 = 0.15; lazy_static! { - pub static ref GLOBAL_OBJECT_API: Arc>> = Arc::new(RwLock::new(None)); + pub static ref GLOBAL_OBJECT_API: OnceLock> = OnceLock::new(); pub static ref GLOBAL_LOCAL_DISK: Arc>>> = Arc::new(RwLock::new(Vec::new())); pub static ref GLOBAL_IsErasure: RwLock = RwLock::new(false); pub static ref GLOBAL_IsDistErasure: RwLock = RwLock::new(false); @@ -44,20 +43,22 @@ pub async fn get_global_deployment_id() -> Uuid { *id_ptr } -pub fn set_global_endpoints(eps: Vec) -> Result<()> { +pub fn set_global_endpoints(eps: Vec) { GLOBAL_Endpoints .set(EndpointServerPools::from(eps)) - .map_err(|_| Error::msg("GLOBAL_Endpoints set faild"))?; - Ok(()) + .expect("GLOBAL_Endpoints set faild") } -pub fn new_object_layer_fn() -> Arc>> { - GLOBAL_OBJECT_API.clone() +pub fn new_object_layer_fn() -> Option> { + if let Some(ec) = GLOBAL_OBJECT_API.get() { + Some(ec.clone()) + } else { + None + } } -pub async fn set_object_layer(o: ECStore) { - let mut global_object_api = GLOBAL_OBJECT_API.write().await; - *global_object_api = Some(o); +pub async fn set_object_layer(o: Arc) { + GLOBAL_OBJECT_API.set(o).expect("set_object_layer fail ") } pub async fn is_dist_erasure() -> bool { diff --git a/ecstore/src/heal/background_heal_ops.rs b/ecstore/src/heal/background_heal_ops.rs index 0c64482a1..f90bab544 100644 --- a/ecstore/src/heal/background_heal_ops.rs +++ b/ecstore/src/heal/background_heal_ops.rs @@ -1,4 +1,3 @@ -use s3s::{S3Error, S3ErrorCode}; use std::{cmp::Ordering, env, path::PathBuf, sync::Arc, time::Duration}; use tokio::{ sync::{ @@ -100,9 +99,8 @@ async fn monitor_local_disks_and_heal() { interval.reset(); continue; } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock.as_ref().expect("errServerNotInitialized"); + + let store = new_object_layer_fn().expect("errServerNotInitialized"); if let (_, Some(err)) = store.heal_format(false).await.expect("heal format failed") { if let Some(DiskError::NoHealRequired) = err.downcast_ref() { } else { @@ -187,12 +185,8 @@ async fn heal_fresh_disk(endpoint: &Endpoint) -> Result<()> { endpoint.to_string() ); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::msg("errServerNotInitialized")), - }; + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; + let mut buckets = store.list_bucket(&BucketOptions::default()).await?; buckets.push(BucketInfo { name: path_join(&[PathBuf::from(RUSTFS_META_BUCKET), PathBuf::from(RUSTFS_CONFIG_PREFIX)]) @@ -274,12 +268,7 @@ async fn heal_fresh_disk(endpoint: &Endpoint) -> Result<()> { error!("delete tracker failed: {}", err.to_string()); } } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::from(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))), - }; + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; let disks = store.get_disks(pool_idx, set_idx).await?; for disk in disks.into_iter() { if disk.is_none() { @@ -377,9 +366,7 @@ impl HealRoutine { Err(err) => d_err = Some(err), } } else { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock.as_ref().expect("Not init"); + let store = new_object_layer_fn().expect("errServerNotInitialized"); if task.object.is_empty() { match store.heal_bucket(&task.bucket, &task.opts).await { Ok(res) => { @@ -429,9 +416,8 @@ impl HealRoutine { // } async fn heal_disk_format(opts: HealOpts) -> Result<(HealResultItem, Option)> { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock.as_ref().expect("Not init"); + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; + let (res, err) = store.heal_format(opts.dry_run).await?; // return any error, ignore error returned when disks have // already healed. diff --git a/ecstore/src/heal/data_scanner.rs b/ecstore/src/heal/data_scanner.rs index c6b1072e5..d45949371 100644 --- a/ecstore/src/heal/data_scanner.rs +++ b/ecstore/src/heal/data_scanner.rs @@ -102,17 +102,14 @@ pub async fn init_data_scanner() { } async fn run_data_scanner() { - let mut cycle_info = CurrentScannerCycle::default(); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => { - info!("errServerNotInitialized"); - return; - } + let Some(store) = new_object_layer_fn() else { + error!("errServerNotInitialized"); + return; }; - let mut buf = read_config(store, &DATA_USAGE_BLOOM_NAME_PATH) + + let mut cycle_info = CurrentScannerCycle::default(); + + let mut buf = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH) .await .map_or(Vec::new(), |buf| buf); match buf.len().cmp(&8) { @@ -146,7 +143,7 @@ async fn run_data_scanner() { globalScannerMetrics.write().await.set_cycle(Some(cycle_info.clone())).await; } - let bg_heal_info = read_background_heal_info(store).await; + let bg_heal_info = read_background_heal_info(store.clone()).await; let scan_mode = get_cycle_scan_mode(cycle_info.current, bg_heal_info.bitrot_start_cycle, bg_heal_info.bitrot_start_time).await; if bg_heal_info.current_scan_mode != scan_mode { @@ -156,7 +153,7 @@ async fn run_data_scanner() { new_heal_info.bitrot_start_time = SystemTime::now(); new_heal_info.bitrot_start_cycle = cycle_info.current; } - save_background_heal_info(store, &new_heal_info).await; + save_background_heal_info(store.clone(), &new_heal_info).await; } // Wait before starting next cycle and wait on startup. let (tx, rx) = mpsc::channel(100); @@ -165,7 +162,7 @@ async fn run_data_scanner() { }); let mut res = HashMap::new(); res.insert("cycle".to_string(), cycle_info.current.to_string()); - match store.ns_scanner(tx, cycle_info.current as usize, scan_mode).await { + match store.clone().ns_scanner(tx, cycle_info.current as usize, scan_mode).await { Ok(_) => { cycle_info.next += 1; cycle_info.current = 0; @@ -178,7 +175,7 @@ async fn run_data_scanner() { globalScannerMetrics.write().await.set_cycle(Some(cycle_info.clone())).await; let mut tmp = Vec::new(); tmp.write_u64::(cycle_info.next).unwrap(); - let _ = save_config(store, &DATA_USAGE_BLOOM_NAME_PATH, &tmp).await; + let _ = save_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH, &tmp).await; } Err(err) => { res.insert("error".to_string(), err.to_string()); @@ -206,7 +203,7 @@ impl Default for BackgroundHealInfo { } } -async fn read_background_heal_info(store: &ECStore) -> BackgroundHealInfo { +async fn read_background_heal_info(store: Arc) -> BackgroundHealInfo { if *GLOBAL_IsErasureSD.read().await { return BackgroundHealInfo::default(); } @@ -220,7 +217,7 @@ async fn read_background_heal_info(store: &ECStore) -> BackgroundHealInfo { serde_json::from_slice::(&buf).map_or(BackgroundHealInfo::default(), |b| b) } -async fn save_background_heal_info(store: &ECStore, info: &BackgroundHealInfo) { +async fn save_background_heal_info(store: Arc, info: &BackgroundHealInfo) { if *GLOBAL_IsErasureSD.read().await { return; } diff --git a/ecstore/src/heal/data_usage.rs b/ecstore/src/heal/data_usage.rs index 067bcc641..17b8ec5d0 100644 --- a/ecstore/src/heal/data_usage.rs +++ b/ecstore/src/heal/data_usage.rs @@ -3,7 +3,7 @@ use std::{collections::HashMap, time::SystemTime}; use lazy_static::lazy_static; use serde::{Deserialize, Serialize}; use tokio::sync::mpsc::Receiver; -use tracing::info; +use tracing::error; use crate::{ config::common::save_config, @@ -107,25 +107,21 @@ pub struct DataUsageInfo { } pub async fn store_data_usage_in_backend(mut rx: Receiver) { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => { - info!("errServerNotInitialized"); - return; - } + let Some(store) = new_object_layer_fn() else { + error!("errServerNotInitialized"); + return; }; + let mut attempts = 1; loop { match rx.recv().await { Some(data_usage_info) => { if let Ok(data) = serde_json::to_vec(&data_usage_info) { if attempts > 10 { - let _ = save_config(store, &format!("{}{}", *DATA_USAGE_OBJ_NAME_PATH, ".bkp"), &data).await; + let _ = save_config(store.clone(), &format!("{}{}", *DATA_USAGE_OBJ_NAME_PATH, ".bkp"), &data).await; attempts += 1; } - let _ = save_config(store, &DATA_USAGE_OBJ_NAME_PATH, &data).await; + let _ = save_config(store.clone(), &DATA_USAGE_OBJ_NAME_PATH, &data).await; attempts += 1; } else { continue; diff --git a/ecstore/src/heal/data_usage_cache.rs b/ecstore/src/heal/data_usage_cache.rs index 371de38cf..09be08b99 100644 --- a/ecstore/src/heal/data_usage_cache.rs +++ b/ecstore/src/heal/data_usage_cache.rs @@ -11,7 +11,6 @@ use path_clean::PathClean; use rand::Rng; use rmp_serde::Serializer; use s3s::dto::ReplicationConfiguration; -use s3s::{S3Error, S3ErrorCode}; use serde::{Deserialize, Serialize}; use std::collections::{HashMap, HashSet}; use std::hash::{DefaultHasher, Hash, Hasher}; @@ -442,18 +441,14 @@ impl DataUsageCache { } pub async fn save(&self, name: &str) -> Result<()> { + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; let buf = self.marshal_msg()?; let buf_clone = buf.clone(); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::from(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))), - }; + let store_clone = store.clone(); let name_clone = name.to_string(); tokio::spawn(async move { - let _ = save_config(&store_clone, &format!("{}{}", &name_clone, ".bkp"), &buf_clone).await; + let _ = save_config(store_clone, &format!("{}{}", &name_clone, ".bkp"), &buf_clone).await; }); save_config(store, name, &buf).await } diff --git a/ecstore/src/heal/heal_commands.rs b/ecstore/src/heal/heal_commands.rs index d27dad414..06f97fa13 100644 --- a/ecstore/src/heal/heal_commands.rs +++ b/ecstore/src/heal/heal_commands.rs @@ -256,12 +256,8 @@ impl HealingTracker { pub async fn save(&mut self) -> Result<()> { let _ = self.mu.write().await; if self.pool_index.is_none() || self.set_index.is_none() || self.disk_index.is_none() { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::from_string("Not init".to_string())), - }; + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; + (self.pool_index, self.set_index, self.disk_index) = store.get_pool_and_set(&self.id).await?; } diff --git a/ecstore/src/heal/heal_ops.rs b/ecstore/src/heal/heal_ops.rs index a271f6c36..22b376e97 100644 --- a/ecstore/src/heal/heal_ops.rs +++ b/ecstore/src/heal/heal_ops.rs @@ -7,8 +7,6 @@ use super::{ HEAL_ITEM_BUCKET_METADATA, }, }; -use crate::heal::heal_commands::{HEAL_ITEM_BUCKET, HEAL_ITEM_OBJECT}; -use crate::store_api::StorageAPI; use crate::{ config::common::CONFIG_PREFIX, disk::RUSTFS_META_BUCKET, @@ -27,8 +25,11 @@ use crate::{ new_object_layer_fn, utils::path::has_profix, }; +use crate::{ + heal::heal_commands::{HEAL_ITEM_BUCKET, HEAL_ITEM_OBJECT}, + store_api::StorageAPI, +}; use lazy_static::lazy_static; -use s3s::{S3Error, S3ErrorCode}; use std::{ collections::HashMap, future::Future, @@ -387,12 +388,7 @@ impl HealSequence { } async fn heal_rustfs_sys_meta(h: Arc>, meta_prefix: &str) -> Result<()> { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::from(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))), - }; + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; let setting = h.read().await.setting; store .heal_objects(RUSTFS_META_BUCKET, meta_prefix, &setting, h.clone(), true) @@ -431,12 +427,7 @@ impl HealSequence { } (hs_w.object.clone(), hs_w.setting) }; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::from(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))), - }; + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; store.heal_objects(bucket, &object, &setting, hs.clone(), false).await } diff --git a/ecstore/src/pools.rs b/ecstore/src/pools.rs index f732fc2ff..a7f95a49d 100644 --- a/ecstore/src/pools.rs +++ b/ecstore/src/pools.rs @@ -58,8 +58,10 @@ impl PoolMeta { self.pools[idx].decommission.is_some() } - pub async fn load(&mut self, store: &ECStore) -> Result<()> { - let data = match read_config(store, POOL_META_NAME).await { + pub async fn load(&mut self) -> Result<()> { + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; + + let data = match read_config(store.clone(), POOL_META_NAME).await { Ok(data) => { if data.is_empty() { return Ok(()); @@ -94,7 +96,7 @@ impl PoolMeta { Ok(()) } - pub async fn save(&self) -> Result<()> { + pub async fn save(&self, _pools: Vec>) -> Result<()> { if self.dont_save { return Ok(()); } @@ -105,12 +107,11 @@ impl PoolMeta { self.serialize(&mut Serializer::new(&mut buf))?; data.write_all(&buf)?; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(Error::from_string("errServerNotInitialized".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(Error::from_string("errServerNotInitialized".to_string())); }; + + // FIXME: save_config(store, &POOL_META_NAME, &data).await } @@ -171,7 +172,10 @@ pub struct PoolSpaceInfo { impl ECStore { pub async fn status(&self, idx: usize) -> Result { let space_info = self.get_decommission_pool_space_info(idx).await?; - let mut pool_info = self.pool_meta.read().unwrap().pools[idx].clone(); + + let pool_meta = self.pool_meta.read().await; + + let mut pool_info = pool_meta.pools[idx].clone(); if let Some(d) = pool_info.decommission.as_mut() { d.total_size = space_info.total; d.current_size = space_info.free; @@ -204,7 +208,7 @@ impl ECStore { } } - pub async fn decommission_cancel(&mut self, idx: usize) -> Result<()> { + pub async fn decommission_cancel(&self, idx: usize) -> Result<()> { if self.single_pool() { return Err(Error::msg("InvalidArgument")); } @@ -217,11 +221,12 @@ impl ECStore { return Err(Error::new(StorageError::DecommissionNotStarted)); } - if self.pool_meta.write().unwrap().decommission_cancel(idx) { - // FIXME: + let mut lock = self.pool_meta.write().await; + if lock.decommission_cancel(idx) { + lock.save(self.pools.clone()).await?; } - unimplemented!() + Ok(()) } } diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 7e3124cf0..9787a4ad9 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -48,14 +48,13 @@ use http::HeaderMap; use lazy_static::lazy_static; use rand::Rng; use s3s::dto::{BucketVersioningStatus, ObjectLockConfiguration, ObjectLockEnabled, VersioningConfiguration}; -use s3s::{S3Error, S3ErrorCode}; use std::cmp::Ordering; use std::process::exit; use std::slice::Iter; use std::time::SystemTime; use std::{ collections::{HashMap, HashSet}, - sync::{Arc, RwLock as std_RwLock}, + sync::Arc, time::Duration, }; use time::OffsetDateTime; @@ -76,30 +75,30 @@ pub struct ECStore { pub pools: Vec>, pub peer_sys: S3PeerSys, // pub local_disks: Vec, - pub pool_meta: std_RwLock, + pub pool_meta: RwLock, pub decommission_cancelers: Vec>, } -impl Clone for ECStore { - fn clone(&self) -> Self { - let pool_meta = match self.pool_meta.read() { - Ok(pool_meta) => pool_meta.clone(), - Err(_) => PoolMeta::default(), - }; - Self { - id: self.id.clone(), - disk_map: self.disk_map.clone(), - pools: self.pools.clone(), - peer_sys: self.peer_sys.clone(), - pool_meta: std_RwLock::new(pool_meta), - decommission_cancelers: self.decommission_cancelers.clone(), - } - } -} +// impl Clone for ECStore { +// fn clone(&self) -> Self { +// let pool_meta = match self.pool_meta.read() { +// Ok(pool_meta) => pool_meta.clone(), +// Err(_) => PoolMeta::default(), +// }; +// Self { +// id: self.id.clone(), +// disk_map: self.disk_map.clone(), +// pools: self.pools.clone(), +// peer_sys: self.peer_sys.clone(), +// pool_meta: std_RwLock::new(pool_meta), +// decommission_cancelers: self.decommission_cancelers.clone(), +// } +// } +// } impl ECStore { #[allow(clippy::new_ret_no_self)] - pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result { + pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result> { // let layouts = DisksLayout::from_volumes(endpoints.as_slice())?; let mut deployment_id = None; @@ -216,14 +215,14 @@ impl ECStore { pool_meta.dont_save = true; let decommission_cancelers = vec![None; pools.len()]; - let ec = ECStore { + let ec = Arc::new(ECStore { id: deployment_id.unwrap(), disk_map, pools, peer_sys, pool_meta: pool_meta.into(), decommission_cancelers, - }; + }); set_object_layer(ec.clone()).await; @@ -234,11 +233,11 @@ impl ECStore { Ok(ec) } - pub async fn init(&self) -> Result<()> { + pub async fn init(api: Arc) -> Result<()> { config::init(); - GLOBAL_ConfigSys.init(self).await?; + GLOBAL_ConfigSys.init(api.clone()).await?; - let buckets_list = self + let buckets_list = api .list_bucket(&BucketOptions { no_metadata: true, ..Default::default() @@ -248,7 +247,9 @@ impl ECStore { let buckets = buckets_list.iter().map(|v| v.name.clone()).collect(); - init_bucket_metadata_sys(self.clone(), buckets).await; + // FIXME: + + init_bucket_metadata_sys(api.clone(), buckets).await; Ok(()) } @@ -495,12 +496,12 @@ impl ECStore { ServerPoolsAvailableSpace(server_pools) } - fn is_suspended(&self, idx: usize) -> bool { + async fn is_suspended(&self, idx: usize) -> bool { // TODO: LOCK - match self.pool_meta.read() { - Ok(pool_meta) => pool_meta.is_suspended(idx), - Err(_) => false, - } + + let pool_meta = self.pool_meta.read().await; + + pool_meta.is_suspended(idx) } async fn get_pool_idx(&self, bucket: &str, object: &str, size: i64) -> Result { @@ -594,8 +595,9 @@ impl ECStore { let mut def_pool = PoolObjInfo::default(); let mut has_def_pool = false; + let pool_meta = self.pool_meta.read().await; for pinfo in ress.iter() { - if opts.skip_decommissioned && self.pool_meta.read().unwrap().is_suspended(pinfo.index) { + if opts.skip_decommissioned && pool_meta.is_suspended(pinfo.index) { continue; } @@ -605,13 +607,13 @@ impl ECStore { // } if pinfo.err.is_none() { - return Ok((pinfo.clone(), self.pools_with_object(&ress, opts))); + return Ok((pinfo.clone(), self.pools_with_object(&ress, opts).await)); } let err = pinfo.err.as_ref().unwrap(); if is_err_read_quorum(err) && !opts.metadata_chg { - return Ok((pinfo.clone(), self.pools_with_object(&ress, opts))); + return Ok((pinfo.clone(), self.pools_with_object(&ress, opts).await)); } def_pool = pinfo.clone(); @@ -633,10 +635,11 @@ impl ECStore { Err(to_object_err(Error::new(DiskError::FileNotFound), vec![bucket, object])) } - fn pools_with_object(&self, pools: &Vec, opts: &ObjectOptions) -> Vec { + async fn pools_with_object(&self, pools: &Vec, opts: &ObjectOptions) -> Vec { let mut errs = Vec::new(); + let pool_meta = self.pool_meta.read().await; for pool in pools.iter() { - if opts.skip_decommissioned && self.pool_meta.read().unwrap().is_suspended(pool.index) { + if opts.skip_decommissioned && pool_meta.is_suspended(pool.index) { continue; } // TODO:SkipRebalancing @@ -892,8 +895,11 @@ impl ECStore { pub async fn reload_pool_meta(&self) -> Result<()> { let mut meta = PoolMeta::default(); - meta.load(self).await?; - *self.pool_meta.write().unwrap() = meta; + meta.load().await?; + + let mut pool_meta = self.pool_meta.write().await; + *pool_meta = meta; + // *self.pool_meta.write().unwrap() = meta; Ok(()) } } @@ -1263,7 +1269,7 @@ impl StorageAPI for ECStore { meta.versioning_config_xml = xml::serialize::(&enableVersioningConfig)?; } - meta.save(self).await.map_err(|e| to_object_err(e, vec![bucket]))?; + meta.save().await.map_err(|e| to_object_err(e, vec![bucket]))?; set_bucket_metadata(bucket.to_string(), meta).await; @@ -1787,7 +1793,7 @@ impl StorageAPI for ECStore { } for (idx, pool) in self.pools.iter().enumerate() { - if self.is_suspended(idx) { + if self.is_suspended(idx).await { continue; } @@ -2037,14 +2043,8 @@ impl StorageAPI for ECStore { }; if opts_clone.remove && !opts_clone.dry_run { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => { - return Err(Error::from(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))) - } - }; + let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; + if let Err(err) = store.check_abandoned_parts(&bucket, &entry.name, &opts_clone).await { info!("unable to check object {}/{} for abandoned data: {}", bucket, entry.name, err.to_string()); } diff --git a/madmin/src/heal_command.rs b/madmin/src/heal_command.rs index ad87352af..0e26f100d 100644 --- a/madmin/src/heal_command.rs +++ b/madmin/src/heal_command.rs @@ -55,17 +55,12 @@ pub async fn get_local_background_heal_status() -> (BgHealState, bool) { heal_disks_map.insert(ep.to_string()); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => { - let healing = GLOBAL_BackgroundHealState.read().await.get_local_healing_disks().await; - for disk in healing.values() { - status.heal_disks.push(disk.endpoint.clone()); - } - return (status, true); + let Some(store) = new_object_layer_fn() else { + let healing = GLOBAL_BackgroundHealState.read().await.get_local_healing_disks().await; + for disk in healing.values() { + status.heal_disks.push(disk.endpoint.clone()); } + return (status, true); }; let si = store.local_storage_info().await; diff --git a/rustfs/src/admin/handlers.rs b/rustfs/src/admin/handlers.rs index 377f355dc..a72a45b36 100644 --- a/rustfs/src/admin/handlers.rs +++ b/rustfs/src/admin/handlers.rs @@ -116,11 +116,8 @@ impl Operation for AccountInfoHandler { warn!("AccountInfoHandler cread {:?}", &cred); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; // test policy @@ -257,11 +254,8 @@ impl Operation for ListPools { async fn call(&self, _req: S3Request, _params: Params<'_, '_>) -> S3Result> { warn!("handle ListPools"); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let Some(endpoints) = GLOBAL_Endpoints.get() else { @@ -341,11 +335,8 @@ impl Operation for StatusPool { return Err(s3_error!(InvalidArgument)); }; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let pools_status = store.status(idx).await.map_err(to_s3_error)?; @@ -410,22 +401,18 @@ impl Operation for CancelDecommission { } }; - let Some(_idx) = has_idx else { + let Some(idx) = has_idx else { warn!("specified pool {} not found, please specify a valid pool", &query.pool); return Err(s3_error!(InvalidArgument)); }; - let layer = new_object_layer_fn(); - let lock = layer.write().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; - // FIXME: - // store.decommission_cancel(idx).await; + store.decommission_cancel(idx).await.map_err(to_s3_error)?; - return Err(s3_error!(NotImplemented)); + Ok(S3Response::new((StatusCode::OK, Body::default()))) } } diff --git a/rustfs/src/grpc.rs b/rustfs/src/grpc.rs index e4be22302..7a6c715c9 100644 --- a/rustfs/src/grpc.rs +++ b/rustfs/src/grpc.rs @@ -3,12 +3,11 @@ use std::{ error::Error, io::{Cursor, ErrorKind}, pin::Pin, - sync::Arc, }; use ecstore::{ admin_server_info::get_local_server_property, - bucket::{metadata::load_bucket_metadata, metadata_sys::GLOBAL_BucketMetadataSys}, + bucket::{metadata::load_bucket_metadata, metadata_sys}, disk::{ DeleteOptions, DiskAPI, DiskInfoOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, Reader, UpdateMetadataOpts, WalkDirOptions, @@ -1519,17 +1518,13 @@ impl Node for NodeService { _request: Request, ) -> Result, Status> { // let request = request.into_inner(); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(LocalStorageInfoResponse { - success: false, - storage_info: vec![], - error_info: Some("errServerNotInitialized".to_string()), - })) - } + + let Some(store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(LocalStorageInfoResponse { + success: false, + storage_info: vec![], + error_info: Some("errServerNotInitialized".to_string()), + })); }; let info = store.local_storage_info().await; @@ -1796,20 +1791,16 @@ impl Node for NodeService { })); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(LoadBucketMetadataResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(LoadBucketMetadataResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; + match load_bucket_metadata(store, &bucket).await { Ok(meta) => { - GLOBAL_BucketMetadataSys.write().await.set(bucket, Arc::new(meta)).await; + metadata_sys::set_bucket_metadata(bucket, meta).await; Ok(tonic::Response::new(LoadBucketMetadataResponse { success: true, error_info: None, @@ -1846,16 +1837,11 @@ impl Node for NodeService { })); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(DeletePolicyResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(DeletePolicyResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() @@ -1870,16 +1856,11 @@ impl Node for NodeService { error_info: Some("policy name is missing".to_string()), })); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(LoadPolicyResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(LoadPolicyResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() } @@ -1898,16 +1879,11 @@ impl Node for NodeService { } let _user_type = request.user_type; let _is_group = request.is_group; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(LoadPolicyMappingResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(LoadPolicyMappingResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() } @@ -1921,16 +1897,11 @@ impl Node for NodeService { error_info: Some("access_key name is missing".to_string()), })); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(DeleteUserResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(DeleteUserResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() @@ -1948,16 +1919,11 @@ impl Node for NodeService { error_info: Some("access_key name is missing".to_string()), })); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(DeleteServiceAccountResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(DeleteServiceAccountResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() } @@ -1973,16 +1939,11 @@ impl Node for NodeService { })); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(LoadUserResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(LoadUserResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() @@ -2001,16 +1962,11 @@ impl Node for NodeService { })); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(LoadServiceAccountResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(LoadServiceAccountResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() } @@ -2025,16 +1981,11 @@ impl Node for NodeService { })); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(LoadGroupResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(LoadGroupResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() } @@ -2043,16 +1994,11 @@ impl Node for NodeService { &self, _request: Request, ) -> Result, Status> { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(ReloadSiteReplicationConfigResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(ReloadSiteReplicationConfigResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; todo!() } @@ -2112,16 +2058,11 @@ impl Node for NodeService { &self, _request: Request, ) -> Result, Status> { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(ReloadPoolMetaResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(ReloadPoolMetaResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; match store.reload_pool_meta().await { Ok(_) => Ok(tonic::Response::new(ReloadPoolMetaResponse { @@ -2136,16 +2077,11 @@ impl Node for NodeService { } async fn stop_rebalance(&self, _request: Request) -> Result, Status> { - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let _store = match lock.as_ref() { - Some(s) => s, - None => { - return Ok(tonic::Response::new(StopRebalanceResponse { - success: false, - error_info: Some("errServerNotInitialized".to_string()), - })) - } + let Some(_store) = new_object_layer_fn() else { + return Ok(tonic::Response::new(StopRebalanceResponse { + success: false, + error_info: Some("errServerNotInitialized".to_string()), + })); }; // todo diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index 4b25b819f..9ddde2bfc 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -87,7 +87,7 @@ async fn run(opt: config::Opt) -> Result<()> { ); } - set_global_endpoints(endpoint_pools.as_ref().clone()).map_err(|err| Error::from_string(err.to_string()))?; + set_global_endpoints(endpoint_pools.as_ref().clone()); update_erasure_type(setup_type).await; // 初始化本地磁盘 @@ -186,7 +186,7 @@ async fn run(opt: config::Opt) -> Result<()> { .await .map_err(|err| Error::from_string(err.to_string()))?; - store.init().await.map_err(|err| { + ECStore::init(store.clone()).await.map_err(|err| { error!("init faild {:?}", &err); Error::from_string(err.to_string()) })?; diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 93d327982..8241b2353 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -95,11 +95,8 @@ impl S3 for FS { .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store @@ -134,11 +131,8 @@ impl S3 for FS { async fn delete_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; // TODO: DeleteBucketInput 没有force参数? - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store @@ -175,11 +169,8 @@ impl S3 for FS { let objects: Vec = vec![dobj]; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let (dobjs, _errs) = store .delete_objects(&bucket, objects, ObjectOptions::default()) @@ -245,11 +236,8 @@ impl S3 for FS { }) .collect(); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let (dobjs, _errs) = store @@ -289,11 +277,8 @@ impl S3 for FS { // mc get 1 let input = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store @@ -320,11 +305,8 @@ impl S3 for FS { let h = HeaderMap::new(); let opts = &ObjectOptions::default(); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let reader = store @@ -365,11 +347,8 @@ impl S3 for FS { async fn head_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store @@ -386,11 +365,8 @@ impl S3 for FS { // mc get 2 let HeadObjectInput { bucket, key, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let info = store @@ -428,11 +404,8 @@ impl S3 for FS { async fn list_buckets(&self, _: S3Request) -> S3Result> { // mc ls - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let bucket_infos = store.list_bucket(&BucketOptions::default()).await.map_err(to_s3_error)?; @@ -486,11 +459,8 @@ impl S3 for FS { let prefix = prefix.unwrap_or_default(); let delimiter = delimiter.unwrap_or_default(); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let object_infos = store @@ -591,11 +561,8 @@ impl S3 for FS { let mut reader = PutObjReader::new(body, content_length as usize); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let mut metadata = extract_metadata(&req.headers); @@ -636,11 +603,8 @@ impl S3 for FS { // debug!("create_multipart_upload meta {:?}", &metadata); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let mut metadata = extract_metadata(&req.headers); @@ -702,11 +666,8 @@ impl S3 for FS { let mut data = PutObjReader::new(body, content_length as usize); let opts = ObjectOptions::default(); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; // TODO: hash_reader @@ -773,11 +734,8 @@ impl S3 for FS { uploaded_parts.push(CompletePart::from(part)); } - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store @@ -802,11 +760,8 @@ impl S3 for FS { bucket, key, upload_id, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let opts = &ObjectOptions::default(); @@ -845,11 +800,8 @@ impl S3 for FS { async fn put_bucket_tagging(&self, req: S3Request) -> S3Result> { let PutBucketTaggingInput { bucket, tagging, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store @@ -889,11 +841,9 @@ impl S3 for FS { .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; let tags = encode_tags(tagging.tag_set); @@ -912,11 +862,9 @@ impl S3 for FS { async fn get_object_tagging(&self, req: S3Request) -> S3Result> { let GetObjectTaggingInput { bucket, key: object, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; // TODO: version let tags = store @@ -939,11 +887,9 @@ impl S3 for FS { ) -> S3Result> { let DeleteObjectTaggingInput { bucket, key: object, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; // TODO: Replicate // TODO: version @@ -961,11 +907,8 @@ impl S3 for FS { req: S3Request, ) -> S3Result> { let GetBucketVersioningInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = match lock.as_ref() { - Some(s) => s, - None => return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())), + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; store @@ -1018,11 +961,9 @@ impl S3 for FS { async fn get_bucket_policy(&self, req: S3Request) -> S3Result> { let GetBucketPolicyInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1047,11 +988,9 @@ impl S3 for FS { async fn put_bucket_policy(&self, req: S3Request) -> S3Result> { let PutBucketPolicyInput { bucket, policy, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1084,11 +1023,9 @@ impl S3 for FS { ) -> S3Result> { let DeleteBucketPolicyInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1109,11 +1046,9 @@ impl S3 for FS { ) -> S3Result> { let GetBucketLifecycleConfigurationInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1169,11 +1104,9 @@ impl S3 for FS { ) -> S3Result> { let DeleteBucketLifecycleInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1193,11 +1126,9 @@ impl S3 for FS { ) -> S3Result> { let GetBucketEncryptionInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1232,11 +1163,9 @@ impl S3 for FS { info!("sse_config {:?}", &server_side_encryption_configuration); - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1258,11 +1187,9 @@ impl S3 for FS { ) -> S3Result> { let DeleteBucketEncryptionInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1308,11 +1235,9 @@ impl S3 for FS { let Some(input_cfg) = object_lock_configuration else { return Err(s3_error!(InvalidArgument)) }; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1334,11 +1259,9 @@ impl S3 for FS { ) -> S3Result> { let GetBucketReplicationInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1368,11 +1291,9 @@ impl S3 for FS { .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1395,11 +1316,9 @@ impl S3 for FS { ) -> S3Result> { let DeleteBucketReplicationInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1420,11 +1339,9 @@ impl S3 for FS { ) -> S3Result> { let GetBucketNotificationConfigurationInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1469,11 +1386,9 @@ impl S3 for FS { .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1494,11 +1409,9 @@ impl S3 for FS { async fn get_bucket_acl(&self, req: S3Request) -> S3Result> { let GetBucketAclInput { bucket, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1532,11 +1445,9 @@ impl S3 for FS { // TODO:checkRequestAuthType - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; store .get_bucket_info(&bucket, &BucketOptions::default()) @@ -1570,11 +1481,9 @@ impl S3 for FS { async fn get_object_acl(&self, req: S3Request) -> S3Result> { let GetObjectAclInput { bucket, key, .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; if let Err(e) = store.get_object_info(&bucket, &key, &ObjectOptions::default()).await { return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e))); @@ -1607,11 +1516,9 @@ impl S3 for FS { .. } = req.input; - let layer = new_object_layer_fn(); - let lock = layer.read().await; - let store = lock - .as_ref() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?; + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; if let Err(e) = store.get_object_info(&bucket, &key, &ObjectOptions::default()).await { return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", e)));