use Arc<ECStroe>

This commit is contained in:
weisd
2024-11-28 14:47:01 +08:00
parent 9e6ec15f32
commit 73efd4f493
24 changed files with 411 additions and 650 deletions
+3 -4
View File
@@ -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<Json<Vec<PoolStatus>>> {
//
// 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 _,
+7 -13
View File
@@ -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
}
+12 -9
View File
@@ -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<u8> = 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<ECStore>) -> 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<BucketMetadata> {
pub async fn load_bucket_metadata(api: Arc<ECStore>, bucket: &str) -> Result<BucketMetadata> {
load_bucket_metadata_parse(api, bucket, true).await
}
pub async fn load_bucket_metadata_parse(api: &ECStore, bucket: &str, parse: bool) -> Result<BucketMetadata> {
let mut bm = match read_bucket_metadata(api, bucket).await {
pub async fn load_bucket_metadata_parse(api: Arc<ECStore>, bucket: &str, parse: bool) -> Result<BucketMetadata> {
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<BucketMetadata> {
async fn read_bucket_metadata(api: Arc<ECStore>, bucket: &str) -> Result<BucketMetadata> {
if bucket.is_empty() {
error!("bucket name empty");
return Err(Error::msg("invalid argument"));
+64 -81
View File
@@ -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<RwLock<BucketMetadataSys>> = Arc::new(RwLock::new(BucketMetadataSys::new()));
pub static ref GLOBAL_BucketMetadataSys: OnceLock<Arc<RwLock<BucketMetadataSys>>> = OnceLock::new();
}
pub async fn init_bucket_metadata_sys(api: ECStore, buckets: Vec<String>) {
let mut sys = GLOBAL_BucketMetadataSys.write().await;
sys.init(api, buckets).await
pub async fn init_bucket_metadata_sys(api: Arc<ECStore>, buckets: Vec<String>) {
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<RwLock<BucketMetadataSys>> {
GLOBAL_BucketMetadataSys.clone()
// panic if not init
pub(super) fn get_bucket_metadata_sys() -> Arc<RwLock<BucketMetadataSys>> {
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<Arc<BucketMetadata>> {
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<u8>) -> Result<OffsetDateTime> {
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<OffsetDateTime> {
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<Option<NotificationConfiguration>> {
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<BucketMetadata> {
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<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.created_at(bucket).await
}
#[derive(Debug, Default)]
#[derive(Debug)]
pub struct BucketMetadataSys {
metadata_map: RwLock<HashMap<String, Arc<BucketMetadata>>>,
api: Option<ECStore>,
api: Arc<ECStore>,
initialized: RwLock<bool>,
}
impl BucketMetadataSys {
pub fn new() -> Self {
Self::default()
pub fn new(api: Arc<ECStore>) -> Self {
Self {
metadata_map: RwLock::new(HashMap::new()),
api,
initialized: RwLock::new(false),
}
}
pub async fn init(&mut self, api: ECStore, buckets: Vec<String>) {
// if api.is_none() {
// return Err(Error::msg("errServerNotInitialized"));
// }
self.api = Some(api);
pub async fn init(&mut self, buckets: Vec<String>) {
let _ = self.init_internal(buckets).await;
}
async fn init_internal(&self, buckets: Vec<String>) -> 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<u8>, parse: bool) -> Result<OffsetDateTime> {
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<BucketMetadata> {
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<BucketMetadata>, 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))
}
}
+1 -1
View File
@@ -22,7 +22,7 @@ impl PolicySys {
args.is_owner
}
pub async fn get(bucket: &str) -> Result<BucketPolicy> {
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?;
+1 -1
View File
@@ -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?;
+16 -15
View File
@@ -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<Vec<u8>> {
pub async fn read_config(api: Arc<ECStore>, file: &str) -> Result<Vec<u8>> {
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<u8>, ObjectInfo)> {
async fn read_config_with_metadata(api: Arc<ECStore>, file: &str, opts: &ObjectOptions) -> Result<(Vec<u8>, 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<ECStore>, 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<ECStore>, 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<Config> {
async fn new_and_save_server_config(api: Arc<ECStore>) -> Result<Config> {
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<Config> {
pub async fn read_config_without_migrate(api: Arc<ECStore>) -> Result<Config> {
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<Config> {
read_server_config(api, data.as_slice()).await
}
async fn read_server_config(api: &ECStore, data: &[u8]) -> Result<Config> {
async fn read_server_config(api: Arc<ECStore>, data: &[u8]) -> Result<Config> {
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<Config> {
Ok(cfg.merge())
}
async fn save_server_config(api: &ECStore, cfg: &Config) -> Result<()> {
async fn save_server_config(api: Arc<ECStore>, 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<ECStore>) {
// 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<ECStore>) -> 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<ECStore>, 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) {
+3 -3
View File
@@ -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<storageclass::Config> = 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<ECStore>) -> Result<()> {
let mut cfg = read_config_without_migrate(api.clone().clone()).await?;
lookup_configs(&mut cfg, api).await;
-1
View File
@@ -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};
+3 -13
View File
@@ -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?);
+11 -10
View File
@@ -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<RwLock<Option<ECStore>>> = Arc::new(RwLock::new(None));
pub static ref GLOBAL_OBJECT_API: OnceLock<Arc<ECStore>> = OnceLock::new();
pub static ref GLOBAL_LOCAL_DISK: Arc<RwLock<Vec<Option<DiskStore>>>> = Arc::new(RwLock::new(Vec::new()));
pub static ref GLOBAL_IsErasure: RwLock<bool> = RwLock::new(false);
pub static ref GLOBAL_IsDistErasure: RwLock<bool> = RwLock::new(false);
@@ -44,20 +43,22 @@ pub async fn get_global_deployment_id() -> Uuid {
*id_ptr
}
pub fn set_global_endpoints(eps: Vec<PoolEndpoints>) -> Result<()> {
pub fn set_global_endpoints(eps: Vec<PoolEndpoints>) {
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<RwLock<Option<ECStore>>> {
GLOBAL_OBJECT_API.clone()
pub fn new_object_layer_fn() -> Option<Arc<ECStore>> {
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<ECStore>) {
GLOBAL_OBJECT_API.set(o).expect("set_object_layer fail ")
}
pub async fn is_dist_erasure() -> bool {
+8 -22
View File
@@ -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<Error>)> {
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.
+13 -16
View File
@@ -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::<LittleEndian>(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<ECStore>) -> 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::<BackgroundHealInfo>(&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<ECStore>, info: &BackgroundHealInfo) {
if *GLOBAL_IsErasureSD.read().await {
return;
}
+7 -11
View File
@@ -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<DataUsageInfo>) {
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;
+3 -8
View File
@@ -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
}
+2 -6
View File
@@ -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?;
}
+6 -15
View File
@@ -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<RwLock<HealSequence>>, 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
}
+18 -13
View File
@@ -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<Arc<Sets>>) -> 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<PoolStatus> {
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(())
}
}
+48 -48
View File
@@ -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<Arc<Sets>>,
pub peer_sys: S3PeerSys,
// pub local_disks: Vec<DiskStore>,
pub pool_meta: std_RwLock<PoolMeta>,
pub pool_meta: RwLock<PoolMeta>,
pub decommission_cancelers: Vec<Option<usize>>,
}
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<Self> {
pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result<Arc<Self>> {
// 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<ECStore>) -> 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<usize> {
@@ -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<PoolObjInfo>, opts: &ObjectOptions) -> Vec<PoolErr> {
async fn pools_with_object(&self, pools: &Vec<PoolObjInfo>, opts: &ObjectOptions) -> Vec<PoolErr> {
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::<VersioningConfiguration>(&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());
}
+5 -10
View File
@@ -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;
+11 -24
View File
@@ -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<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
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())))
}
}
+70 -134
View File
@@ -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<LocalStorageInfoRequest>,
) -> Result<Response<LocalStorageInfoResponse>, 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<ReloadSiteReplicationConfigRequest>,
) -> Result<Response<ReloadSiteReplicationConfigResponse>, 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<ReloadPoolMetaRequest>,
) -> Result<Response<ReloadPoolMetaResponse>, 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<StopRebalanceRequest>) -> Result<Response<StopRebalanceResponse>, 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
+2 -2
View File
@@ -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())
})?;
+97 -190
View File
@@ -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<DeleteBucketInput>) -> S3Result<S3Response<DeleteBucketOutput>> {
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<ObjectToDelete> = 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<HeadBucketInput>) -> S3Result<S3Response<HeadBucketOutput>> {
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<ListBucketsInput>) -> S3Result<S3Response<ListBucketsOutput>> {
// 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<PutBucketTaggingInput>) -> S3Result<S3Response<PutBucketTaggingOutput>> {
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<GetObjectTaggingInput>) -> S3Result<S3Response<GetObjectTaggingOutput>> {
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<S3Response<DeleteObjectTaggingOutput>> {
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<GetBucketVersioningInput>,
) -> S3Result<S3Response<GetBucketVersioningOutput>> {
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<GetBucketPolicyInput>) -> S3Result<S3Response<GetBucketPolicyOutput>> {
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<PutBucketPolicyInput>) -> S3Result<S3Response<PutBucketPolicyOutput>> {
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<S3Response<DeleteBucketPolicyOutput>> {
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<S3Response<GetBucketLifecycleConfigurationOutput>> {
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<S3Response<DeleteBucketLifecycleOutput>> {
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<S3Response<GetBucketEncryptionOutput>> {
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<S3Response<DeleteBucketEncryptionOutput>> {
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<S3Response<GetBucketReplicationOutput>> {
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<S3Response<DeleteBucketReplicationOutput>> {
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<S3Response<GetBucketNotificationConfigurationOutput>> {
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<GetBucketAclInput>) -> S3Result<S3Response<GetBucketAclOutput>> {
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<GetObjectAclInput>) -> S3Result<S3Response<GetObjectAclOutput>> {
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)));