From 91c25922130198b3fb1478b1088db79c392c29e9 Mon Sep 17 00:00:00 2001 From: weisd Date: Thu, 7 Nov 2024 17:15:23 +0800 Subject: [PATCH] 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"