change pool idx to i32

This commit is contained in:
weisd
2024-11-07 17:15:23 +08:00
parent 476d93d076
commit 91c2592213
10 changed files with 300 additions and 158 deletions
+21 -21
View File
@@ -21,9 +21,9 @@ pub struct Endpoint {
pub url: url::Url,
pub is_local: bool,
pub pool_idx: Option<usize>,
pub set_idx: Option<usize>,
pub disk_idx: Option<usize>,
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,
+24 -11
View File
@@ -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);
+21 -3
View File
@@ -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)
}
},
}
}
+104 -106
View File
@@ -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),
+1
View File
@@ -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;
+91
View File
@@ -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<PoolDecommissionInfo>,
}
#[derive(Debug, Clone)]
pub struct PoolMeta {
pub pools: Vec<PoolStatus>,
pub dont_save: bool,
}
impl PoolMeta {
pub fn new(pools: Vec<Arc<Sets>>) -> 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<String>,
pub decommissioned_buckets: Vec<String>,
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<PoolStatus> {
unimplemented!()
}
async fn get_decommission_pool_space_info(&self, idx: usize) -> Result<PoolSpaceInfo> {
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<StorageDisk>, info: &StorageInfo) -> usize {
for disk in disks.iter() {
// if disk.pool_index < 0 || info.backend.standard_scdata.len() <= disk.pool_index {
// continue;
// }
}
unimplemented!()
}
+24 -3
View File
@@ -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<Vec<BucketInfo>> {
unimplemented!()
+8 -8
View File
@@ -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<Arc<Sets>>,
pub peer_sys: S3PeerSys,
// pub local_disks: Vec<DiskStore>,
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<Vec<BucketInfo>> {
+4 -4
View File
@@ -811,15 +811,15 @@ pub struct StorageDisk {
pub used_inodes: u64,
pub free_inodes: u64,
pub local: bool,
pub pool_index: Option<usize>,
pub set_index: Option<usize>,
pub disk_index: Option<usize>,
pub pool_index: i32,
pub set_index: i32,
pub disk_index: i32,
}
#[derive(Debug, Default)]
pub struct StorageInfo {
pub disks: Vec<StorageDisk>,
pub backend: Option<BackendInfo>,
pub backend: BackendInfo,
}
#[derive(Debug, Default)]
+2 -2
View File
@@ -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"