diff --git a/ecstore/src/bucket/metadata.rs b/ecstore/src/bucket/metadata.rs index 6d53e1b4d..8f349a983 100644 --- a/ecstore/src/bucket/metadata.rs +++ b/ecstore/src/bucket/metadata.rs @@ -2,6 +2,7 @@ use super::{ encryption::BucketSSEConfig, event, lifecycle::lifecycle::Lifecycle, objectlock, policy::bucket_policy::BucketPolicy, quota::BucketQuota, replication, tags::Tags, target::BucketTargets, versioning::Versioning, }; + use byteorder::{BigEndian, ByteOrder, LittleEndian}; use rmp_serde::Serializer as rmpSerializer; use serde::Serializer; @@ -10,7 +11,7 @@ use std::collections::HashMap; use std::fmt::Display; use std::str::FromStr; use time::OffsetDateTime; -use tracing::error; +use tracing::{error, warn}; use crate::bucket::tags; use crate::config; @@ -172,7 +173,19 @@ impl BucketMetadata { return Err(Error::msg("read_bucket_metadata: data invalid")); } - // TODO: check version + let format = LittleEndian::read_u16(&buf[0..2]); + let version = LittleEndian::read_u16(&buf[2..4]); + + match format { + BUCKET_METADATA_FORMAT => {} + _ => return Err(Error::msg("read_bucket_metadata: format invalid")), + } + + match version { + BUCKET_METADATA_VERSION => {} + _ => return Err(Error::msg("read_bucket_metadata: version invalid")), + } + Ok(()) } @@ -252,6 +265,7 @@ impl BucketMetadata { fn parse_all_configs(&mut self, _api: &ECStore) -> Result<()> { if !self.tagging_config_xml.is_empty() { + warn!("self.tagging_config_xml {:?}", &self.tagging_config_xml); self.tagging_config = Some(tags::Tags::unmarshal(&self.tagging_config_xml)?); } @@ -267,6 +281,7 @@ pub async fn load_bucket_metadata_parse(api: &ECStore, bucket: &str, parse: bool let mut bm = match read_bucket_metadata(api, bucket).await { Ok(res) => res, Err(err) => { + warn!("load_bucket_metadata_parse err {:?}", &err); if !config::error::is_not_found(&err) { return Err(err); } @@ -299,7 +314,9 @@ async fn read_bucket_metadata(api: &ECStore, bucket: &str) -> Result(deserializer: D) -> core::result::Result diff --git a/ecstore/src/bucket/metadata_sys.rs b/ecstore/src/bucket/metadata_sys.rs index a43eda5df..05a60f8f4 100644 --- a/ecstore/src/bucket/metadata_sys.rs +++ b/ecstore/src/bucket/metadata_sys.rs @@ -1,3 +1,4 @@ +use std::collections::HashSet; use std::{collections::HashMap, sync::Arc}; use crate::bucket::error::BucketMetadataError; @@ -7,21 +8,26 @@ use crate::config; use crate::config::error::ConfigError; use crate::disk::error::DiskError; use crate::error::{Error, Result}; -use crate::global::{is_dist_erasure, is_erasure, new_object_layer_fn}; +use crate::global::{is_dist_erasure, is_erasure, new_object_layer_fn, GLOBAL_Endpoints}; use crate::store::ECStore; use futures::future::join_all; use lazy_static::lazy_static; use time::OffsetDateTime; use tokio::sync::RwLock; +use tracing::{error, info, warn}; use super::metadata::{load_bucket_metadata, BucketMetadata}; use super::tags; lazy_static! { - static ref GLOBAL_BucketMetadataSys: Arc = Arc::new(BucketMetadataSys::new()); + static ref GLOBAL_BucketMetadataSys: Arc> = Arc::new(RwLock::new(BucketMetadataSys::new())); } -pub fn get_bucket_metadata_sys() -> Arc { +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 get_bucket_metadata_sys() -> Arc> { GLOBAL_BucketMetadataSys.clone() } @@ -37,46 +43,77 @@ impl BucketMetadataSys { Self::default() } - pub async fn init(&mut self, api: Option, buckets: Vec<&str>) -> Result<()> { - if api.is_none() { - return Err(Error::msg("errServerNotInitialized")); - } - self.api = api; + pub async fn init(&mut self, api: ECStore, buckets: Vec) { + // if api.is_none() { + // return Err(Error::msg("errServerNotInitialized")); + // } + self.api = Some(api); let _ = self.init_internal(buckets).await; + } + async fn init_internal(&self, buckets: Vec) -> Result<()> { + let count = { + let endpoints = GLOBAL_Endpoints.read().await; + endpoints.es_count() * 10 + }; + + let mut failed_buckets: HashSet = HashSet::new(); + let mut buckets = buckets.as_slice(); + + loop { + if buckets.len() < count { + self.concurrent_load(buckets, &mut failed_buckets).await; + break; + } + + self.concurrent_load(&buckets[..count], &mut failed_buckets).await; + + buckets = &buckets[count..] + } + + let mut initialized = self.initialized.write().await; + *initialized = true; + + if is_dist_erasure().await { + // TODO: refresh_buckets_metadata_loop + } Ok(()) } - async fn init_internal(&self, buckets: Vec<&str>) -> Result<()> { - if self.api.is_none() { - return Err(Error::msg("errServerNotInitialized")); - } - let mut futures = Vec::new(); - let mut errs = Vec::new(); - let mut ress = Vec::new(); - for &bucket in buckets.iter() { - futures.push(load_bucket_metadata(self.api.as_ref().unwrap(), bucket)); + async fn concurrent_load(&self, buckets: &[String], failed_buckets: &mut HashSet) { + let mut futures = Vec::new(); + + for bucket in buckets.iter() { + futures.push(load_bucket_metadata(self.api.as_ref().unwrap(), bucket.as_str())); } let results = join_all(futures).await; + let mut idx = 0; + + let mut mp = self.metadata_map.write().await; + for res in results { match res { - Ok(entrys) => { - ress.push(Some(entrys)); - errs.push(None); + Ok(res) => { + if let Some(bucket) = buckets.get(idx) { + mp.insert(bucket.clone(), res); + } } Err(e) => { - ress.push(None); - errs.push(Some(e)); + error!("Unable to load bucket metadata, will be retried: {:?}", e); + if let Some(bucket) = buckets.get(idx) { + failed_buckets.insert(bucket.clone()); + } } } + + idx += 1; } - unimplemented!() } - async fn concurrent_load(&self, buckets: Vec<&str>) -> Result> { + async fn refresh_buckets_metadata_loop(&self, failed_buckets: &HashSet) -> Result<()> { unimplemented!() } @@ -105,7 +142,7 @@ impl BucketMetadataSys { map.clear(); } - pub async fn update(&mut self, bucket: &str, config_file: &str, data: Vec) -> Result<(OffsetDateTime)> { + pub async fn update(&mut self, bucket: &str, config_file: &str, data: Vec) -> Result { self.update_and_parse(bucket, config_file, data, true).await } diff --git a/ecstore/src/bucket/mod.rs b/ecstore/src/bucket/mod.rs index b2df60cc2..9be464946 100644 --- a/ecstore/src/bucket/mod.rs +++ b/ecstore/src/bucket/mod.rs @@ -10,7 +10,7 @@ mod quota; mod replication; mod tags; mod target; -mod versioning; pub mod utils; +mod versioning; -pub use metadata_sys::get_bucket_metadata_sys; +pub use metadata_sys::{get_bucket_metadata_sys, init_bucket_metadata_sys}; diff --git a/ecstore/src/config/common.rs b/ecstore/src/config/common.rs index e7132ed65..6dedfbed8 100644 --- a/ecstore/src/config/common.rs +++ b/ecstore/src/config/common.rs @@ -47,7 +47,7 @@ async fn save_config_with_opts(api: &ECStore, file: &str, data: &[u8], opts: &Ob RUSTFS_META_BUCKET, file, PutObjReader::new(StreamingBlob::from(Body::from(data.to_vec())), data.len()), - &ObjectOptions::default(), + opts, ) .await?; Ok(()) diff --git a/ecstore/src/endpoints.rs b/ecstore/src/endpoints.rs index ad533140b..14547bbfd 100644 --- a/ecstore/src/endpoints.rs +++ b/ecstore/src/endpoints.rs @@ -388,7 +388,7 @@ pub struct PoolEndpoints { /// list of list of endpoints #[derive(Debug, Clone)] -pub struct EndpointServerPools(Vec); +pub struct EndpointServerPools(pub Vec); impl From> for EndpointServerPools { fn from(v: Vec) -> Self { @@ -409,6 +409,9 @@ impl AsMut> for EndpointServerPools { } impl EndpointServerPools { + pub fn reset(&mut self, eps: Vec) { + self.0 = eps; + } pub fn from_volumes(server_addr: &str, endpoints: Vec) -> Result<(EndpointServerPools, SetupType)> { let layouts = DisksLayout::try_from(endpoints.as_slice())?; diff --git a/ecstore/src/global.rs b/ecstore/src/global.rs index c7ac7b00b..1b24924c0 100644 --- a/ecstore/src/global.rs +++ b/ecstore/src/global.rs @@ -1,17 +1,27 @@ -use crate::error::Result; use lazy_static::lazy_static; use std::{collections::HashMap, sync::Arc}; -use tokio::{fs, sync::RwLock}; +use tokio::sync::RwLock; use crate::{ - disk::{new_disk, DiskOption, DiskStore}, - endpoints::{EndpointServerPools, SetupType}, + disk::DiskStore, + endpoints::{EndpointServerPools, PoolEndpoints, SetupType}, store::ECStore, }; lazy_static! { pub static ref GLOBAL_OBJECT_API: Arc>> = Arc::new(RwLock::new(None)); pub static ref GLOBAL_LOCAL_DISK: Arc>>> = Arc::new(RwLock::new(Vec::new())); + static ref GLOBAL_IsErasure: RwLock = RwLock::new(false); + static ref GLOBAL_IsDistErasure: RwLock = RwLock::new(false); + static ref GLOBAL_IsErasureSD: RwLock = RwLock::new(false); + pub static ref GLOBAL_LOCAL_DISK_MAP: Arc>>> = Arc::new(RwLock::new(HashMap::new())); + pub static ref GLOBAL_LOCAL_DISK_SET_DRIVES: Arc> = Arc::new(RwLock::new(Vec::new())); + pub static ref GLOBAL_Endpoints: RwLock = RwLock::new(EndpointServerPools(Vec::new())); +} + +pub async fn set_global_endpoints(eps: Vec) { + let mut endpoints = GLOBAL_Endpoints.write().await; + endpoints.reset(eps); } pub fn new_object_layer_fn() -> Arc>> { @@ -23,12 +33,6 @@ pub async fn set_object_layer(o: ECStore) { *global_object_api = Some(o); } -lazy_static! { - static ref GLOBAL_IsErasure: RwLock = RwLock::new(false); - static ref GLOBAL_IsDistErasure: RwLock = RwLock::new(false); - static ref GLOBAL_IsErasureSD: RwLock = RwLock::new(false); -} - pub async fn is_dist_erasure() -> bool { let lock = GLOBAL_IsDistErasure.read().await; *lock == true @@ -55,8 +59,3 @@ pub async fn update_erasure_type(setup_type: SetupType) { } type TypeLocalDiskSetDrives = Vec>>>; - -lazy_static! { - pub static ref GLOBAL_LOCAL_DISK_MAP: Arc>>> = Arc::new(RwLock::new(HashMap::new())); - pub static ref GLOBAL_LOCAL_DISK_SET_DRIVES: Arc> = Arc::new(RwLock::new(Vec::new())); -} diff --git a/ecstore/src/lib.rs b/ecstore/src/lib.rs index 4210db4c7..bd8bced90 100644 --- a/ecstore/src/lib.rs +++ b/ecstore/src/lib.rs @@ -21,4 +21,5 @@ mod utils; pub mod bucket; pub use global::new_object_layer_fn; +pub use global::set_global_endpoints; pub use global::update_erasure_type; diff --git a/ecstore/src/peer.rs b/ecstore/src/peer.rs index 27256de8a..a1b6deb04 100644 --- a/ecstore/src/peer.rs +++ b/ecstore/src/peer.rs @@ -26,7 +26,7 @@ pub trait PeerS3Client: Debug + Sync + Send + 'static { fn get_pools(&self) -> Option>; } -#[derive(Debug)] +#[derive(Debug, Clone)] pub struct S3PeerSys { pub clients: Vec, pub pools_count: usize, diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 0e563adae..772d87e94 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -34,7 +34,7 @@ use tokio::sync::Semaphore; use tracing::{debug, info}; use uuid::Uuid; -#[derive(Debug)] +#[derive(Debug, Clone)] pub struct ECStore { pub id: uuid::Uuid, // pub disks: Vec, @@ -46,7 +46,7 @@ pub struct ECStore { 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::try_from(endpoints.as_slice())?; let mut deployment_id = None; @@ -145,9 +145,9 @@ impl ECStore { peer_sys, }; - set_object_layer(ec).await; + set_object_layer(ec.clone()).await; - Ok(()) + Ok(ec) } pub fn init_local_disks() {} @@ -507,28 +507,29 @@ impl StorageAPI for ECStore { // TODO: delete created bucket when error self.peer_sys.make_bucket(bucket, opts).await?; - let meta = BucketMetadata::new(bucket); - let data = meta.marshal_msg()?; - let file_path = meta.save_file_path(); + let mut meta = BucketMetadata::new(bucket); + meta.save(self).await?; + // let data = meta.marshal_msg()?; + // let file_path = meta.save_file_path(); - // TODO: wrap hash reader + // // TODO: wrap hash reader - let content_len = data.len(); + // let content_len = data.len(); - let body = Body::from(data); + // let body = Body::from(data); - let reader = PutObjReader::new(StreamingBlob::from(body), content_len); + // let reader = PutObjReader::new(StreamingBlob::from(body), content_len); - self.put_object( - RUSTFS_META_BUCKET, - &file_path, - reader, - &ObjectOptions { - max_parity: true, - ..Default::default() - }, - ) - .await?; + // self.put_object( + // RUSTFS_META_BUCKET, + // &file_path, + // reader, + // &ObjectOptions { + // max_parity: true, + // ..Default::default() + // }, + // ) + // .await?; // TODO: toObjectErr diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index d05bfea6b..8289ec22a 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -414,7 +414,11 @@ pub struct ObjectOptions { // } #[derive(Debug, Default, Serialize, Deserialize)] -pub struct BucketOptions {} +pub struct BucketOptions { + pub deleted: bool, // true only when site replication is enabled + pub cached: bool, // true only when we are requesting a cached response instead of hitting the disk for example ListBuckets() call. + pub no_metadata: bool, +} #[derive(Debug, Clone, Serialize, Deserialize)] pub struct BucketInfo { diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index 802832a83..c23e04a6f 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -6,8 +6,11 @@ mod storage; use clap::Parser; use common::error::{Error, Result}; use ecstore::{ + bucket::init_bucket_metadata_sys, endpoints::EndpointServerPools, + set_global_endpoints, store::{init_local_disks, ECStore}, + store_api::{BucketOptions, StorageAPI}, update_erasure_type, }; use grpc::make_server; @@ -90,8 +93,9 @@ async fn run(opt: config::Opt) -> Result<()> { // 用于rpc let (endpoint_pools, setup_type) = EndpointServerPools::from_volumes(opt.address.clone().as_str(), opt.volumes.clone()) .map_err(|err| Error::from_string(err.to_string()))?; - + set_global_endpoints(endpoint_pools.as_ref().clone()).await; update_erasure_type(setup_type).await; + // 初始化本地磁盘 init_local_disks(endpoint_pools.clone()) .await @@ -179,11 +183,23 @@ async fn run(opt: config::Opt) -> Result<()> { }); // init store - ECStore::new(opt.address.clone(), endpoint_pools.clone()) + let store = ECStore::new(opt.address.clone(), endpoint_pools.clone()) .await .map_err(|err| Error::from_string(err.to_string()))?; info!(" init store success!"); + let buckets_list = store + .list_bucket(&BucketOptions { + no_metadata: true, + ..Default::default() + }) + .await + .map_err(|err| Error::from_string(err.to_string()))?; + + let buckets = buckets_list.iter().map(|v| v.name.clone()).collect(); + + init_bucket_metadata_sys(store.clone(), buckets).await; + tokio::select! { _ = tokio::signal::ctrl_c() => { diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 2c9a0325d..b701404bb 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -250,7 +250,7 @@ impl S3 for FS { None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), }; - if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions {}).await { + if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions::default()).await { if DiskError::VolumeNotFound.is(&e) { return Err(s3_error!(NoSuchBucket)); } else { @@ -324,7 +324,7 @@ impl S3 for FS { None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), }; - if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions {}).await { + if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions::default()).await { if DiskError::VolumeNotFound.is(&e) { return Err(s3_error!(NoSuchBucket)); } else { @@ -375,7 +375,7 @@ impl S3 for FS { None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), }; - let bucket_infos = try_!(store.list_bucket(&BucketOptions {}).await); + let bucket_infos = try_!(store.list_bucket(&BucketOptions::default()).await); let buckets: Vec = bucket_infos .iter()