From 476d93d0760e38c4c48ce82fb3f73851bb4da469 Mon Sep 17 00:00:00 2001 From: weisd Date: Wed, 6 Nov 2024 17:24:52 +0800 Subject: [PATCH 1/4] init storage info api --- ecstore/src/sets.rs | 28 +++++++++++++-- ecstore/src/store.rs | 71 ++++++++++++++++++++++++++++++++++++- ecstore/src/store_api.rs | 76 ++++++++++++++++++++++++++++++++++++++-- 3 files changed, 169 insertions(+), 6 deletions(-) diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index c5dc2c07c..384983441 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -22,9 +22,9 @@ use crate::{ }, set_disk::SetDisks, store_api::{ - BucketInfo, BucketOptions, CompletePart, DeleteBucketOptions, DeletedObject, GetObjectReader, HTTPRangeSpec, + BackendInfo, BucketInfo, BucketOptions, CompletePart, DeleteBucketOptions, DeletedObject, GetObjectReader, HTTPRangeSpec, ListMultipartsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartUploadResult, ObjectIO, ObjectInfo, ObjectOptions, - ObjectToDelete, PartInfo, PutObjReader, StorageAPI, + ObjectToDelete, PartInfo, PutObjReader, StorageAPI, StorageInfo, }, utils::hash, }; @@ -47,6 +47,7 @@ pub struct Sets { pub partiy_count: usize, pub set_count: usize, pub set_drive_count: usize, + pub default_parity_count: usize, pub distribution_algo: DistributionAlgoVersion, ctx: CancellationToken, } @@ -152,6 +153,7 @@ impl Sets { partiy_count, set_count, set_drive_count, + default_parity_count: partiy_count, distribution_algo: fm.erasure.distribution_algo.clone(), ctx: CancellationToken::new(), }); @@ -286,6 +288,28 @@ impl ObjectIO for Sets { #[async_trait::async_trait] impl StorageAPI for Sets { + async fn backend_info(&self) -> BackendInfo { + unimplemented!() + } + async fn storage_info(&self) -> StorageInfo { + unimplemented!() + } + async fn local_storage_info(&self) -> StorageInfo { + let mut futures = Vec::with_capacity(self.disk_set.len()); + + for set in self.disk_set.iter() { + futures.push(set.local_storage_info()) + } + + let results = join_all(futures).await; + + let mut disks = Vec::new(); + + for res in results.into_iter() { + disks.extend_from_slice(&res.disks); + } + StorageInfo { backend: None, disks } + } async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { unimplemented!() } diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 98627c709..7e0e1773c 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -3,6 +3,7 @@ use crate::bucket::metadata; use crate::bucket::metadata_sys::{self, init_bucket_metadata_sys, set_bucket_metadata}; use crate::bucket::utils::{check_valid_bucket_name, check_valid_bucket_name_strict, is_meta_bucketname}; +use crate::config::GLOBAL_StorageClass; use crate::config::{self, storageclass, GLOBAL_ConfigSys}; use crate::disk::endpoint::EndpointType; use crate::disk::{DiskAPI, DiskInfo, DiskInfoOptions, MetaCacheEntry}; @@ -13,7 +14,7 @@ use crate::global::{ use crate::heal::heal_commands::{HealOpts, HealResultItem, HealScanMode}; use crate::heal::heal_ops::HealObjectFn; use crate::new_object_layer_fn; -use crate::store_api::{ListMultipartsInfo, ObjectIO}; +use crate::store_api::{BackendByte, BackendDisks, BackendInfo, ListMultipartsInfo, ObjectIO, StorageInfo}; use crate::store_err::{ is_err_bucket_exists, is_err_invalid_upload_id, is_err_object_not_found, is_err_version_not_found, StorageError, }; @@ -791,6 +792,74 @@ lazy_static! { #[async_trait::async_trait] impl StorageAPI for ECStore { + async fn backend_info(&self) -> BackendInfo { + let (standard_scparities, rrscparities) = { + if let Some(sc) = GLOBAL_StorageClass.get() { + let sc_parity = sc + .get_parity_for_sc(storageclass::CLASS_STANDARD) + .or(Some(self.pools[0].default_parity_count)); + + let rrs_sc_parity = sc.get_parity_for_sc(storageclass::RRS); + + (sc_parity, rrs_sc_parity) + } else { + (Some(self.pools[0].default_parity_count), None) + } + }; + + let mut standard_scdata = Vec::new(); + let mut rrscdata = Vec::new(); + let mut drives_per_set = Vec::new(); + let mut total_sets = Vec::new(); + + for (idx, set_count) in self.set_drive_counts().iter().enumerate() { + if let Some(sc_parity) = standard_scparities { + standard_scdata.push(set_count - sc_parity); + } + if let Some(rr_sc_parity) = rrscparities { + rrscdata.push(set_count - rr_sc_parity); + } + total_sets.push(self.pools[idx].set_count); + drives_per_set.push(*set_count); + } + + BackendInfo { + backend_type: BackendByte::Erasure, + online_disks: BackendDisks::new(), + offline_disks: BackendDisks::new(), + standard_scdata, + standard_scparities, + rrscdata, + rrscparities, + total_sets, + drives_per_set, + } + } + async fn storage_info(&self) -> StorageInfo { + unimplemented!() + } + async fn local_storage_info(&self) -> StorageInfo { + let mut futures = Vec::with_capacity(self.pools.len()); + + for pool in self.pools.iter() { + futures.push(pool.local_storage_info()) + } + + let results = join_all(futures).await; + + let mut disks = Vec::new(); + + for res in results.into_iter() { + disks.extend_from_slice(&res.disks); + } + + let backend = self.backend_info().await; + StorageInfo { + backend: Some(backend), + disks, + } + } + async fn list_bucket(&self, opts: &BucketOptions) -> Result> { // TODO: opts.cached diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index 21df384ed..9cc6517a5 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -778,6 +778,75 @@ pub struct DeletedObject { // pub replication_state: ReplicationState, } +#[derive(Debug, Default)] +pub enum BackendByte { + #[default] + Unknown, + FS, + Erasure, +} + +#[derive(Serialize, Deserialize, Debug, Default, Clone)] +pub struct StorageDisk { + pub endpoint: String, + pub root_disk: bool, + pub drive_path: String, + pub healing: bool, + pub scanning: bool, + pub state: String, + pub uuid: String, + pub major: u32, + pub minor: u32, + pub model: Option, + pub total_space: u64, + pub used_space: u64, + pub available_space: u64, + pub read_throughput: f64, + pub write_throughput: f64, + pub read_latency: f64, + pub write_latency: f64, + pub utilization: f64, + // pub metrics: Option, + // pub heal_info: Option, + pub used_inodes: u64, + pub free_inodes: u64, + pub local: bool, + pub pool_index: Option, + pub set_index: Option, + pub disk_index: Option, +} + +#[derive(Debug, Default)] +pub struct StorageInfo { + pub disks: Vec, + pub backend: Option, +} + +#[derive(Debug, Default)] +pub struct BackendDisks(HashMap); + +impl BackendDisks { + pub fn new() -> Self { + Self(HashMap::new()) + } + pub fn sum(&self) -> usize { + self.0.values().sum() + } +} + +#[derive(Debug, Default)] +pub struct BackendInfo { + pub backend_type: BackendByte, + pub online_disks: BackendDisks, + pub offline_disks: BackendDisks, + pub standard_scdata: Vec, + pub standard_scparities: Option, + pub rrscdata: Vec, + pub rrscparities: Option, + pub total_sets: Vec, + pub drives_per_set: Vec, +} + #[async_trait::async_trait] pub trait ObjectIO: Send + Sync + 'static { // GetObjectNInfo @@ -798,9 +867,10 @@ pub trait StorageAPI: ObjectIO { // NewNSLock // Shutdown // NSScanner - // BackendInfo - // StorageInfo - // LocalStorageInfo + + async fn backend_info(&self) -> BackendInfo; + async fn storage_info(&self) -> StorageInfo; + async fn local_storage_info(&self) -> StorageInfo; async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()>; async fn get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result; From 91c25922130198b3fb1478b1088db79c392c29e9 Mon Sep 17 00:00:00 2001 From: weisd Date: Thu, 7 Nov 2024 17:15:23 +0800 Subject: [PATCH 2/4] change pool idx to i32 --- ecstore/src/disk/endpoint.rs | 42 +++---- ecstore/src/disk/local.rs | 35 ++++-- ecstore/src/disk/remote.rs | 24 +++- ecstore/src/endpoints.rs | 210 +++++++++++++++++------------------ ecstore/src/lib.rs | 1 + ecstore/src/pools.rs | 91 +++++++++++++++ ecstore/src/sets.rs | 27 ++++- ecstore/src/store.rs | 16 +-- ecstore/src/store_api.rs | 8 +- scripts/run.sh | 4 +- 10 files changed, 300 insertions(+), 158 deletions(-) create mode 100644 ecstore/src/pools.rs diff --git a/ecstore/src/disk/endpoint.rs b/ecstore/src/disk/endpoint.rs index 6a6ba743f..254a1213b 100644 --- a/ecstore/src/disk/endpoint.rs +++ b/ecstore/src/disk/endpoint.rs @@ -21,9 +21,9 @@ pub struct Endpoint { pub url: url::Url, pub is_local: bool, - pub pool_idx: Option, - pub set_idx: Option, - pub disk_idx: Option, + pub pool_idx: i32, + pub set_idx: i32, + pub disk_idx: i32, } impl Display for Endpoint { @@ -122,9 +122,9 @@ impl TryFrom<&str> for Endpoint { Ok(Endpoint { url, is_local, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: -1, + set_idx: -1, + disk_idx: -1, }) } } @@ -141,17 +141,17 @@ impl Endpoint { /// sets a specific pool number to this node pub fn set_pool_index(&mut self, idx: usize) { - self.pool_idx = Some(idx) + self.pool_idx = idx as i32 } /// sets a specific set number to this node pub fn set_set_index(&mut self, idx: usize) { - self.set_idx = Some(idx) + self.set_idx = idx as i32 } /// sets a specific disk number to this node pub fn set_disk_index(&mut self, idx: usize) { - self.disk_idx = Some(idx) + self.disk_idx = idx as i32 } /// resolves the host and updates if it is local or not. @@ -231,9 +231,9 @@ mod test { expected_endpoint: Some(Endpoint { url: root_slash_foo, is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: -1, + set_idx: -1, + disk_idx: -1, }), expected_type: Some(EndpointType::Path), expected_err: None, @@ -243,9 +243,9 @@ mod test { expected_endpoint: Some(Endpoint { url: u2, is_local: false, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: -1, + set_idx: -1, + disk_idx: -1, }), expected_type: Some(EndpointType::Url), expected_err: None, @@ -255,9 +255,9 @@ mod test { expected_endpoint: Some(Endpoint { url: u4, is_local: false, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: -1, + set_idx: -1, + disk_idx: -1, }), expected_type: Some(EndpointType::Url), expected_err: None, @@ -315,9 +315,9 @@ mod test { expected_endpoint: Some(Endpoint { url: u6, is_local: false, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: -1, + set_idx: -1, + disk_idx: -1, }), expected_type: Some(EndpointType::Url), expected_err: None, diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 4d5b1faf9..d77f1322d 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -122,7 +122,7 @@ impl LocalDisk { let fm = FormatV3::try_from(s)?; let (set_idx, disk_idx) = fm.find_disk_index_by_disk_id(fm.erasure.this)?; - if Some(set_idx) != ep.set_idx || Some(disk_idx) != ep.disk_idx { + if set_idx as i32 != ep.set_idx || disk_idx as i32 != ep.disk_idx { return Err(Error::from(DiskError::InconsistentDisk)); } @@ -768,9 +768,27 @@ impl DiskAPI for LocalDisk { fn get_disk_location(&self) -> DiskLocation { DiskLocation { - pool_idx: self.endpoint.pool_idx, - set_idx: self.endpoint.set_idx, - disk_idx: self.endpoint.pool_idx, + pool_idx: { + if self.endpoint.pool_idx < 0 { + None + } else { + Some(self.endpoint.pool_idx as usize) + } + }, + set_idx: { + if self.endpoint.set_idx < 0 { + None + } else { + Some(self.endpoint.set_idx as usize) + } + }, + disk_idx: { + if self.endpoint.disk_idx < 0 { + None + } else { + Some(self.endpoint.disk_idx as usize) + } + }, } } @@ -811,13 +829,8 @@ impl DiskAPI for LocalDisk { let disk_id = fm.erasure.this; - match (self.endpoint.set_idx, self.endpoint.disk_idx) { - (Some(set_idx), Some(disk_idx)) => { - if m != set_idx || n != disk_idx { - return Err(Error::new(DiskError::InconsistentDisk)); - } - } - _ => return Err(Error::new(DiskError::InconsistentDisk)), + if m as i32 != self.endpoint.set_idx || n as i32 != self.endpoint.disk_idx { + return Err(Error::new(DiskError::InconsistentDisk)); } format_info.id = Some(disk_id); diff --git a/ecstore/src/disk/remote.rs b/ecstore/src/disk/remote.rs index 0eefe1c58..80fdbb5b7 100644 --- a/ecstore/src/disk/remote.rs +++ b/ecstore/src/disk/remote.rs @@ -83,9 +83,27 @@ impl DiskAPI for RemoteDisk { fn get_disk_location(&self) -> DiskLocation { DiskLocation { - pool_idx: self.endpoint.pool_idx, - set_idx: self.endpoint.set_idx, - disk_idx: self.endpoint.pool_idx, + pool_idx: { + if self.endpoint.pool_idx < 0 { + None + } else { + Some(self.endpoint.pool_idx as usize) + } + }, + set_idx: { + if self.endpoint.set_idx < 0 { + None + } else { + Some(self.endpoint.set_idx as usize) + } + }, + disk_idx: { + if self.endpoint.disk_idx < 0 { + None + } else { + Some(self.endpoint.disk_idx as usize) + } + }, } } diff --git a/ecstore/src/endpoints.rs b/ecstore/src/endpoints.rs index c5ff9cbf0..844fe936c 100644 --- a/ecstore/src/endpoints.rs +++ b/ecstore/src/endpoints.rs @@ -488,10 +488,8 @@ impl EndpointServerPools { grid_host: ep.grid_host(), }); - if let Some(pool_idx) = ep.pool_idx { - if !n.pools.contains(&pool_idx) { - n.pools.push(pool_idx); - } + if !n.pools.contains(&(ep.pool_idx as usize)) { + n.pools.push(ep.pool_idx as usize); } } } @@ -730,9 +728,9 @@ mod test { expected_endpoints: Some(Endpoints(vec![Endpoint { url: must_file_path("/d1"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }])), expected_setup_type: Some(SetupType::ErasureSD), ..Default::default() @@ -744,9 +742,9 @@ mod test { expected_endpoints: Some(Endpoints(vec![Endpoint { url: must_file_path("/d1"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }])), expected_setup_type: Some(SetupType::ErasureSD), ..Default::default() @@ -772,30 +770,30 @@ mod test { Endpoint { url: must_file_path("/d1"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: must_file_path("/d2"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: must_file_path("/d3"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: must_file_path("/d4"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, ])), expected_setup_type: Some(SetupType::Erasure), @@ -815,30 +813,30 @@ mod test { Endpoint { url: must_url("http://localhost:9000/d1"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: must_url("http://localhost:9000/d2"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: must_url("http://localhost:9000/d3"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: must_url("http://localhost:9000/d4"), is_local: true, - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, ])), expected_setup_type: Some(SetupType::Erasure), @@ -897,30 +895,30 @@ mod test { Endpoint { url: case1_ur_ls[0].clone(), is_local: case1_local_flags[0], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case1_ur_ls[1].clone(), is_local: case1_local_flags[1], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case1_ur_ls[2].clone(), is_local: case1_local_flags[2], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case1_ur_ls[3].clone(), is_local: case1_local_flags[3], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, ])), expected_setup_type: Some(SetupType::DistErasure), @@ -939,30 +937,30 @@ mod test { Endpoint { url: case2_ur_ls[0].clone(), is_local: case2_local_flags[0], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case2_ur_ls[1].clone(), is_local: case2_local_flags[1], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case2_ur_ls[2].clone(), is_local: case2_local_flags[2], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case2_ur_ls[3].clone(), is_local: case2_local_flags[3], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, ])), expected_setup_type: Some(SetupType::DistErasure), @@ -981,30 +979,30 @@ mod test { Endpoint { url: case3_ur_ls[0].clone(), is_local: case3_local_flags[0], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case3_ur_ls[1].clone(), is_local: case3_local_flags[1], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case3_ur_ls[2].clone(), is_local: case3_local_flags[2], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case3_ur_ls[3].clone(), is_local: case3_local_flags[3], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, ])), expected_setup_type: Some(SetupType::DistErasure), @@ -1023,30 +1021,30 @@ mod test { Endpoint { url: case4_ur_ls[0].clone(), is_local: case4_local_flags[0], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case4_ur_ls[1].clone(), is_local: case4_local_flags[1], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case4_ur_ls[2].clone(), is_local: case4_local_flags[2], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case4_ur_ls[3].clone(), is_local: case4_local_flags[3], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, ])), expected_setup_type: Some(SetupType::DistErasure), @@ -1065,30 +1063,30 @@ mod test { Endpoint { url: case5_ur_ls[0].clone(), is_local: case5_local_flags[0], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case5_ur_ls[1].clone(), is_local: case5_local_flags[1], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case5_ur_ls[2].clone(), is_local: case5_local_flags[2], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case5_ur_ls[3].clone(), is_local: case5_local_flags[3], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, ])), expected_setup_type: Some(SetupType::DistErasure), @@ -1107,30 +1105,30 @@ mod test { Endpoint { url: case6_ur_ls[0].clone(), is_local: case6_local_flags[0], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case6_ur_ls[1].clone(), is_local: case6_local_flags[1], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case6_ur_ls[2].clone(), is_local: case6_local_flags[2], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, Endpoint { url: case6_ur_ls[3].clone(), is_local: case6_local_flags[3], - pool_idx: None, - set_idx: None, - disk_idx: None, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, }, ])), expected_setup_type: Some(SetupType::DistErasure), diff --git a/ecstore/src/lib.rs b/ecstore/src/lib.rs index 2f53377e3..95bf4451b 100644 --- a/ecstore/src/lib.rs +++ b/ecstore/src/lib.rs @@ -22,6 +22,7 @@ mod utils; pub mod bucket; pub mod file_meta_inline; pub mod options; +pub mod pools; pub(crate) mod store_err; pub mod xhttp; diff --git a/ecstore/src/pools.rs b/ecstore/src/pools.rs new file mode 100644 index 000000000..e8d7f49ae --- /dev/null +++ b/ecstore/src/pools.rs @@ -0,0 +1,91 @@ +use crate::error::{Error, Result}; +use crate::store_api::{StorageAPI, StorageDisk, StorageInfo}; +use crate::{sets::Sets, store::ECStore}; +use std::sync::Arc; +use time::OffsetDateTime; + +#[derive(Debug, Clone)] +pub struct PoolStatus { + pub id: usize, + pub cmd_line: String, + pub last_update: OffsetDateTime, + pub decommission: Option, +} + +#[derive(Debug, Clone)] +pub struct PoolMeta { + pub pools: Vec, + pub dont_save: bool, +} + +impl PoolMeta { + pub fn new(pools: Vec>) -> Self { + let mut status = Vec::with_capacity(pools.len()); + for (idx, pool) in pools.iter().enumerate() { + status.push(PoolStatus { + id: idx, + cmd_line: pool.endpoints.cmd_line.clone(), + last_update: OffsetDateTime::now_utc(), + decommission: None, + }); + } + + Self { + pools: status, + dont_save: false, + } + } +} + +#[derive(Debug, Clone)] +pub struct PoolDecommissionInfo { + pub start_time: OffsetDateTime, + pub start_size: usize, + pub total_size: usize, + pub current_size: usize, + pub complete: bool, + pub failed: bool, + pub canceled: bool, + pub queued_buckets: Vec, + pub decommissioned_buckets: Vec, + pub bucket: String, + pub prefix: String, + pub object: String, + + pub items_decommissioned: usize, + pub items_decommission_failed: usize, + pub bytes_done: usize, + pub bytes_failed: usize, +} + +struct PoolSpaceInfo { + pub free: usize, + pub total: usize, + pub used: usize, +} + +impl ECStore { + pub fn status(&self, idx: usize) -> Result { + unimplemented!() + } + + async fn get_decommission_pool_space_info(&self, idx: usize) -> Result { + if let Some(sets) = self.pools.get(idx) { + let mut info = sets.storage_info().await; + info.backend = self.backend_info().await; + + unimplemented!() + } else { + Err(Error::msg("InvalidArgument")) + } + } +} + +fn get_total_usable_capacity(disks: &Vec, info: &StorageInfo) -> usize { + for disk in disks.iter() { + // if disk.pool_index < 0 || info.backend.standard_scdata.len() <= disk.pool_index { + // continue; + // } + } + unimplemented!() +} diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index 384983441..203e3a3ad 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -90,7 +90,8 @@ impl Sets { let idx = i * set_drive_count + j; let mut disk = disks[idx].clone(); - let endpoint = endpoints.endpoints.as_ref().get(idx).cloned(); + let endpoint = endpoints.endpoints.as_ref()[idx].clone(); + // let endpoint = endpoints.endpoints.as_ref().get(idx).cloned(); set_endpoints.push(endpoint); if disk.is_none() { @@ -292,7 +293,24 @@ impl StorageAPI for Sets { unimplemented!() } async fn storage_info(&self) -> StorageInfo { - unimplemented!() + let mut futures = Vec::with_capacity(self.disk_set.len()); + + for set in self.disk_set.iter() { + futures.push(set.storage_info()) + } + + let results = join_all(futures).await; + + let mut disks = Vec::new(); + + for res in results.into_iter() { + disks.extend_from_slice(&res.disks); + } + + StorageInfo { + disks, + ..Default::default() + } } async fn local_storage_info(&self) -> StorageInfo { let mut futures = Vec::with_capacity(self.disk_set.len()); @@ -308,7 +326,10 @@ impl StorageAPI for Sets { for res in results.into_iter() { disks.extend_from_slice(&res.disks); } - StorageInfo { backend: None, disks } + StorageInfo { + disks, + ..Default::default() + } } async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { unimplemented!() diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 7e0e1773c..871dfdafb 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -14,6 +14,7 @@ use crate::global::{ use crate::heal::heal_commands::{HealOpts, HealResultItem, HealScanMode}; use crate::heal::heal_ops::HealObjectFn; use crate::new_object_layer_fn; +use crate::pools::PoolMeta; use crate::store_api::{BackendByte, BackendDisks, BackendInfo, ListMultipartsInfo, ObjectIO, StorageInfo}; use crate::store_err::{ is_err_bucket_exists, is_err_invalid_upload_id, is_err_object_not_found, is_err_version_not_found, StorageError, @@ -52,7 +53,7 @@ use std::{ }; use time::OffsetDateTime; use tokio::fs; -use tokio::sync::Semaphore; +use tokio::sync::{RwLock, Semaphore}; use tracing::{debug, info}; use uuid::Uuid; @@ -67,6 +68,7 @@ pub struct ECStore { pub pools: Vec>, pub peer_sys: S3PeerSys, // pub local_disks: Vec, + pub pool_meta: PoolMeta, } impl ECStore { @@ -167,12 +169,15 @@ impl ECStore { } let peer_sys = S3PeerSys::new(&endpoint_pools); + let mut pool_meta = PoolMeta::new(pools.clone()); + pool_meta.dont_save = true; let ec = ECStore { id: deployment_id.unwrap(), disk_map, pools, peer_sys, + pool_meta, }; set_object_layer(ec.clone()).await; @@ -626,9 +631,7 @@ pub async fn init_local_disks(endpoint_pools: EndpointServerPools) -> Result<()> set_drives.insert(ep.disk_idx, Some(disk.clone())); - if ep.pool_idx.is_some() && ep.set_idx.is_some() && ep.disk_idx.is_some() { - global_set_drives[ep.pool_idx.unwrap()][ep.set_idx.unwrap()][ep.disk_idx.unwrap()] = Some(disk.clone()); - } + global_set_drives[ep.pool_idx as usize][ep.set_idx as usize][ep.disk_idx as usize] = Some(disk.clone()); } } @@ -854,10 +857,7 @@ impl StorageAPI for ECStore { } let backend = self.backend_info().await; - StorageInfo { - backend: Some(backend), - disks, - } + StorageInfo { backend: backend, disks } } async fn list_bucket(&self, opts: &BucketOptions) -> Result> { diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index 9cc6517a5..026c109dd 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -811,15 +811,15 @@ pub struct StorageDisk { pub used_inodes: u64, pub free_inodes: u64, pub local: bool, - pub pool_index: Option, - pub set_index: Option, - pub disk_index: Option, + pub pool_index: i32, + pub set_index: i32, + pub disk_index: i32, } #[derive(Debug, Default)] pub struct StorageInfo { pub disks: Vec, - pub backend: Option, + pub backend: BackendInfo, } #[derive(Debug, Default)] diff --git a/scripts/run.sh b/scripts/run.sh index e7bd04f0f..c67370d10 100755 --- a/scripts/run.sh +++ b/scripts/run.sh @@ -14,8 +14,8 @@ fi export RUSTFS_STORAGE_CLASS_INLINE_BLOCK="512 KB" -# DATA_DIR_ARG="./target/volume/test{0...4}" -DATA_DIR_ARG="./target/volume/test" +DATA_DIR_ARG="./target/volume/test{0...4}" +# DATA_DIR_ARG="./target/volume/test" if [ -n "$1" ]; then DATA_DIR_ARG="$1" From b4008605385ebce7fa059c6422d357da0dfdccc6 Mon Sep 17 00:00:00 2001 From: weisd Date: Tue, 12 Nov 2024 17:38:36 +0800 Subject: [PATCH 3/4] test accountinfo --- Cargo.lock | 1 + ecstore/src/bucket/policy/action.rs | 2 +- ecstore/src/bucket/policy/resource.rs | 8 +++- ecstore/src/store.rs | 23 +++++----- ecstore/src/store_api.rs | 23 +++++++--- router/Cargo.toml | 1 + router/src/handlers.rs | 66 +++++++++++++++++++++++++-- scripts/run.sh | 4 +- 8 files changed, 102 insertions(+), 26 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 926af7b0a..d565ff3b0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2115,6 +2115,7 @@ dependencies = [ "s3s", "serde", "serde-xml-rs", + "serde_json", "serde_urlencoded", "time", "tracing", diff --git a/ecstore/src/bucket/policy/action.rs b/ecstore/src/bucket/policy/action.rs index 2b87bb1a4..d444e52d5 100644 --- a/ecstore/src/bucket/policy/action.rs +++ b/ecstore/src/bucket/policy/action.rs @@ -14,7 +14,7 @@ use super::condition::{ #[derive(Debug, Deserialize, Serialize, Default, Clone, PartialEq, Eq)] -pub struct ActionSet(HashSet); +pub struct ActionSet(pub HashSet); impl ActionSet { pub fn is_match(&self, act: &Action) -> bool { diff --git a/ecstore/src/bucket/policy/resource.rs b/ecstore/src/bucket/policy/resource.rs index cfa8e1cdf..7e2c6472e 100644 --- a/ecstore/src/bucket/policy/resource.rs +++ b/ecstore/src/bucket/policy/resource.rs @@ -41,6 +41,12 @@ pub struct Resource { } impl Resource { + pub fn new(pattern: &str) -> Self { + Self { + pattern: pattern.to_owned(), + rtype: ResourceARNType::ResourceARNS3, + } + } pub fn validate_bucket(&self, bucket: &str) -> Result<()> { self.validate()?; if !wildcard::match_pattern(&self.pattern, bucket) @@ -187,7 +193,7 @@ impl<'de> Deserialize<'de> for Resource { #[derive(Debug, Default, Serialize, Deserialize, Clone, PartialEq, Eq)] #[serde(transparent)] -pub struct ResourceSet(HashSet); +pub struct ResourceSet(pub HashSet); impl ResourceSet { pub fn validate_bucket(&self, bucket: &str) -> Result<()> { diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 21ccceb3b..1057c707d 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -819,7 +819,7 @@ lazy_static! { #[async_trait::async_trait] impl StorageAPI for ECStore { async fn backend_info(&self) -> BackendInfo { - let (standard_scparities, rrscparities) = { + let (standard_sc_parity, rr_sc_parity) = { if let Some(sc) = GLOBAL_StorageClass.get() { let sc_parity = sc .get_parity_for_sc(storageclass::CLASS_STANDARD) @@ -833,17 +833,17 @@ impl StorageAPI for ECStore { } }; - let mut standard_scdata = Vec::new(); - let mut rrscdata = Vec::new(); + let mut standard_sc_data = Vec::new(); + let mut rr_sc_data = Vec::new(); let mut drives_per_set = Vec::new(); let mut total_sets = Vec::new(); for (idx, set_count) in self.set_drive_counts().iter().enumerate() { - if let Some(sc_parity) = standard_scparities { - standard_scdata.push(set_count - sc_parity); + if let Some(sc_parity) = standard_sc_parity { + standard_sc_data.push(set_count - sc_parity); } - if let Some(rr_sc_parity) = rrscparities { - rrscdata.push(set_count - rr_sc_parity); + if let Some(sc_parity) = rr_sc_parity { + rr_sc_data.push(set_count - sc_parity); } total_sets.push(self.pools[idx].set_count); drives_per_set.push(*set_count); @@ -853,12 +853,13 @@ impl StorageAPI for ECStore { backend_type: BackendByte::Erasure, online_disks: BackendDisks::new(), offline_disks: BackendDisks::new(), - standard_scdata, - standard_scparities, - rrscdata, - rrscparities, + standard_sc_data, + standard_sc_parity, + rr_sc_data, + rr_sc_parity, total_sets, drives_per_set, + ..Default::default() } } async fn storage_info(&self) -> StorageInfo { diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index 332de108a..16cf36b38 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -780,7 +780,7 @@ pub struct DeletedObject { // pub replication_state: ReplicationState, } -#[derive(Debug, Default)] +#[derive(Debug, Default, Serialize)] pub enum BackendByte { #[default] Unknown, @@ -824,7 +824,7 @@ pub struct StorageInfo { pub backend: BackendInfo, } -#[derive(Debug, Default)] +#[derive(Debug, Default, Serialize)] pub struct BackendDisks(HashMap); impl BackendDisks { @@ -836,15 +836,24 @@ impl BackendDisks { } } -#[derive(Debug, Default)] +#[derive(Debug, Default, Serialize)] +#[serde(rename_all = "PascalCase", default)] pub struct BackendInfo { pub backend_type: BackendByte, pub online_disks: BackendDisks, pub offline_disks: BackendDisks, - pub standard_scdata: Vec, - pub standard_scparities: Option, - pub rrscdata: Vec, - pub rrscparities: Option, + #[serde(rename = "StandardSCData")] + pub standard_sc_data: Vec, + #[serde(rename = "StandardSCParities")] + pub standard_sc_parities: Vec, + #[serde(rename = "StandardSCParity")] + pub standard_sc_parity: Option, + #[serde(rename = "RRSCData")] + pub rr_sc_data: Vec, + #[serde(rename = "RRSCParities")] + pub rr_sc_parities: Vec, + #[serde(rename = "RRSCParity")] + pub rr_sc_parity: Option, pub total_sets: Vec, pub drives_per_set: Vec, } diff --git a/router/Cargo.toml b/router/Cargo.toml index 27dd7e412..447cd0ecf 100644 --- a/router/Cargo.toml +++ b/router/Cargo.toml @@ -21,3 +21,4 @@ quick-xml = "0.37.0" serde-xml-rs = "0.6.0" ecstore.workspace = true time.workspace = true +serde_json.workspace = true diff --git a/router/src/handlers.rs b/router/src/handlers.rs index 111cacf08..6206dcb8f 100644 --- a/router/src/handlers.rs +++ b/router/src/handlers.rs @@ -1,16 +1,24 @@ +use std::collections::HashSet; + use crate::router::Operation; +use ecstore::bucket::policy::action::{Action, ActionSet}; +use ecstore::bucket::policy::bucket_policy::{BPStatement, BucketPolicy}; +use ecstore::bucket::policy::effect::Effect; +use ecstore::bucket::policy::resource::{Resource, ResourceSet}; +use ecstore::store_api::StorageAPI; use ecstore::utils::xml; +use ecstore::{new_object_layer_fn, store_api::BackendInfo}; use hyper::StatusCode; use matchit::Params; +use s3s::S3ErrorCode; use s3s::{ dto::{AssumeRoleOutput, Credentials, Timestamp}, - s3_error, Body, S3Request, S3Response, S3Result, + s3_error, Body, S3Error, S3Request, S3Response, S3Result, }; -use serde::Deserialize; +use serde::{Deserialize, Serialize}; use serde_urlencoded::from_bytes; use time::{Duration, OffsetDateTime}; use tracing::warn; - #[derive(Deserialize, Debug, Default)] #[serde(rename_all = "PascalCase", default)] pub struct AssumeRoleRequest { @@ -88,6 +96,14 @@ impl Operation for AssumeRoleHandle { } } +#[derive(Debug, Serialize, Default)] +#[serde(rename_all = "PascalCase", default)] +pub struct AccountInfo { + pub account_name: String, + pub server: BackendInfo, + pub policy: BucketPolicy, +} + pub struct AccountInfoHandler {} #[async_trait::async_trait] impl Operation for AccountInfoHandler { @@ -98,7 +114,49 @@ impl Operation for AccountInfoHandler { warn!("AccountInfoHandler cread {:?}", &cred); - return Err(s3_error!(NotImplemented)); + 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())), + }; + + // test policy + + let mut s3_all_act = HashSet::with_capacity(1); + s3_all_act.insert(Action::AllActions); + + let mut all_res = HashSet::with_capacity(1); + all_res.insert(Resource::new("*")); + + let bucket_policy = BucketPolicy { + id: "".to_owned(), + version: "2012-10-17".to_owned(), + statements: vec![BPStatement { + sid: "".to_owned(), + effect: Effect::Allow, + actions: ActionSet(s3_all_act), + resources: ResourceSet(all_res), + ..Default::default() + }], + }; + + // let policy = bucket_policy + // .marshal_msg() + // .map_err(|_e| S3Error::with_message(S3ErrorCode::InternalError, "parse policy failed"))?; + + let backend_info = store.backend_info().await; + + let info = AccountInfo { + account_name: cred.access_key, + server: backend_info, + policy: bucket_policy, + }; + + let output = serde_json::to_string(&info) + .map_err(|_e| S3Error::with_message(S3ErrorCode::InternalError, "parse accountInfo failed"))?; + + Ok(S3Response::new((StatusCode::OK, Body::from(output)))) } } diff --git a/scripts/run.sh b/scripts/run.sh index c67370d10..e7bd04f0f 100755 --- a/scripts/run.sh +++ b/scripts/run.sh @@ -14,8 +14,8 @@ fi export RUSTFS_STORAGE_CLASS_INLINE_BLOCK="512 KB" -DATA_DIR_ARG="./target/volume/test{0...4}" -# DATA_DIR_ARG="./target/volume/test" +# DATA_DIR_ARG="./target/volume/test{0...4}" +DATA_DIR_ARG="./target/volume/test" if [ -n "$1" ]; then DATA_DIR_ARG="$1" From 26b6d3a1ec0908876d65600f6a661dbd83d99d2e Mon Sep 17 00:00:00 2001 From: weisd Date: Thu, 14 Nov 2024 09:06:42 +0800 Subject: [PATCH 4/4] test accountinfo --- ecstore/src/bucket/policy/bucket_policy.rs | 7 ++++--- ecstore/src/bucket/policy/condition/function.rs | 3 +++ ecstore/src/bucket/policy/resource.rs | 4 ++++ router/src/handlers.rs | 2 +- 4 files changed, 12 insertions(+), 4 deletions(-) diff --git a/ecstore/src/bucket/policy/bucket_policy.rs b/ecstore/src/bucket/policy/bucket_policy.rs index 9c38168d1..5dedbf0dc 100644 --- a/ecstore/src/bucket/policy/bucket_policy.rs +++ b/ecstore/src/bucket/policy/bucket_policy.rs @@ -2,6 +2,7 @@ use crate::error::{Error, Result}; // use rmp_serde::Serializer as rmpSerializer; use serde::{Deserialize, Serialize}; use std::collections::HashMap; +use std::collections::HashSet; use super::{ action::{Action, ActionSet, IAMActionConditionKeyMap}, @@ -35,11 +36,11 @@ pub struct BPStatement { pub principal: Principal, #[serde(rename = "Action")] pub actions: ActionSet, - #[serde(rename = "NotAction", default)] + #[serde(rename = "NotAction", skip_serializing_if = "ActionSet::is_empty")] pub not_actions: ActionSet, - #[serde(rename = "Resource")] + #[serde(rename = "Resource", skip_serializing_if = "ResourceSet::is_empty")] pub resources: ResourceSet, - #[serde(rename = "Condition", default)] + #[serde(rename = "Condition", skip_serializing_if = "Functions::is_empty")] pub conditions: Functions, } diff --git a/ecstore/src/bucket/policy/condition/function.rs b/ecstore/src/bucket/policy/condition/function.rs index c744331c7..929b2d4eb 100644 --- a/ecstore/src/bucket/policy/condition/function.rs +++ b/ecstore/src/bucket/policy/condition/function.rs @@ -106,6 +106,9 @@ impl Functions { } set } + pub fn is_empty(&self) -> bool { + self.0.is_empty() + } } impl Debug for Functions { diff --git a/ecstore/src/bucket/policy/resource.rs b/ecstore/src/bucket/policy/resource.rs index 7e2c6472e..409cd9cb7 100644 --- a/ecstore/src/bucket/policy/resource.rs +++ b/ecstore/src/bucket/policy/resource.rs @@ -228,6 +228,10 @@ impl ResourceSet { } false } + + pub fn is_empty(&self) -> bool { + self.0.is_empty() + } } impl AsRef> for ResourceSet { diff --git a/router/src/handlers.rs b/router/src/handlers.rs index 6206dcb8f..ef334d12c 100644 --- a/router/src/handlers.rs +++ b/router/src/handlers.rs @@ -135,7 +135,7 @@ impl Operation for AccountInfoHandler { statements: vec![BPStatement { sid: "".to_owned(), effect: Effect::Allow, - actions: ActionSet(s3_all_act), + actions: ActionSet(s3_all_act.clone()), resources: ResourceSet(all_res), ..Default::default() }],